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

Elasticsearch:将一个转换作为另一个转换的数据源

Elasticsearch转换同步与摄取管道使用方案

一、实现品类转换与产品转换的同步

根据你的场景(单品类含大量产品,依赖产品转换结果做品类统计),可以通过以下两种方式实现同步:

  • 连续模式实时联动
    如果产品转换采用连续运行模式(通过"sync": {"time": {"field": "update_time"}}配置实现增量更新),直接将品类转换的数据源设为产品转换的目标索引,同时给品类转换配置相同的连续同步策略。当产品转换更新目标索引的产品收入数据时,品类转换会自动捕获这些增量变化,实时更新品类级统计结果。
    注意:确保产品转换的目标索引包含可用于同步的字段(如更新时间戳update_time),品类转换的同步配置需基于该字段跟踪数据变化。

  • 一次性转换的调度控制
    若仅需计算特定时间段的数据,需保证产品转换完全执行后再启动品类转换。可通过Elasticsearch API或外部调度工具实现:

    1. 调用POST _transform/<product-transform-id>/<product-transform-id>/_start启动产品转换,通过GET _transform/<product-transform-id>/_stats轮询状态,确认转换完成(status为stopped且documents_processed符合预期)。
    2. 产品转换完成后,再调用POST _transform/<category-transform-id>/_start启动品类转换。
  • 增量同步的字段优化
    建议在产品转换的目标索引中添加transformed_at字段,通过转换配置中的script处理器生成(ctx.transformed_at = new Date())。品类转换的同步策略基于该字段,确保只处理产品转换最新生成/更新的数据,避免重复计算。

二、转换中使用摄取管道添加摄取时间字段

Elasticsearch转换完全支持通过摄取管道添加摄取时间字段,具体操作如下:

  1. 创建摄取管道
    先创建包含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"
          }
        }
      ]
    }
    
  2. 在转换配置中关联管道
    创建或更新品类转换时,在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 11:17:25