使用Kafka JDBC Sink Connector将JSON数据导入PostgreSQL问题咨询
报错原因解析
- 你当前配置中将
value.converter.schemas.enable设为了false,Kafka Connect解析你手动发送的纯JSON消息时,只会将内容转换为HashMap对象,不会生成对应的结构化Schema。而JDBC Sink连接器默认要求消息必须为带Schema的Struct类型,才能完成消息字段到数据库表列的映射,这是触发报错的核心原因。 - 目标表
dup_emp配置了主键emp_id,但你的连接器配置中未设置pk.mode参数,连接器无法识别主键映射规则,也是报错的触发条件之一。
可行配置方案
方案1:继续使用无Schema纯JSON消息(无需调整消息发送格式,推荐使用)
直接修改连接器配置,新增3个核心参数即可,修改后的完整配置如下:
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", "value.converter":"org.apache.kafka.connect.json.JsonConverter", "key.converter.schemas.enable":"false", "value.converter.schemas.enable":"false", "tasks.max" : "1", "topics":"dup_emp", "table.name.format":"dup_emp", "insert.mode":"insert", "quote.sql.identifiers":"never", "schema.ignore":"true", "pk.mode":"record_value", "pk.fields":"emp_id" }'
新增参数说明:
schema.ignore:设为true后,连接器会跳过Struct Schema校验逻辑,直接使用消息字段名和数据库表列名进行匹配映射pk.mode:设为record_value表示主键从消息的value部分取值pk.fields:指定主键对应的消息字段名,和表主键emp_id对齐
如果需要实现主键重复时自动更新数据,可以将insert.mode修改为upsert。
方案2:使用带Schema的JSON消息(无需调整核心转换器配置,需修改消息格式)
如果不希望修改连接器的原有逻辑,可以调整你手动发送的消息结构,增加schema定义部分,示例消息格式如下:
{ "schema": { "type": "struct", "fields": [ { "type": "int32", "optional": false, "field": "emp_id" }, { "type": "string", "optional": true, "field": "emp_name" }, { "type": "int32", "optional": true, "field": "emp_salary" } ], "optional": false, "name": "dup_emp" }, "payload": { "emp_id": 1, "emp_name": "bheem", "emp_salary": 2000 } }
该格式下连接器可以直接从消息的schema字段获取结构化定义,同时仍需要在配置中补充pk.mode和pk.fields两个参数,即可正常完成数据写入。
内容的提问来源于stack exchange,提问作者vigneshwar
相关产品推荐
相关产品推荐

