Elasticsearch Python脚本循环无法切换至下一个forwarder求助
问题分析与修复:Elasticsearch Python脚本无法遍历所有forwarder
原始代码核心问题
最初的build_query函数中,循环第一次迭代就执行了return query,函数直接返回第一个forwarder对应的查询语句,后续forwarder完全没有机会进入循环处理。
修改后代码的遗留问题
你调整了循环逻辑,将forwarder遍历移到main函数中,但仍存在几个关键错误:
- 参数硬编码覆盖:
build_query里重新赋值start_date = "now-1d/d"和end_date = "now/d",完全忽略了main中传入的时间切片参数,导致时间范围固定,失去了时间切片的意义。 - 参数名称不匹配:
main中调用build_query时传的参数是forwader,但函数定义的参数是forwarder,直接触发参数错误。 - 结果保存错误:最后保存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
相关产品推荐
相关产品推荐

