MongoDB Source Connector配置pipeline后无消息推送问题排查
核心问题原因
1. Pipeline匹配路径错误
MongoDB Source Connector的pipeline是作用在Change Stream事件结构上,而非直接作用在业务文档上。Change Stream的标准结构如下,业务字段全部嵌套在fullDocument节点下:
{ "operationType": "insert", "fullDocument": { "_id": "xxx", "jobStatus": 5, // 其余业务字段 }, // 其他Change Stream元字段 }
你原配置中的匹配规则直接匹配根节点的jobStatus字段,该字段在事件结构中不存在,自然过滤不到任何数据。
2. 可选配置缺失(更新场景必填)
如果要捕获update操作的最新完整文档,需要额外配置Change Stream的全量文档查询规则,否则update事件的fullDocument默认不会返回最新内容,也会导致匹配失败。
3. 配置JSON语法错误
你粘贴的配置中config对象在pipeline行后提前闭合,导致transforms相关配置不在config块内,若提交的实际请求也存在该语法问题,会导致配置不生效。
修正后完整配置
{ "name": "MongoSourceConn", "config": { "name": "MongoSourceConn", "connector.class": "com.mongodb.kafka.connect.MongoSourceConnector", "key.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "key.converter.schemas.enable": false, "value.converter.schemas.enable": false, "value.converter.schema.registry.url":"http://schema-registry:8081", "publish.full.document.only": true, "topics": "test_topic", "connection.uri": "mongodb://siteUserAdmin:rstatools@rsgadcmgo5:27017", "database": "kafka", "collection": "test_topic", "change.stream.full.document": "updateLookup", "pipeline": "[{ \"$match\": { \"$and\": [ {\"operationType\": { \"$in\": [ \"update\",\"insert\" ]}}, {\"fullDocument.jobStatus\": {\"$eq\": 5}} ] }} ]", "transforms":"dropPrefix", "transforms.dropPrefix.regex":"kafka.test_topic", "transforms.dropPrefix.type":"org.apache.kafka.connect.transforms.RegexRouter", "transforms.dropPrefix.replacement":"test_topic" } }
校验方法
如果修改后仍未生效,可临时关闭publish.full.document.only配置,将完整的Change Stream事件输出到Topic,确认事件结构和字段值符合预期后再调整pipeline规则即可。
内容的提问来源于stack exchange,提问作者jackyroo
相关产品推荐
相关产品推荐

