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

如何在OpenSearch Data Prepper中按可选键聚合CDR与QoE数据?

用OpenSearch Data Prepper实现多键关联的CDR+QoE事件聚合

不用额外引入Kafka,Data Prepper完全能搞定你的需求,核心是先给所有事件统一一个聚合标识,再基于这个标识完成三方合并。下面是具体实现步骤:

1. 前置处理:给所有事件打统一聚合键

在聚合之前,用mutate处理器给每个事件新增一个aggregation_key字段,不管是CDR还是QoE事件,都把对应的通话ID赋值给这个字段:

  • CDR事件直接用自身的唯一通话ID(假设字段是callId)作为aggregation_key
  • QoE事件判断是带ingressCallId还是egressCallId,把存在的那个ID赋值给aggregation_key

配置片段如下:

processor:
  - mutate:
      entries:
        # 处理CDR事件,用自身callId做聚合键
        - key: aggregation_key
          value: "{{callId}}"
          condition: '{{callId exists}}'
        # QoE事件优先用ingressCallId,没有就用egressCallId
        - key: aggregation_key
          value: "{{ingressCallId}}"
          condition: '{{ingressCallId exists}}'
        - key: aggregation_key
          value: "{{egressCallId}}"
          condition: '{{egressCallId exists}}'

2. 配置聚合处理器,按统一键合并事件

基于刚才生成的aggregation_key做聚合,设置超时时间(比如15秒,覆盖SBC的10秒间隔),同时设置凑齐3个事件就立即聚合:

processor:
  - aggregate:
      identification_keys: ["aggregation_key"]
      action:
        merge:
          merge_strategy: merge_entries
      timeout: 15s
      max_size: 3
  • identification_keys指定用统一后的aggregation_key做匹配
  • merge_entries会把同键的三个事件字段合并成一个JSON对象(如果字段冲突,后续事件的字段会覆盖前面的,你也可以根据需求换成keep_first或keep_last)
  • timeout设15秒,确保能等齐三个事件;max_size设3,凑够数就立即输出,不用等超时

3. 输出到OpenSearch

聚合完成后,直接用opensearch sink把合并好的对象写入目标索引就行。

额外提示

  • 如果担心有事件丢失(比如某路通话只收到2个事件),可以在聚合处理器里加discard_on_timeout: false,让超时后的部分聚合结果也能输出,避免丢数据。
  • 要是需要调整字段结构或者处理字段冲突,聚合之后再加个mutate处理器就行。

内容的提问来源于stack exchange,提问作者Damo

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 00:22:43