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

如何在Elasticsearch中自动计算日累计用水量差值并存储

实现单日总用水量自动存储的几种方案

注意:确保原始数据中的时间字符串已正确解析为Elasticsearch的date类型字段(如@timestamp),否则时间相关的聚合和查询无法正常工作。可以使用Ingest Pipeline的date处理器完成解析,示例:

{
  "date": {
    "field": "time_str",
    "target_field": "@timestamp",
    "formats": ["yyyyMMdd HH'h'mm"]
  }
}

方案1:使用Elasticsearch Transform创建持久化日汇总索引

Transform可以定时将小时级原始数据聚合为日级统计数据,自动计算单日总用水量并存储到独立的汇总索引中,适合需要长期保留日统计结果的场景。

配置步骤:

  1. 创建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的数值)的差值,得到单日总用水量。
  2. 启动Transform任务

    POST _transform/daily-water-consumption/_start
    

    任务启动后会自动同步新写入的原始数据,持续更新日汇总索引。

方案2:使用Ingest Pipeline在写入次日数据时实时更新汇总文档

这种方式在写入次日00:00的累计数据时,自动计算前一天的总用水量并写入专门的汇总文档,适合需要实时获取当日统计结果的场景。

配置步骤:

  1. 创建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
                  }
                });
              }
            """
          }
        }
      ]
    }
    
  2. 关联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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 23:26:28