
Elasticsearch作为一款强大的搜索和分析引擎,在日志分析、全文检索、数据分析等领域有着广泛应用。本文将介绍如何通过Python自动化操作Elasticsearch,包括覆盖 ES 8.x 认证、CRUD、滚动查询、reindex、异步 IO、数据流、机器学习、Transform 等场景,全部提供可直接执行的代码。
Python操作Elasticsearch的典型应用场景包括:
1、日志分析系统:自动化收集、索引和分析应用日志。
2、数据管道:作为ETL流程的一部分,将处理后的数据批量导入ES。
3、监控告警:定期查询ES数据,触发异常告警。
4、集群管理:自动化执行索引维护、副本调整等运维任务。
5、数据迁移:在不同集群或索引间迁移数据。
pip install elasticsearch==8.8.2 # 版本应与ES版本保持一致from elasticsearch import Elasticsearch
from loguru import logger
class ES:
def __init__(self):
self.client = Elasticsearch(
hosts=' https://localhost:9200',
ca_certs="certs/http_ca.crt",
basic_auth=('elastic', 'your_password')
)
def info(self):
logger.info(self.client.info())
if __name__ == '__main__':
es = ES()
es.info()class ES:
def __init__(self):
self.client = Elasticsearch(
hosts=' https://localhost:9200',
ssl_assert_fingerprint="your_fingerprint",
basic_auth=('elastic', 'your_password')
)class ES:
def __init__(self):
self.client = Elasticsearch(
hosts=' https://localhost:9200',
ca_certs="certs/http_ca.crt",
bearer_auth="your_token"
)def query_count(self, index: str, query: dict) -> None:
"""
DSL查询结果统计
:param index: 查询索引
:param query: 查询语句
"""
result = self.client.count(index=index, query=query)
logger.info(result)
# 使用示例
query = {
"bool": {
"filter": [
{"terms": {"access_status": ["502", "503", "504"]}},
{"range": {"@timestamp": {"gte": "2023-08-07T17:20:00.000+08:00"}}}
]
}
}
es.query_count('logs-myapp-default', query)def query_dsl(self, index: str, query: dict, sort: list, size: int) -> None:
"""
获取DSL查询结果的内容
:param index: 查询索引
:param query: 查询语句
:param sort: 排序参数
:param size: 分页参数
"""
result = self.client.search(index=index, query=query, sort=sort, size=size)
logger.info(result)
# 使用示例
sort = [{"@timestamp": {"order": "asc"}}]
es.query_dsl('logs-myapp-default', query, sort, 20)def insert_data(self, index: str, data: dict) -> None:
"""
ES中插入单条数据
:param index: 索引
:param data: 数据内容
"""
logger.info(self.client.index(index=index, document=data))
# 使用示例
data = {"username": "alex", "age": 30}
es.insert_data('user-info', data)from elasticsearch import helpers
def insert_bulk(self, data: list) -> None:
"""
ES中批量插入数据
:param data: 数据列表,每条需包含'_index'字段
"""
logger.info(helpers.bulk(self.client, data))
# 使用示例
data = [
{'_index': 'user-info', "username": "alex1", "age": 11},
{'_index': 'user-info', "username": "alex2", "age": 12}
]
es.insert_bulk(data)# 更新数据
def update(self, index: str, id: str, data: dict) -> None:
logger.info(self.client.update(index=index, id=id, doc=data))
# 删除数据
def delete(self, index: str, id: str) -> None:
logger.info(self.client.delete(index=index, id=id))def query_scan(self, index: str, query: dict) -> None:
"""
获取DSL滚动查询结果的内容
:param index: 查询索引
:param query: 查询语句
"""
result = helpers.scan(self.client, query=query, index=index, scroll='2m', size=1000)
data = [i for i in result]
logger.info(f"共查询到{len(data)}条记录")def reindex(self, source_index: str, target_index: str, query: dict) -> None:
"""
ES reindex指定条件的数据到新的index中
:param source_index: 原数据索引
:param target_index: 目的数据索引
:param query: 查询条件
"""
logger.info(helpers.reindex(self.client,
source_index=source_index,
target_index=target_index,
query=query))通过HTTP API实现集群层面的管理:
import httpx
import ssl
class ESCluster:
def __init__(self):
self.base_url = " https://localhost:9200 "
self.auth = ('elastic', 'your_password')
self.context = ssl.create_default_context()
self.context.load_verify_locations(cafile='certs/http_ca.crt')
def get_cluster_health(self):
"""获取集群健康状态"""
res = httpx.get(f"{self.base_url}/_cluster/health",
auth=self.auth, verify=self.context)
return res.json()
def adjust_replicas(self, index: str, replicas: int):
"""调整索引副本数"""
data = {"index": {"number_of_replicas": replicas}}
res = httpx.put(f"{self.base_url}/{index}/_settings",
json=data, auth=self.auth, verify=self.context)
return res.json()
# 使用示例
cluster = ESCluster()
health = cluster.get_cluster_health()
if health['status'] != 'green':
# 自动修复yellow状态
cluster.adjust_replicas('problem-index', 0)# es_client.py
fromdatetimeimporttimedelta
fromelasticsearchimportElasticsearch
fromloguruimportlogger
classESClient:
def__init__(self, hosts: str|list,
ca_certs: str|None = None,
basic_auth: tuple[str, str] |None = None,
bearer_auth: str|None = None,
fingerprint: str|None = None):
"""统一客户端:支持证书、指纹、Token 多种认证"""
self.client = Elasticsearch(
hosts,
ca_certs=ca_certs,
ssl_assert_fingerprint=fingerprint,
basic_auth=basic_auth,
bearer_auth=bearer_auth,
retry_on_timeout=True,
raise_on_error=True,
verify_certs=True,
)
logger.info("ES cluster: {}", self.client.info())
# ------------------ 基础 CRUD ------------------
defcount(self, index: str, query: dict) ->int:
returnself.client.count(index=index, query=query).raw["count"]
defsearch(self, index: str, query: dict, sort: list|None = None, size: int = 10) ->list[dict]:
res = self.client.search(index=index, query=query, sort=sortor [], size=size)
return [hit["_source"] forhitinres.raw["hits"]["hits"]]
defindex_doc(self, index: str, doc: dict, doc_id: str|None = None) ->str:
resp = self.client.index(index=index, id=doc_id, document=doc)
returnresp.meta["id"]
defbulk_docs(self, actions: list[dict]) ->tuple[int, list]:
"""返回 (成功数, 失败列表)"""
stats, *details = self.client.bulk(operations=actions, refresh=True)
fails = [itemforitemindetailsifitem.get("index", {}).get("error")]
returnstats["items"], fails
# ------------------ 滚动 & reindex ------------------
defscroll_all(self, index: str, query: dict, scroll_time: str = "2m"):
"""生成器,逐条 yield"""
forhitinself.client.helpers.scan(
index=index, query=query, scroll=timedelta(minutes=2)
):
yieldhit["_source"]
defreindex_by_query(self, source: str, target: str, query: dict):
body = {"source": {"index": source, "query": query}, "dest": {"index": target}}
returnself.client.reindex(body, wait_for_completion=True)
# ------------------ 集群运维 ------------------
defhealth(self) ->dict:
returnself.client.cluster.health()
defset_replicas(self, index: str, replicas: int):
returnself.client.indices.put_settings(
index=index, settings={"number_of_replicas": replicas}
)fromes_clientimportESClient
fromdatetimeimportdatetime, timezone
es = ESClient(
hosts="https://localhost:9200",
ca_certs="certs/http_ca.crt",
basic_auth=("elastic", "your_password")
)
# 1. 统计 5xx 日志
query = {
"bool": {
"filter": [
{"terms": {"access_status": ["502", "503", "504"]}},
{"range": {"@timestamp": {"gte": "now-1h"}}}
]
}
}
print("5xx count:", es.count("logs-myapp-default", query))
# 2. DSL 查询
sort = [{"@timestamp": {"order": "desc"}}]
hits = es.search("logs-myapp-default", query, sort, 20)
forhinhits:
logger.info("path={}, status={}", h.get("url"), h.get("access_status"))
# 3. 单条插入
doc = {"username": "alex", "age": 30, "@timestamp": datetime.now(timezone.utc)}
es.index_doc("user-info", doc)
# 4. bulk 插入
actions = [
{"index": {"_index": "user-info"}},
{"username": f"user-{i}", "age": 20+i, "@timestamp": datetime.now(timezone.utc)}
foriinrange(100)
]
succ, fails = es.bulk_docs(actions)
logger.info("bulk success={} fails={}", succ, fails)# 滚动导出全量
withopen("export.jsonl", "w", encoding="utf8") asf:
fordocines.scroll_all("logs-myapp-default", {"match_all": {}}):
f.write(f"{doc}\n")
# reindex 近 1 天数据到新索引
es.reindex_by_query(
source="logs-myapp-default",
target="logs-myapp-archive-2024.06",
query={"range": {"@timestamp": {"gte": "now-1d"}}}
)# 08-async-bulk.py
importasyncio
fromelasticsearchimportAsyncElasticsearch
fromdatetimeimportdatetime, timezone
asyncdefmain():
es = AsyncElasticsearch(
"https://localhost:9200",
ca_certs="certs/http_ca.crt",
basic_auth=("elastic", "your_password"),
retry_on_timeout=True,
)
actions = [
{"index": {"_index": "async-demo"}},
{"user": f"u{i}", "counter": i, "@timestamp": datetime.now(timezone.utc)}
foriinrange(10_000)
]
awaites.bulk(operations=actions, refresh=True)
awaites.close()
if__name__ == "__main__":
asyncio.run(main())# 09-data-stream.py
PUT_index_template/myapp-logs
{
"index_patterns": ["myapp-logs-*"],
"data_stream": {},
"template": {
"settings": {
"number_of_shards": 1,
"number_of_replicas": 1,
"index.lifecycle.name": "myapp-logs-policy"
}
}
}
# Python 写入数据流
es.client.index(index="myapp-logs-default", document={
"message": "order placed",
"@timestamp": datetime.now(timezone.utc)
})# 10-ml-anomaly.py
es.client.ml.put_job(job_id="order-count-anomaly", body={
"analysis_config": {
"bucket_span": "10m",
"detectors": [{"function": "count"}]
},
"data_description": {"time_field": "@timestamp"}
})
es.client.ml.open_job(job_id="order-count-anomaly")
# 后续通过 datafeed 持续喂数据,省略# 11-transform.py
es.client.transform.put_transform(
transform_id="order-hourly-stats",
body={
"source": {"index": "orders-*"},
"dest": {"index": "order-hourly-rollup"},
"frequency": "1h",
"pivot": {
"group_by": {
"hour": {"date_histogram": {"field": "@timestamp", "calendar_interval": "1h"}}
},
"aggregations": {
"total_amount": {"sum": {"field": "amount"}}
}
}
}
)
es.client.transform.start(transform_id="order-hourly-stats")health = es.health()
logger.info("cluster status: {}", health["status"])
# 临时把副本降到 0 快速恢复 yellow
ifhealth["status"] == "yellow":
es.set_replicas("problem-index", 0)本文介绍了Python操作Elasticsearch的多种方式,包括:
1、多种认证连接方式确保安全访问。
2、基础CRUD操作满足日常数据管理需求。
3、高级功能如滚动查询、批量操作提升处理效率。
4、集群API调用实现自动化运维。
通过Python与Elasticsearch的结合,开发者可以构建强大的数据分析和搜索应用,同时实现运维自动化,大幅提升工作效率。
“无他,惟手熟尔”!有需要的用起来。
本文分享自 Nicholas与Pypi 微信公众号,前往查看
如有侵权,请联系 cloudcommunity@tencent.com 删除。
本文参与 腾讯云自媒体同步曝光计划 ,欢迎热爱写作的你一起参与!