Debezium Connect更新MongoDB后before/after字段为空问题求助
环境配置
- MongoDB:docker镜像
mongo:6 - Kafka:
confluentinc/cp-kafka:latest - Debezium Connect:
debezium/connect:2.4 - ZooKeeper:
confluentinc/cp-zookeeper:latest
需求与当前配置
需要捕获MongoDB的更新/删除操作,获取变更前后的完整文档状态用于日志记录。当前使用的Debezium Connector配置如下:
{ "name": "management-workspaces-deneme", "config": { "connector.class": "io.debezium.connector.mongodb.MongoDbConnector", "mongodb.connection.string": "mongodb://192.168.2.15:27030", "topic.prefix": "eys", "capture.mode":"change_streams_with_pre_image", "mongodb.poll.interval.ms": "5000", "database.history.kafka.bootstrap.servers": "172.23.0.3:29092", "heartbeat.interval.ms": "5000", "database.include.list":"wins", "collection.include.list": "wins.workspaces,wins.tasks" } }
根据Debezium官方文档:
The value of a change event for an update in the sample customers collection has the same schema as a create event for that collection. Likewise, the event value’s payload has the same structure. However, the event value payload contains different values in an update event. An update event includes an after value only if the capture.mode option is set to change_streams_update_full. A before value is provided if the capture.mode option is set to one of the *_with_pre_image option.
但当前消费到的变更消息中,before和after字段均为null:
"payload":{"before":null,"after":null,"updateDescription":{"removedFields":null,"updatedFields":"{\"description\": \"666teste\"}",
解决步骤
1. 启用MongoDB集合的预镜像功能
MongoDB 6.0+默认不会为集合开启预镜像支持,而Debezium的*_with_pre_image模式依赖此功能。针对目标集合执行以下命令开启:
// 开启wins.workspaces集合的预镜像 db.runCommand({ collMod: "workspaces", changeStreamPreAndPostImages: { enabled: true } }) // 开启wins.tasks集合的预镜像 db.runCommand({ collMod: "tasks", changeStreamPreAndPostImages: { enabled: true } })
注意:该操作需要对应数据库的collMod权限,且需逐个集合配置,无法批量操作。
2. 修改Debezium的capture.mode配置
要同时获取before(变更前完整文档)和after(变更后完整文档),需将capture.mode调整为change_streams_update_full_with_pre_image。原配置的change_streams_with_pre_image仅提供before字段,不会自动填充after;而change_streams_update_full_with_pre_image是同时支持前后状态的模式。
3. 重启Debezium Connector
修改配置后,需重启Connector使变更生效。可通过Kafka Connect的REST API操作:
# 暂停Connector curl -X PUT http://<connect-host>:<port>/connectors/management-workspaces-deneme/pause # 更新Connector配置(替换为你的Connect服务地址) curl -X PUT http://<connect-host>:<port>/connectors/management-workspaces-deneme/config \ -H "Content-Type: application/json" \ -d '{ "connector.class": "io.debezium.connector.mongodb.MongoDbConnector", "mongodb.connection.string": "mongodb://192.168.2.15:27030", "topic.prefix": "eys", "capture.mode":"change_streams_update_full_with_pre_image", "mongodb.poll.interval.ms": "5000", "database.history.kafka.bootstrap.servers": "172.23.0.3:29092", "heartbeat.interval.ms": "5000", "database.include.list":"wins", "collection.include.list": "wins.workspaces,wins.tasks" }' # 重启Connector curl -X PUT http://<connect-host>:<port>/connectors/management-workspaces-deneme/resume
4. 验证结果
对目标集合执行更新或删除操作后,重新消费Kafka主题,检查消息的payload字段是否已包含非null的before和after值。
内容的提问来源于stack exchange,提问作者Furkan YIlmaZ

