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

使用Kafka JDBC Sink Connector加载纯JSON数据到PostgreSQL报错

问题解答

理解正误判断

你的理解部分正确:

  • 若需要将消息的key作为PostgreSQL表的主键,确实需要将pk.mode设置为record_key,同时要求消息key为包含所有主键字段的Struct类型
  • delete.enabled:true不是必选配置,只有当你需要同步Kafka中的墓碑消息(即value为null的记录)来触发数据库行删除时才需要开启,纯写入场景不需要开启

核心报错原因

你遇到的报错本质是JDBC Sink连接器要求记录必须携带非空的Schema和Struct类型值,你关闭value.converter.schemas.enable后,Schema推断功能未正常生效,导致连接器拿到的是HashMap类型的值和空Schema,触发校验失败。

可选解决方案

方案1:无需修改消息结构,使用value字段作为主键(更简便)

要求Kafka版本 >= 2.6 / Confluent Platform版本 >= 6.0 (支持KIP-301的Schema推断能力),调整连接器配置如下:

curl -X PUT http://localhost:8083/connectors/load_test/config \
-H "Content-Type: application/json" \
-d '{
 "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector",
 "connection.url":"jdbc:postgresql://localhost:5432/somedb",
 "connection.user":"user",
 "connection.password":"passwd",
 "value.converter":"org.apache.kafka.connect.json.JsonConverter",
 "value.converter.schemas.enable":"false",
 "value.converter.schemas.infer.enable": "true",
 "value.converter.schemas.infer.union.enable": "true",
 "tasks.max" : "1",
 "topics":"dup_emp",
 "table.name.format":"dup_emp",
 "insert.mode":"insert",
 "quote.sql.identifiers":"never",
 "pk.mode": "record_value",
 "pk.fields": "id",
 "auto.create": "true",
 "auto.evolve": "true"
}'

注:pk.fields替换为你表中实际的主键字段名,多主键用逗号分隔;auto.create和auto.evolve为可选配置,不需要自动维护表结构可关闭。

方案2:按你的预期使用消息key作为主键

第一步:生产带Struct类型key的消息

以Java生产者为例,示例代码如下:

import org.apache.kafka.common.schema.Schema;
import org.apache.kafka.common.schema.SchemaBuilder;
import org.apache.kafka.common.struct.Struct;
import com.google.gson.JsonObject;
import org.apache.kafka.clients.producer.ProducerRecord;

// 1. 定义key的Schema
Schema empKeySchema = SchemaBuilder.struct()
    .name("com.example.EmployeeKey")
    .field("emp_id", Schema.INT32_SCHEMA) // 替换为你的主键字段
    .build();

// 2. 构建Struct类型的key
Struct recordKey = new Struct(empKeySchema)
    .put("emp_id", 1001);

// 3. 构建消息value
JsonObject recordValue = new JsonObject();
recordValue.addProperty("emp_id", 1001);
recordValue.addProperty("emp_name", "张三");
recordValue.addProperty("dept", "研发部");

// 4. 发送消息
ProducerRecord<Struct, JsonObject> record = new ProducerRecord<>("dup_emp", recordKey, recordValue);
producer.send(record);

第二步:调整连接器配置

curl -X PUT http://localhost:8083/connectors/load_test/config \
-H "Content-Type: application/json" \
-d '{
 "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector",
 "connection.url":"jdbc:postgresql://localhost:5432/somedb",
 "connection.user":"user",
 "connection.password":"passwd",
 "key.converter":"org.apache.kafka.connect.json.JsonConverter",
 "key.converter.schemas.enable":"true",
 "value.converter":"org.apache.kafka.connect.json.JsonConverter",
 "value.converter.schemas.enable":"false",
 "value.converter.schemas.infer.enable": "true",
 "tasks.max" : "1",
 "topics":"dup_emp",
 "table.name.format":"dup_emp",
 "insert.mode":"insert",
 "quote.sql.identifiers":"never",
 "pk.mode": "record_key",
 "pk.fields": "emp_id"
}'

注:pk.fields替换为key Struct中定义的主键字段名。

排查注意事项

如果配置后仍然报错,需要确认:

  • 连接器配置没有被Kafka Connect服务端的全局Converter配置覆盖,所有Converter相关参数都明确写在连接器配置中
  • 消息value的结构统一,没有出现同字段类型不一致的情况,若有可开启value.converter.schemas.infer.union.enable兼容

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 07:57:02