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

如何通过MongoDB Kafka Sink Connector实现消息过滤与字段转换?

使用MongoDB Kafka Sink Connector实现消息过滤与字段选择

MongoDB Kafka Sink Connector完全支持通过**Single Message Transformations (SMT)**实现你需要的消息过滤和字段选择需求,无需自行开发额外组件。

1. 实现消息条件过滤

使用org.apache.kafka.connect.transforms.Filter这个内置SMT,就能基于消息内容的条件过滤掉不符合要求的消息。针对你需要保留colour=green消息的场景,配置示例如下:

transforms=FilterGreen
transforms.FilterGreen.type=org.apache.kafka.connect.transforms.Filter
transforms.FilterGreen.condition=value.colour == 'green'
transforms.FilterGreen.filter.negate=false
  • condition:设置判断规则,这里通过value.colour == 'green'匹配目标字段值
  • filter.negate=false:表示保留满足条件的消息,若设为true则会过滤掉符合条件的消息

2. 仅插入指定字段

可以通过两种SMT实现字段筛选,按需选择即可:

方式一:提取指定字段(推荐,明确保留所需字段)

比如你只想保留id、name、colour三个字段,配置如下:

transforms=FilterGreen,ExtractFields
transforms.FilterGreen.type=org.apache.kafka.connect.transforms.Filter
transforms.FilterGreen.condition=value.colour == 'green'
transforms.FilterGreen.filter.negate=false
transforms.ExtractFields.type=org.apache.kafka.connect.transforms.ExtractField$Value
transforms.ExtractFields.fields=id,name,colour

方式二:屏蔽不需要的字段

如果需要排除的字段较少,可使用此方式,比如屏蔽description、price字段:

transforms=FilterGreen,MaskFields
transforms.FilterGreen.type=org.apache.kafka.connect.transforms.Filter
transforms.FilterGreen.condition=value.colour == 'green'
transforms.FilterGreen.filter.negate=false
transforms.MaskFields.type=org.apache.kafka.connect.transforms.MaskField$Value
transforms.MaskFields.fields=description,price

关键注意点

  • 确保消息为JSON等结构化格式,SMT才能正确解析字段
  • SMT的执行顺序很重要:先过滤消息,再处理字段,避免对无用消息做多余操作
  • 若消息是嵌套结构,条件表达式需对应路径,比如value.details.colour == 'green'

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 16:54:51