聚合查询的代价随数据量线性增长:对十亿级文档做一次 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 的差异
| 维度 | pivot | latest |
|---|---|---|
| 输出语义 | 分组统计 | 实体最新状态 |
| 行数 | 基数组合数 | 唯一键数量 |
| 支持的聚合 | 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 能力对比
| 维度 | Transform | rollup |
|---|---|---|
| 分组维度 | 任意多桶组合 | 时间桶 + 有限维度 |
| 聚合类型 | 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》。
延伸阅读
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。