如何用Kafka Connect JDBC Sink导入纯JSON数据至数据库(不可修改生产者配置)
可以通过Kafka Connect实现需求,只需修正Connector配置即可
错误原因分析
你遇到的反序列化/未知魔术字节错误,核心问题是Connector的消息值转换器配置缺失或错误:
- 当前配置仅指定了
key.converter,未设置value.converter,Kafka Connect会默认使用io.confluent.connect.avro.AvroConverter,而该转换器期望消息是带特定魔术字节的Avro序列化数据,但你的消息是纯文本JSON,因此触发错误。 - 额外添加的
key.converter.schema.registry.url完全多余,StringConverter不需要Schema Registry支持。
修正后的Connector配置
以下是适配纯文本JSON消息的正确配置:
{ "name": "JdbcSinkConnector", "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector", "table.name.format": "people", "topics": "people", "tasks.max": "1", "transforms": "RenameField", "transforms.RenameField.renames": "town:city,mobile_number:msisdn", "transforms.RenameField.type": "org.apache.kafka.connect.transforms.ReplaceField$Value", "connection.url": "jdbc:postgresql://localhost:5432/postgres", "connection.password": "postgres", "connection.user": "postgres", "insert.mode": "insert", "key.converter": "org.apache.kafka.connect.storage.StringConverter", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter.schemas.enable": "false" }
关键配置说明
value.converter设为JsonConverter:专门处理JSON格式的消息,无需提前序列化value.converter.schemas.enable=false:因为你的消息是无Schema的纯文本JSON,不需要启用Schema校验- 移除了多余的
key.converter.schema.registry.url配置 - 修正了原配置中
transforms.RenameField.renames的多余引号,避免字段重命名失效
验证步骤
- 停止原Connector并删除配置
- 提交上述修正后的配置到Kafka Connect集群
- 用
kafka-console-producer发送测试JSON数据 - 检查PostgreSQL的
people表是否成功插入数据
内容的提问来源于stack exchange,提问作者sina mohsenifar
相关产品推荐
相关产品推荐

