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

如何用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的多余引号,避免字段重命名失效

验证步骤

  1. 停止原Connector并删除配置
  2. 提交上述修正后的配置到Kafka Connect集群
  3. 用kafka-console-producer发送测试JSON数据
  4. 检查PostgreSQL的people表是否成功插入数据

内容的提问来源于stack exchange,提问作者sina mohsenifar

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 05:55:16