如何在Elasticsearch中自动计算日累计用水量差值并存储
实现单日总用水量自动存储的几种方案
注意:确保原始数据中的时间字符串已正确解析为Elasticsearch的
date类型字段(如@timestamp),否则时间相关的聚合和查询无法正常工作。可以使用Ingest Pipeline的date处理器完成解析,示例:{ "date": { "field": "time_str", "target_field": "@timestamp", "formats": ["yyyyMMdd HH'h'mm"] } }
方案1:使用Elasticsearch Transform创建持久化日汇总索引
Transform可以定时将小时级原始数据聚合为日级统计数据,自动计算单日总用水量并存储到独立的汇总索引中,适合需要长期保留日统计结果的场景。
配置步骤:
创建Transform任务
定义源索引(你的用水数据索引)、目标索引(存储日统计的索引),通过scripted_metric聚合计算当日00:00与次日00:00的累计用水量差值:PUT _transform/daily-water-consumption { "source": { "index": "water-consumption-*" }, "dest": { "index": "daily-water-summary" }, "pivot": { "group_by": { "day": { "date_histogram": { "field": "@timestamp", "calendar_interval": "day", "time_zone": "UTC" } } }, "aggregations": { "daily_total": { "scripted_metric": { "init_script": "state.consumptions = []", "map_script": "state.consumptions.add(doc['consumption'].value)", "combine_script": "if (state.consumptions.isEmpty()) return 0; def sorted = state.consumptions.sort(); return sorted[-1] - sorted[0]", "reduce_script": "def total = 0; for (s in states) { total += s; } return total" } } } }, "sync": { "time": { "field": "@timestamp", "delay": "1h" } } }scripted_metric会收集当日所有累计用水量值,取最大值(次日00:00的数值)与最小值(当日00:00的数值)的差值,得到单日总用水量。
启动Transform任务
POST _transform/daily-water-consumption/_start任务启动后会自动同步新写入的原始数据,持续更新日汇总索引。
方案2:使用Ingest Pipeline在写入次日数据时实时更新汇总文档
这种方式在写入次日00:00的累计数据时,自动计算前一天的总用水量并写入专门的汇总文档,适合需要实时获取当日统计结果的场景。
配置步骤:
创建Ingest Pipeline
通过脚本查询前一天00:00的累计值,计算差值后写入汇总索引:PUT _ingest/pipeline/calculate-daily-consumption { "processors": [ { "script": { "if": "ctx['@timestamp'].getHourOfDay() == 0", "source": """ // 生成前一天00:00的时间戳 def prevDay = ctx['@timestamp'].minusDays(1).withHourOfDay(0).withMinuteOfHour(0).withSecondOfMinute(0).withMillisOfSecond(0); // 查询前一天00:00的累计用水量 def searchResp = ctx._ingest.client.search({ index: 'water-consumption-*', query: { term: { '@timestamp': prevDay } }, size: 1 }); if (searchResp.hits.hits.length > 0) { def prevConsumption = searchResp.hits.hits[0]._source.consumption; def dailyTotal = ctx.consumption - prevConsumption; // 写入汇总文档 ctx._ingest.client.index({ index: 'daily-water-summary', id: prevDay.format('yyyyMMdd'), body: { date: prevDay, daily_total: dailyTotal, start_consumption: prevConsumption, end_consumption: ctx.consumption } }); } """ } } ] }关联Pipeline到原始数据写入流程
可以在创建原始索引时设置默认Pipeline,或写入数据时指定Pipeline:PUT water-consumption/_doc/1?pipeline=calculate-daily-consumption { "@timestamp": "2023-09-15T00:00:00Z", "consumption": 199 }
方案3:使用Runtime Fields动态计算(无需持久化存储)
如果不需要持久化存储单日总用水量,仅在查询时需要结果,可以使用Runtime Fields动态计算,避免额外存储开销:
配置索引的Runtime Field
PUT water-consumption/_mapping { "runtime": { "daily_total": { "type": "double", "script": { "source": """ def dayStart = doc['@timestamp'].value.truncatedTo(ChronoUnit.DAYS); def dayEnd = dayStart.plusDays(1); // 查询当日及次日00:00的累计值 def searchReq = new SearchRequest('water-consumption-*'); searchReq.source().query( boolQuery().filter(rangeQuery('@timestamp').gte(dayStart).lte(dayEnd)) ).size(2); def resp = ctx._index.search(searchReq); def values = resp.getHits().getHits().stream() .map(hit -> hit.getSourceAsMap().get('consumption')) .sorted() .collect(Collectors.toList()); emit(values.size() >= 2 ? values.get(1) - values.get(0) : 0); """ } } } }
查询时直接调用daily_total字段即可获取结果,但该方式每次查询都会执行脚本,性能逊于前两种方案。
内容的提问来源于stack exchange,提问作者Jeff Carabin
相关产品推荐
相关产品推荐

