如何在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
相关产品推荐
相关产品推荐

