如何无需自定义Java类,配置Kafka Source Connector过滤IBMMQ字符串消息?
基于Confluent Filter Transform过滤字符串消息的配置方案
完全不需要编写自定义Transform类,用Confluent官方的Filter Transform结合内置表达式语言就能实现需求。
核心配置思路
利用Connect表达式语言直接对字符串类型的消息体(payload)进行判断,保留包含ABC123的消息,丢弃不符合条件的消息。
具体配置示例
如果IBM MQ Source Connector输出的消息值已经是字符串类型,直接添加以下Transform配置:
# 启用Filter Transform transforms=filter # 指定Filter Transform的类型 transforms.filter.type=io.confluent.connect.transforms.Filter # 设置过滤条件:仅保留包含ABC123的消息 transforms.filter.condition=$value contains 'ABC123' # 可选:如果需要明确排除含PQR123的消息(避免两种标识同时出现的极端情况),条件可改为 # transforms.filter.condition=$value contains 'ABC123' and not ($value contains 'PQR123')
注意事项
如果IBM MQ Source Connector默认输出的是字节数组类型的消息值,需要先添加Convert Transform将其转为字符串,再执行过滤:
# 先转字符串,再过滤 transforms=convert,filter # 配置Convert Transform:将消息值转为字符串 transforms.convert.type=org.apache.kafka.connect.transforms.Convert$Value transforms.convert.target.type=string # 配置Filter Transform transforms.filter.type=io.confluent.connect.transforms.Filter transforms.filter.condition=$value contains 'ABC123'
内容的提问来源于stack exchange,提问作者bluelabel
相关产品推荐
相关产品推荐

