Elasticsearch:将一个转换作为另一个转换的数据源
一、实现品类转换与产品转换的同步
根据你的场景(单品类含大量产品,依赖产品转换结果做品类统计),可以通过以下两种方式实现同步:
连续模式实时联动
如果产品转换采用连续运行模式(通过"sync": {"time": {"field": "update_time"}}配置实现增量更新),直接将品类转换的数据源设为产品转换的目标索引,同时给品类转换配置相同的连续同步策略。当产品转换更新目标索引的产品收入数据时,品类转换会自动捕获这些增量变化,实时更新品类级统计结果。
注意:确保产品转换的目标索引包含可用于同步的字段(如更新时间戳update_time),品类转换的同步配置需基于该字段跟踪数据变化。一次性转换的调度控制
若仅需计算特定时间段的数据,需保证产品转换完全执行后再启动品类转换。可通过Elasticsearch API或外部调度工具实现:- 调用
POST _transform/<product-transform-id>/<product-transform-id>/_start启动产品转换,通过GET _transform/<product-transform-id>/_stats轮询状态,确认转换完成(status为stopped且documents_processed符合预期)。 - 产品转换完成后,再调用
POST _transform/<category-transform-id>/_start启动品类转换。
- 调用
增量同步的字段优化
建议在产品转换的目标索引中添加transformed_at字段,通过转换配置中的script处理器生成(ctx.transformed_at = new Date())。品类转换的同步策略基于该字段,确保只处理产品转换最新生成/更新的数据,避免重复计算。
二、转换中使用摄取管道添加摄取时间字段
Elasticsearch转换完全支持通过摄取管道添加摄取时间字段,具体操作如下:
创建摄取管道
先创建包含date处理器的管道,用于生成摄取时间字段(示例为自定义ingest_time字段):PUT _ingest/pipeline/add-ingest-time { "processors": [ { "date": { "field": "_ingest.timestamp", "target_field": "ingest_time", "format": "yyyy-MM-dd'T'HH:mm:ss.SSSZ" } } ] }在转换配置中关联管道
创建或更新品类转换时,在dest配置中添加pipeline参数指向上述管道:PUT _transform/category-revenue-transform { "source": { "index": "product-revenue-transform-result" // 产品转换的目标索引 }, "pivot": { "group_by": { "category_id": { "terms": { "field": "category_id" } } }, "aggregations": { "total_revenue": { "sum": { "field": "revenue" } } } }, "dest": { "index": "category-revenue-result", "pipeline": "add-ingest-time" // 指定摄取管道 }, "sync": { "time": { "field": "transformed_at" } } }转换向目标索引写入文档时,会自动通过摄取管道添加
ingest_time字段,记录数据写入时间。
内容的提问来源于stack exchange,提问作者Mathijs

