Debezium MongoDB源数据无法同步至MySQL Sink问题求助
问题分析
错误Unsupported source data type: STRUCT的根源是MongoDB源文档的_id是嵌套结构体({"id":6}),经过UnwrapFromMongoDbEnvelope转换后,该字段仍以STRUCT类型保留在record中。而目标MySQL表的id是INT类型主键,JDBC Sink连接器无法将STRUCT类型映射到INT字段,导致SQL参数绑定失败。
解决方案
1. 修改Sink连接器配置,添加字段转换
在Sink配置中新增字段提升、重命名与删除的Transforms,将嵌套在_id中的id值提取为顶级INT类型字段,匹配MySQL表的主键要求:
{ "name": "JdbcSinkConnector", "config": { "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector", "tasks.max": "1", "key.converter": "io.confluent.connect.avro.AvroConverter", "value.converter": "io.confluent.connect.avro.AvroConverter", "transforms": "unwrap,hoistId,renameId,dropOldId", "topics": "dbserver2.uts_production_kafkaconnect.booking_journal_production", "transforms.unwrap.type": "io.debezium.connector.mongodb.transforms.UnwrapFromMongoDbEnvelope", // 将_id.id嵌套字段提升为顶级字段 "transforms.hoistId.type": "org.apache.kafka.connect.transforms.HoistField$Value", "transforms.hoistId.field": "_id.id", // 把提升后的字段重命名为id "transforms.renameId.type": "org.apache.kafka.connect.transforms.ReplaceField$Value", "transforms.renameId.renames": "_id.id:id", // 删除原有的_id结构体字段 "transforms.dropOldId.type": "org.apache.kafka.connect.transforms.ReplaceField$Value", "transforms.dropOldId.blacklist": "_id", "connection.url": "jdbc:mysql://localhost:3306/uts_thinclient_kafkaconnect", "connection.user": "root", "connection.password": "password", "dialect.name": "MySqlDatabaseDialect", "insert.mode": "upsert", "table.name.format": "uts_thinclient_kafkaconnect.booking_journal_thinclient", "pk.mode": "record_value", "pk.fields": "id", "db.timezone": "Asia/Kolkata", "auto.create": "false", // 目标表已存在,关闭自动创建避免冲突 "auto.evolve": "true", "value.converter.schema.registry.url": "http://localhost:8081", "key.converter.schema.registry.url": "http://localhost:8081" } }
2. 可选:在源端提前处理字段
如果希望减少Sink端的转换逻辑,可以在MongoDB源连接器中添加相同的Transforms,直接输出扁平化的id字段:
{ "name": "MongoDbSourceConnector", "config": { "connector.class": "io.debezium.connector.mongodb.MongoDbConnector", "tasks.max": "1", "key.converter": "io.confluent.connect.avro.AvroConverter", "value.converter": "io.confluent.connect.avro.AvroConverter", "mongodb.hosts": "debezium/localhost:27017", "mongodb.user": "sudhir", "mongodb.password": "password", "mongodb.name": "dbserver2", "value.converter.schema.registry.url": "http://localhost:8081", "key.converter.schema.registry.url": "http://localhost:8081", "transforms": "hoistId,renameId,dropOldId", "transforms.hoistId.type": "org.apache.kafka.connect.transforms.HoistField$Value", "transforms.hoistId.field": "_id.id", "transforms.renameId.type": "org.apache.kafka.connect.transforms.ReplaceField$Value", "transforms.renameId.renames": "_id.id:id", "transforms.dropOldId.type": "org.apache.kafka.connect.transforms.ReplaceField$Value", "transforms.dropOldId.blacklist": "_id" } }
3. 日期字段兼容性说明
MongoDB的$date类型经过Debezium转换后会处理为标准Timestamp类型,结合Sink配置中的db.timezone: Asia/Kolkata,可直接映射到MySQL的datetime(6)字段,无需额外配置。
内容的提问来源于stack exchange,提问作者Sudhir Daga
相关产品推荐
相关产品推荐

