使用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
相关产品推荐
相关产品推荐

