Kafka到MySQL Sink Connector配置Transform报错:STRUCT类型无对应SQL列类型
Kafka到MySQL Sink Connector 报错解决指南
问题根源
你的配置里有两个致命问题,直接导致STRUCT类型无法映射到MySQL列:
- 重复定义
transforms字段:先写了transforms: "dropPrefix,unwrap",后面又写transforms: "ReplaceField",后者会完全覆盖前者,导致unwrap和时间转换的Transform根本没生效。Avro消息默认是嵌套的STRUCT结构,unwrap的作用就是把嵌套内容展开成扁平字段,没它的话,Connect会把整个STRUCT当成一个字段传给MySQL,而MySQL不支持这种类型,自然报错。 unwrap只有名称无配置:就算没被覆盖,你只把unwrap列在transforms列表里,没指定它的类型参数,这个Transform根本无法运行。
修正后的完整配置
把所有需要的Transform合并到一个transforms字段里,补上unwrap的配置,同时保留原有功能:
{ "name": "mysql-conf-sink", "config": { "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector", "tasks.max": "3", "value.converter": "io.confluent.connect.avro.AvroConverter", "value.converter.schema.registry.url": "http://localhost:8081", "topics": "mysql.cars.prices", "transforms": "dropPrefix,unwrap,timestamp,ReplaceField", "transforms.dropPrefix.type": "org.apache.kafka.connect.transforms.RegexRouter", "transforms.dropPrefix.regex": "mysql.cars.prices", "transforms.dropPrefix.replacement": "prices", "transforms.unwrap.type": "io.confluent.connect.transforms.ExtractField$Value", "transforms.timestamp.type": "org.apache.kafka.connect.transforms.TimestampConverter$Value", "transforms.timestamp.target.type": "Timestamp", "transforms.timestamp.field": "date_time", "transforms.timestamp.format": "yyyy-MM-dd HH:mm:ss", "transforms.ReplaceField.type": "org.apache.kafka.connect.transforms.ReplaceField$Value", "transforms.ReplaceField.blacklist": "id, brand", "errors.tolerance": "all", "errors.log.enable": "true", "errors.log.include.messages": "true", "connection.url": "jdbc:mysql://localhost:3306/product", "connection.user": "kafka", "connection.password": "123456", "insert.mode": "insert", "auto.create": "true", "auto.evolve": "true", "batch.size": 50000 } }
关键调整说明
- 合并Transform列表:所有要启用的Transform用逗号分隔,放在同一个
transforms配置项中,避免被后续配置覆盖。 - 补全
unwrap配置:使用ExtractField$Value类型的unwrap,将Avro消息中嵌套的value字段提取出来,转换成扁平的字段结构,这样Connect就能正确识别每个字段的类型,映射到MySQL的对应列。 - 保留原有功能:
ReplaceField拉黑id和brand、时间格式转换、Topic前缀替换这些功能都正常保留,不影响业务逻辑。
验证操作步骤
- 删除旧的错误连接器配置:
curl -X DELETE http://localhost:8083/connectors/mysql-conf-sink
- 提交修正后的配置(假设配置文件名为
mysql-sink-config.json):
curl -X POST -H "Content-Type: application/json" --data @mysql-sink-config.json http://localhost:8083/connectors
- 检查连接器运行状态,确认无报错:
curl http://localhost:8083/connectors/mysql-conf-sink/status
内容的提问来源于stack exchange,提问作者Anahita barzegar
相关产品推荐
相关产品推荐

