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

Elasticsearch Python脚本循环无法切换至下一个forwarder求助

问题分析与修复:Elasticsearch Python脚本无法遍历所有forwarder

原始代码核心问题

最初的build_query函数中,循环第一次迭代就执行了return query,函数直接返回第一个forwarder对应的查询语句,后续forwarder完全没有机会进入循环处理。

修改后代码的遗留问题

你调整了循环逻辑,将forwarder遍历移到main函数中,但仍存在几个关键错误:

  1. 参数硬编码覆盖:build_query里重新赋值start_date = "now-1d/d"和end_date = "now/d",完全忽略了main中传入的时间切片参数,导致时间范围固定,失去了时间切片的意义。
  2. 参数名称不匹配:main中调用build_query时传的参数是forwader,但函数定义的参数是forwarder,直接触发参数错误。
  3. 结果保存错误:最后保存CSV时用的是parsed_results,这只是最后一次循环的结果,而非累计的total_results_counter,导致大部分数据丢失。

修复后的完整代码

from collections import Counter
import datetime

# 假设以下为已定义的依赖函数/常量,根据实际情况调整
LOCAL_UTC_OFF_SET = "+08:00"
def parse_args():
    return type('Args', (), {'config_file': None, 'days': 1, 'start': None, 'end': None, 'filename': 'results.csv'})()
def load_config(file, keys):
    return {'server': 'localhost', 'port': 9200, 'username': 'user', 'password': 'pass'}
def time_range_slices(hour_slice, start, end, days):
    now = datetime.datetime.now()
    for i in range(days):
        yield (now - datetime.timedelta(days=i+1)).strftime("%Y-%m-%dT00:00:00Z"), now.strftime("%Y-%m-%dT23:59:59Z")
class DataLakeSession:
    def __init__(self, server, port, username, password):
        self.server = server
        self.port = port
        self.username = username
        self.password = password
def run_agg_query(session, query):
    return {'aggregations': {'1': {'buckets': [{'key': 'dest1', 'doc_count': 10}, {'key': 'dest2', 'doc_count': 5}]}}}
def parse_agg_results(results):
    parsed = {}
    for bucket in results['aggregations']['1']['buckets']:
        parsed[bucket['key']] = bucket['doc_count']
    return parsed
def to_csv(counter, filename):
    with open(filename, 'w') as f:
        f.write('dest_host,count\n')
        for key, count in counter.items():
            f.write(f'{key},{count}\n')
    return filename

def main():
    args = parse_args()
    if args.config_file:
        config = load_config(
            args.config_file, ("server", "port", "username", "password")
        )
    else:
        config = {}
    server = config["server"]
    port = config["port"]
    username = config["username"]
    password = config["password"]
    days = args.days
    start = args.start
    end = args.end

    forwarders = (
        "SVGCEXABEC01",
        "SVGCEXABEC03",
        "SV04EXABEC02",
        "SV04EXABEC03",
        "SV07EXABEC01",
        "SV07EXABEC02",
        "SV07EXABEC03",
        "SVGCEXABEC02"
    )

    total_results_counter = Counter()
    for start_date, end_date in time_range_slices(
        hour_slice=1,
        start=start,
        end=end,
        days=days,
    ):
        for forwarder in forwarders:
            # 参数名称统一为forwarder,匹配函数定义
            query = build_query(start_date=start_date, end_date=end_date, forwarder=forwarder)

            try:
                time_stamp = datetime.datetime.now()
                print(
                    f"Time Stamp: {time_stamp} | Start: {start_date} | End: {end_date} | Forwarder: {forwarder}",
                    end="",
                    flush=True,
                )
                session = DataLakeSession(server, port, username, password)
                
                results = run_agg_query(session=session, query=query)
                parsed_results = parse_agg_results(results)

                print(f" | Total: {len(parsed_results)}", flush=True)
                total_results_counter.update(parsed_results)

            except Exception as e:
                print(f" | Error getting results. Error: {e}", flush=True)

    # 使用累计的total_results_counter保存所有数据
    save_path = to_csv(total_results_counter, args.filename)
    
    print(f"CSV file saved to: {save_path}")


def build_query(start_date, end_date, forwarder):
    # 移除硬编码,使用传入的时间参数
    print(f"Processing forwarder: {forwarder}")
    query = {
        "size": 10_000,
        "query": {
            "bool": {
                "must": [
                    {
                        "query_string": {
                            "analyze_wildcard": True,
                            "default_field": "message",
                            "query": f'forwarder:"{forwarder}" AND (host.keyword:/LN.*/ OR host.keyword:/PC.*/ OR host.keyword:/TB.*/)',  # noqa: E501
                        }
                    },
                    {
                        "range": {
                            "@timestamp": {
                                "gte": start_date,
                                "lte": end_date,
                                "time_zone": LOCAL_UTC_OFF_SET,
                            }
                        }
                    },
                ],
                "must_not": [],
            }
        },
        "aggs": {
            "1": {
                "terms": {
                    "field": "dest_host.keyword",
                    "size": 9_999,
                    "shard_size": 9_999,
                    "order": {"_count": "asc"},
                },
            },
        },
    }
    return query

if __name__ == "__main__":
    main()

关键修复点总结

  • 移除build_query中对start_date和end_date的硬编码,保留传入的参数值,确保时间切片生效。
  • 统一参数名称:将forwader改为forwarder,确保函数调用与定义一致。
  • 保存CSV时使用total_results_counter,确保所有循环的结果都被累计保存。
  • 移除build_query内部的forwarder循环,改为在main中遍历forwarder,每次调用build_query生成对应forwarder的查询语句。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 14:13:08