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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 00:37:34