首页
学习
活动
专区
圈层
工具
发布
社区首页 >专栏 >Python 自动化管理Elasticsearch集群

Python 自动化管理Elasticsearch集群

作者头像
用户11081884
发布2026-07-20 18:50:16
发布2026-07-20 18:50:16
260
举报

Elasticsearch作为一款强大的搜索和分析引擎,在日志分析、全文检索、数据分析等领域有着广泛应用。本文将介绍如何通过Python自动化操作Elasticsearch,包括覆盖 ES 8.x 认证、CRUD、滚动查询、reindex、异步 IO、数据流、机器学习、Transform 等场景,全部提供可直接执行的代码。

使用场景

Python操作Elasticsearch的典型应用场景包括:

1、日志分析系统:自动化收集、索引和分析应用日志。

2、数据管道:作为ETL流程的一部分,将处理后的数据批量导入ES。

3、监控告警:定期查询ES数据,触发异常告警。

4、集群管理:自动化执行索引维护、副本调整等运维任务。

5、数据迁移:在不同集群或索引间迁移数据。

安装与连接

安装Python客户端

代码语言:javascript
复制
pip install elasticsearch==8.8.2  # 版本应与ES版本保持一致

多种认证连接方式

1. CA证书+Basic Auth认证
代码语言:javascript
复制
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()
2. 证书指纹+Basic Auth认证
代码语言:javascript
复制
class ES:
    def __init__(self):
        self.client = Elasticsearch(
            hosts=' https://localhost:9200',
            ssl_assert_fingerprint="your_fingerprint",
            basic_auth=('elastic', 'your_password')
        )
3. CA证书+Token认证
代码语言:javascript
复制
class ES:
    def __init__(self):
        self.client = Elasticsearch(
            hosts=' https://localhost:9200',
            ca_certs="certs/http_ca.crt",
            bearer_auth="your_token"
        )

基础数据操作

查询数据统计

代码语言:javascript
复制
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)

DSL查询

代码语言:javascript
复制
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)

数据插入操作

单条数据插入
代码语言:javascript
复制
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)
批量数据插入
代码语言:javascript
复制
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)

数据更新与删除

代码语言:javascript
复制
# 更新数据
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))

进阶操作

滚动查询(处理大数据量)

代码语言:javascript
复制
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)}条记录")

索引重建(Reindex)

代码语言:javascript
复制
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实现集群层面的管理:

代码语言:javascript
复制
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)

实践案例

代码语言:javascript
复制
# 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}
        )

1、基础数据操作

代码语言:javascript
复制
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)

2、进阶操作

代码语言:javascript
复制
# 滚动导出全量
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"}}}
)

3、异步大批量写入

代码语言:javascript
复制
# 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())

4、数据流与 ILM

代码语言:javascript
复制
# 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)
})

5、机器学习异常检测

代码语言:javascript
复制
# 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 持续喂数据,省略

6、Transform 持续聚合

代码语言:javascript
复制
# 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")

7、集群管理

代码语言:javascript
复制
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调用实现自动化运维。

通过PythonElasticsearch的结合,开发者可以构建强大的数据分析和搜索应用,同时实现运维自动化,大幅提升工作效率。

“无他,惟手熟尔”!有需要的用起来。

本文参与 腾讯云自媒体同步曝光计划,分享自微信公众号。
原始发表:2025-10-15,如有侵权请联系 cloudcommunity@tencent.com 删除

本文分享自 Nicholas与Pypi 微信公众号,前往查看

如有侵权,请联系 cloudcommunity@tencent.com 删除。

本文参与 腾讯云自媒体同步曝光计划  ,欢迎热爱写作的你一起参与!

评论
登录后参与评论
0 条评论
热度
最新
推荐阅读
目录
  • 使用场景
  • 安装与连接
    • 安装Python客户端
    • 多种认证连接方式
      • 1. CA证书+Basic Auth认证
      • 2. 证书指纹+Basic Auth认证
      • 3. CA证书+Token认证
  • 基础数据操作
    • 查询数据统计
    • DSL查询
    • 数据插入操作
      • 单条数据插入
      • 批量数据插入
    • 数据更新与删除
  • 进阶操作
    • 滚动查询(处理大数据量)
    • 索引重建(Reindex)
  • 集群管理操作
    • 实践案例
  • 1、基础数据操作
  • 2、进阶操作
  • 3、异步大批量写入
  • 4、数据流与 ILM
  • 5、机器学习异常检测
  • 6、Transform 持续聚合
  • 7、集群管理
领券
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档