如何用Elasticsearch按天聚合CSV数据并在Python中分析提交峰值
按天拆分Elasticsearch数据并分析每日峰值方案
一、Elasticsearch端按天聚合数据
首先确认FinalSubmissionTime字段为date类型(若导入时未配置,需先更新索引映射),通过嵌套聚合实现按天拆分+小时级统计:
聚合查询语句
POST /your_index_name/_search?size=0 { "aggs": { "daily_groups": { "date_histogram": { "field": "FinalSubmissionTime", "calendar_interval": "day", "format": "yyyy-MM-dd" }, "aggs": { "hourly_groups": { "date_histogram": { "field": "FinalSubmissionTime", "calendar_interval": "hour", "format": "HH:mm" }, "aggs": { "total_product_count": { "sum": { "field": "number of product" } } } } } } } }
size=0:仅返回聚合结果,不加载原始数据,提升查询效率- 外层
date_histogram按天分组,内层嵌套按小时拆分 sum聚合计算每小时的待提交产品总数量
二、Python端获取聚合结果并分析峰值
使用elasticsearch库连接ES,提取聚合数据后筛选每日峰值时段:
代码示例
from elasticsearch import Elasticsearch # 连接Elasticsearch(根据实际配置调整地址、端口、认证信息) es = Elasticsearch(["http://localhost:9200"]) # 定义聚合查询 agg_query = { "aggs": { "daily_groups": { "date_histogram": { "field": "FinalSubmissionTime", "calendar_interval": "day", "format": "yyyy-MM-dd" }, "aggs": { "hourly_groups": { "date_histogram": { "field": "FinalSubmissionTime", "calendar_interval": "hour", "format": "HH:mm" }, "aggs": { "total_product_count": { "sum": { "field": "number of product" } } } } } } }, "size": 0 } # 执行查询 response = es.search(index="your_index_name", body=agg_query) # 提取并分析每日峰值 daily_peak_data = [] for day_bucket in response["aggregations"]["daily_groups"]["buckets"]: current_date = day_bucket["key_as_string"] hourly_buckets = day_bucket["hourly_groups"]["buckets"] if hourly_buckets: # 筛选出当日产品数量最多的时段 peak_hour = max(hourly_buckets, key=lambda x: x["total_product_count"]["value"]) daily_peak_data.append({ "date": current_date, "peak_hour": peak_hour["key_as_string"], "total_products": peak_hour["total_product_count"]["value"] }) # 输出结果 for item in daily_peak_data: print(f"日期: {item['date']}, 峰值时段: {item['peak_hour']}, 产品数量: {item['total_products']}")
扩展说明
- 若需按
KindOfProduct细分峰值,可在小时聚合下嵌套terms聚合按产品类型分组 - 针对超大规模数据,可在查询中添加
range条件限定时间范围,避免全量聚合
内容的提问来源于stack exchange,提问作者Andreas K
相关产品推荐
相关产品推荐

