Transform 与物化聚合

讲解 Elasticsearch Transform 物化聚合:pivot 与 latest 两种模式、group_by 与聚合定义、连续转换的 sync 调度与 checkpoint 机制、历史数据回填与目标索引重建,以及 Transform 与 rollup 在时序降采样场景下的取舍。

聚合查询的代价随数据量线性增长:对十亿级文档做一次 date_histogram 加嵌套 terms,即使有 doc_values 与查询缓存兜底,也可能要几秒到几十秒。当同一个聚合被仪表盘每分钟刷新一次时,这种代价会被放大成持续的集群压力。Transform 的思路是把聚合结果物化到另一个索引:查询从源索引搬到结果索引,响应从秒级降到毫秒级。本文讲清 pivot 与 latest 两种模式、连续转换的调度与 checkpoint、历史回填方式,以及它与 rollup 的分工。

1. Transform 解决的问题

1.1 从实时聚合到物化视图

关系型数据库早就有物化视图的概念:把昂贵的连接与聚合结果预先算好、存成表,查询时直接读结果。Elasticsearch 的聚合是「每次查询现算」的模型,没有内置物化。Transform 补上了这一环:它按定义周期性执行聚合,把结果写入目标索引,并提供查询入口。

1.2 Transform 的定位

Transform 适合三类场景:一是仪表盘高频刷新的固定聚合;二是需要长期保留的汇总指标(如按小时的活跃用户数);三是把「以事件为中心」的原始数据,转换成「以实体为中心」的宽表,便于后续检索。它不是实时管道,默认有分钟级延迟,需要毫秒级一致的场景应直接查源索引。

1.3 三种运行方式

方式说明适用
batch一次性执行,跑完即停历史数据一次性汇总
continuous持续运行,按 sync 配置增量仪表盘、长期指标
手工触发用 _start 接口显式启动回填、重跑

选择哪种取决于数据是否还在写入:只处理存量数据用 batch,处理增量用 continuous。

2. pivot 模式

2.1 基本定义

pivot 是 Transform 的主模式:按 group_by 分组,对每组算聚合,结果写一行。

curl -X PUT "localhost:9200/_transform/orders_by_day" \
  -H "Content-Type: application/json" -d'
{
  "source": { "index": "orders" },
  "dest":   { "index": "orders_by_day" },
  "pivot": {
    "group_by": {
      "day":   { "date_histogram": { "field": "created_at", "calendar_interval": "1d" } },
      "status": { "terms": { "field": "status" } }
    },
    "aggregations": {
      "total_amount": { "sum": { "field": "amount" } },
      "order_count":  { "value_count": { "field": "order_id" } },
      "avg_amount":   { "avg": { "field": "amount" } }
    }
  }
}'

执行后目标索引里每行是一个「日期 + 状态」组合,带三个指标字段,行数等于源数据的基数组合数,而不是源文档数。

2.2 group_by 支持的桶

group_by 支持 terms、histogram、date_histogram、range、geotile_grid、script 六种。常用的三种:

{
  "group_by": {
    "user":  { "terms": { "field": "user.id", "missing_bucket": true } },
    "hour":  { "date_histogram": { "field": "@timestamp", "fixed_interval": "1h" } },
    "price": { "histogram": { "field": "amount", "interval": 100 } }
  }
}

missing_bucket: true 让字段缺失的文档归到一个单独的桶,而不是被丢弃;这在排查「为什么总量对不上」时非常关键。

2.3 聚合与 scripted_metric 限制

pivot 的 aggregations 支持大部分 metric 聚合(sum、avg、min、max、value_count、cardinality、percentiles、scripted_metric),但不支持 bucket 聚合的嵌套下钻。要表达多层分组,把每一层都写进 group_by,而不是嵌套 aggs。cardinality 在 Transform 里用 HLL++ 近似算法,跨批次的结果会有小幅误差,做精确去重计数需要换方案。

2.4 目标索引的映射生成

不指定 dest.index 的映射时,Transform 会根据 group_by 与聚合字段自动推导:分组字段的类型从源索引复制,聚合结果统一为 long 或 double。自动推导省事,但有两点要注意:date_histogram 分组会生成 date 字段,时区按 UTC 存储,展示层需自行转换;cardinality 生成 long,数值可能因 HLL++ 近似而与精确值有偏差。若需要给结果字段加 doc_values: false、指定 format 或补 keyword 子字段,应先用 _preview 看推导结果,再手工建目标索引并显式给出映射。

3. latest 模式

3.1 语义

latest 模式不做聚合,而是按 unique_key 取每个实体的最新一条文档:

curl -X PUT "localhost:9200/_transform/user_latest" \
  -H "Content-Type: application/json" -d'
{
  "source": { "index": "user_events" },
  "dest":   { "index": "user_latest" },
  "latest": { "unique_key": ["user.id"], "sort": "@timestamp" }
}'

目标索引里每个 user.id 只有一行,是该用户时间戳最大的事件内容。

3.2 与 pivot 的差异

维度pivotlatest
输出语义分组统计实体最新状态
行数基数组合数唯一键数量
支持的聚合metric 聚合无
典型用途指标看板用户画像、设备状态
增量更新累加/重算分组覆盖同键旧行

latest 的关键约束是:目标索引里同一 unique_key 只有一行,新事件到达时旧行被替换(实际实现是 upsert 加删除旧文档)。

3.3 使用场景

latest 常用于三类需求:用户最近一次登录的设备与 IP、IoT 设备的最新上报值、订单的最新状态。它把「从事件流里捞最新一条」这个在 DSL 里要写 top_hits 或 collapse 的操作,变成了一次普通的 term 查询。

curl -X POST "localhost:9200/user_latest/_search?pretty" -H "Content-Type: application/json" -d'
{ "query": { "term": { "user.id": "u-1024" } } }'

4. 连续转换与调度

4.1 sync 配置

continuous 模式靠 sync 块驱动:

{
  "sync": {
    "time": {
      "field": "@timestamp",
      "delay": "60s"
    }
  }
}

field 是用于检测增量的时间字段,delay 是「等多久才处理」的延迟窗口,用来容忍乱序到达的事件。若源数据没有时间字段,可用 sync.time 之外的方式不可行——此时只能用 batch 模式手工重跑。

4.2 checkpoint 与增量边界

Transform 内部维护 checkpoint:记录上次处理到的 @timestamp 边界。每次调度时,引擎读取 checkpoint < @timestamp <= now - delay 的文档做聚合。理解 checkpoint 是排查「数据没更新」的关键:如果新写入文档的时间戳落在 delay 窗口内,本次不会处理,要等下一轮。

4.3 调度频率

{ "frequency": "1m" }

frequency 决定检查增量的频率,默认 1 分钟。设置过密会给源索引带来持续的小查询压力;过疏则数据新鲜度差。经验值:仪表盘刷新频率就是新鲜度的下限,frequency 取该值的一半到相等即可。

4.4 失败与重试

Transform 是「至少一次」语义:任务失败重启后从 checkpoint 重跑,可能导致同一时间窗被处理两次。对 pivot 的累加型指标,重复处理会造成重复计数吗?不会——pivot 每次对窗口内的源文档重新计算并 upsert 结果行,是幂等的;但 latest 模式在重跑时会用旧数据覆盖新数据,需注意源数据的时间戳单调性。

4.5 与 ILM 配合

目标索引同样会随时间增长。给它挂上 ILM 策略,让老化的汇总数据转冷或删除,是标准做法:

{
  "policy": "transform_results_30d",
  "rollover_alias": "orders_by_day"
}

注意 Transform 的 dest.index 若是别名,rollover 后写入会自动切到新索引,查询仍走别名,无需改 Transform 配置。

5. 回填与重跑

5.1 一次性回填历史

新建 Transform 时默认从「当前时间」开始,历史数据不会自动补。回填要显式给起点:

curl -X POST "localhost:9200/_transform/orders_by_day/_start?pretty" \
  -H "Content-Type: application/json" -d'
{
  "start": "2024-01-01T00:00:00Z"
}'

start 是 ISO 时间,引擎会从该时间点开始按 sync 的字段推进。

5.2 重建目标索引

改了 group_by 或 aggregations 之后,旧结果与新定义不兼容,必须重建:

curl -X POST "localhost:9200/_transform/orders_by_day/_stop"
curl -X DELETE "localhost:9200/orders_by_day"
curl -X PUT    "localhost:9200/_transform/orders_by_day" -d @new_def.json
curl -X POST   "localhost:9200/_transform/orders_by_day/_start?timeout=1m"

重建期间目标索引不可查,若上层是仪表盘,要先准备临时索引或切换别名,避免查询报 index_not_found_exception。

5.3 回填的性能

回填是重负载操作:它会全量扫描源索引并按时间窗分批聚合。建议在业务低峰执行,并限制资源:

{
  "settings": {
    "max_page_search_size": 500,
    "docs_per_second": 1000
  }
}

max_page_search_size 控制每批拉取的桶数,docs_per_second 限制吞吐。回填大索引时先把 docs_per_second 压到千级,跑通后再逐步放开,避免把源集群的查询线程池打满。

5.4 预览结果再落地

正式创建前,先用 _preview 看前 N 行结果,验证聚合口径:

curl -X POST "localhost:9200/_transform/_preview?pretty" \
  -H "Content-Type: application/json" -d'
{
  "source": { "index": "orders" },
  "pivot": {
    "group_by": { "status": { "terms": { "field": "status" } } },
    "aggregations": { "total": { "sum": { "field": "amount" } } }
  }
}'

_preview 不写目标索引、不消耗 checkpoint,是改定义时的安全试验台。它能暴露两类问题:一是字段名拼写错误导致的 unknown field;二是 terms 分组未设 size 时只返回前 10 个桶,让你误以为分组维度很少。

6. 与 rollup 的取舍

6.1 rollup 的定位

rollup 也是把聚合物化,但目标更窄:它专为时序指标降采样设计,只能按时间桶加少量维度分组,支持的是 metrics 里的数值聚合,且不能做非时间维度的灵活分组。rollup 的优势是写入路径更轻、压缩比更高。

6.2 能力对比

维度Transformrollup
分组维度任意多桶组合时间桶 + 有限维度
聚合类型metric 全集数值指标为主
增量语义checkpoint 推进按时间桶滚动
查询方式直接查目标索引走 _rollup_search
适用灵活汇总、实体宽表纯时序降采样
维护状态持续演进已进入维护模式

6.3 选型建议

新项目优先用 Transform:它的分组更灵活、查询就是普通 _search、社区方向也明确。只有一种情况值得考虑 rollup:数据是纯时序指标、降采样粒度固定、且需要极致的存储压缩。即便如此,Transform 配合 ILM 的 shrink 与 force_merge 也能拿到接近的效果。

7. 运维与监控

7.1 查看运行状态

curl -s "localhost:9200/_transform/orders_by_day/_stats?pretty"

返回 state(started/stopped/failed)、checkpoint(当前处理到的时间)、stats.documents_processed、stats.documents_indexed、stats.search_failures、stats.index_failures。持续观察 documents_processed 是否推进,是判断任务是否卡住的最直接方式。

7.2 常见故障

现象原因处理
目标索引不更新checkpoint 落后、delay 过大检查 sync.delay 与 frequency
任务变 failed目标索引被删、权限不足看 reason 字段,重建索引
结果行数暴涨分组字段基数高收窄 group_by,加 missing_bucket
内存压力大max_page_search_size 过大调小分页,降低并发
查询结果对不上源聚合有乱序事件落在 delay 外加大 delay 或改事件时间

7.3 权限与多租户

Transform 需要读取源索引、写入目标索引、操作自身状态的权限。开启安全后,要授予 manage_transform 集群权限与相应的索引权限,否则任务会以 security_exception 失败。多租户场景下,Transform 的目标索引要放在对应租户的索引命名空间里,避免跨租户数据泄漏。

8. 总结

环节要点
核心价值把聚合结果物化,查询从秒级降到毫秒级
pivot按 group_by 分组算指标,输出基数组合行
latest按唯一键取最新文档,输出实体宽表
连续转换sync.time + delay 检测增量,checkpoint 推进
回填_start 指定起点,重建索引需停任务
rollup 取舍灵活汇总选 Transform,纯时序降采样可考虑 rollup
监控看 _stats 的 checkpoint 与 documents_processed
幂等性pivot 重跑幂等,latest 重跑依赖时间戳单调

Transform 是把「实时聚合」转成「预计算查询」的最直接手段。它的成本是分钟级延迟与额外的存储,收益是稳定的查询延迟与可控的集群负载。落地时先用 batch 模式验证聚合定义,确认结果口径无误后再切 continuous,并把目标索引的 ILM 一并规划好。聚合语法细节可阅读《聚合分析与统计》,时序数据治理可阅读《时序数据与 TSDB》。

延伸阅读

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「elasticsearch」更多文章

  1. 组合模板与索引生命周期
  2. 批量写入调优与背压
  3. 相关性调优与离线评测