You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何高效导出250亿条Elastic Security事件?

大规模Elastic Security事件备份问题

作为MSSP人员,正在为客户执行下线操作,需备份250亿条Elastic Security事件以满足合规要求。目前通过Python API查询索引并提取事件JSON,但脚本耗时极长;尝试过Discover标签的分享功能导出,频繁超时报错;API导出同样耗时过久。当前脚本可完成认证、查询、持久化存储,计划用searchAfter实现中断后重跑并封装为bash脚本,但searchAfter无法正常工作。

当前使用的Python脚本

#!/usr/bin/env python3

from elasticsearch import Elasticsearch
import sys
import time

username = "USER"
password = "REDACTED"

elasticsearchEndpoint = "https://elastic.com:9243"

es = Elasticsearch(
    elasticsearchEndpoint, 
    basic_auth=(username, password), 
    request_timeout=120
)

SAI = "searchAfterIndex.txt"
file_path = "dump.out"
options = ["/rr", "rerun", "rr", "-rr", "--rr", "--rerun", "/rerun", "-rerun"]


query = {
    "bool": {
        "filter": [
            {
                "range": {
                    "@timestamp": {
                        "gte": "now-1y",
                        "lt": "now"
                    }
                }
            }
        ],
        "must": [],
        "must_not": [],
        "should": []
    }
}

def getEvents(searchAfter=None):

    while True:
        if searchAfter is None:
            classVar = es.search(expand_wildcards='all', index="*-abc**", size=10000, sort="@timestamp:asc", query=query) # IDK if I need query=query
        else:
            time.sleep(.025)
            classVar = es.search(expand_wildcards='all', index="*-abc**", size=10000, search_after=searchAfter, sort="@timestamp:asc", query=query)
        
        hits = classVar.get("hits", {}).get("hits", [])

        with open(file_path, "a") as fp:
            for hit in hits:
                fp.write(str(hit))

        # Break condition
        if len(hits) == 10000:
            searchAfter = hits[-1]["sort"]
            with open(SAI, "w") as fp:
                fp.write(str(searchAfter).strip("[]"))
        else:
            break

if __name__ == "__main__":

    if len(sys.argv)>=2 and sys.argv[1] in options:
        try:
            searchAfterIndex = open(SAI, "r").read()
            getEvents(searchAfterIndex)
            print("Done!", file=sys.stderr)
        except Exception as e:
            print(f"An error occurred: {e}", file=sys.stderr)

    try:
        getEvents()
        print("Done!", file=sys.stderr)
    except Exception as e:
        print(f"An error occurred: {e}", file=sys.stderr)

问题修复与性能优化方案

1. searchAfter失效问题修复

当前脚本重跑时,读取的searchAfter是字符串格式,但Elasticsearch的search_after参数需要数组类型(对应hits[-1]["sort"]的原生格式),格式不匹配导致功能失效。

修复步骤:

  • 写入searchAfter时用JSON格式保存数组:
    import json
    # ...
    with open(SAI, "w") as fp:
        json.dump(searchAfter, fp)
    
  • 读取时解析JSON数组:
    import json
    # ...
    searchAfterIndex = json.load(open(SAI, "r"))
    getEvents(searchAfterIndex)
    

2. 性能优化建议

针对250亿条数据的大规模导出,单脚本分页效率不足,可从以下方向优化:

  • 改用Scroll API:对于超大规模全量导出,Scroll API初始阶段性能更稳定,适合一次性批量获取数据(注意设置合理的滚动过期时间)。
  • 批量写入文件:积累一定数量文档后批量写入,减少IO操作次数,降低磁盘开销。
  • 并行导出:按时间范围拆分查询任务,用多线程/多进程同时导出不同时间段的数据,利用Elasticsearch集群的并行处理能力。
  • 调整size参数:在集群性能允许的前提下,适当增大size值(如50000),减少API请求次数。
  • 优化序列化逻辑:直接写入原始JSON字符串,避免str(hit)的格式问题和性能损耗:
    with open(file_path, "a") as fp:
        for hit in hits:
            fp.write(json.dumps(hit) + "\n")
    

3. 脚本逻辑修正

主函数中无论是否传入重跑参数,都会执行一次getEvents(),导致重复导出。调整逻辑为仅执行对应分支:

if __name__ == "__main__":
    try:
        if len(sys.argv)>=2 and sys.argv[1] in options:
            searchAfterIndex = json.load(open(SAI, "r"))
            getEvents(searchAfterIndex)
        else:
            getEvents()
        print("Done!", file=sys.stderr)
    except Exception as e:
        print(f"An error occurred: {e}", file=sys.stderr)

内容的提问来源于stack exchange,提问作者negativeTenEng

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.08 17:07:52