Kafka到PostgreSQL的JDBC Sink报错Value schema must be of type Struct求助
问题:Kafka JDBC Sink Connector写入PostgreSQL失败,报错“Value schema must be of type Struct”
环境与现状
- ZooKeeper、Kafka、Schema Registry、Kafka Connect均运行正常,Mongo Sink Connector可正常工作
- Kafka主题:
data-test-topic - 主题数据示例:
Key: "A" { "name": "abhishek", "lastname": "Rathore" }
- PostgreSQL表
init结构:
Column | Type | Collation | Nullable | Default | Storage | Compression | Stats target | Description ----------+------+-----------+----------+---------+----------+-------------+--------------+------------- name | text | | not null | | extended | | | lastname | text | | not null | | extended | | | key | text | | not null | | extended | | | Access method: heap
JDBC Sink Connector配置
{ "name": "postgres-jdbc-sink-connector", "config": { "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector", "tasks.max": 1, "topics": "data-test-topic", "connection.url": "jdbc:postgresql://postgresdb:5432/garuna-rgs", "connection.user": "postgres", "connection.password": "arathore", "auto.create": true, "auto.evolve": true, "insert.mode": "upsert", "table.name.format": "init", "pk.mode": "record_key", "pk.fields": "name", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "key.converter": "org.apache.kafka.connect.storage.StringConverter", "value.converter.schemas.enable": false, "key.converter.schemas.enable": false } }
错误日志
org.apache.kafka.connect.errors.ConnectException: Exiting WorkerSinkTask due to unrecoverable exception. at org.apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:618) at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:336) at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:237) at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:206) at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:202) at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:257) at org.apache.kafka.connect.runtime.isolation.Plugins.lambda$withClassLoader$1(Plugins.java:181) at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515) at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) at java.base/java.lang.Thread.run(Thread.java:829) Caused by: org.apache.kafka.connect.errors.ConnectException: Value schema must be of type Struct at io.confluent.connect.jdbc.sink.metadata.FieldsMetadata.extract(FieldsMetadata.java:86) at io.confluent.connect.jdbc.sink.metadata.FieldsMetadata.extract(FieldsMetadata.java:67) at io.confluent.connect.jdbc.sink.BufferedRecords.add(BufferedRecords.java:115) at io.confluent.connect.jdbc.sink.JdbcDbWriter.write(JdbcDbWriter.java:74) at io.confluent.connect.jdbc.sink.JdbcSinkTask.put(JdbcSinkTask.java:90) at org.apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:587) ... 11 more
解决方案
错误核心原因:JDBC Sink Connector需要结构化的Schema信息映射数据库字段,但当前配置value.converter.schemas.enable=false,导致Kafka Connect无法将无Schema的JSON解析为Struct类型,触发报错。
以下是两种可行解决方法:
方法1:使用带Schema的JSON转换器(推荐)
改用Confluent的JsonSchemaConverter,结合Schema Registry管理数据Schema,让连接器正确识别结构化数据:
修改连接器配置中value.converter相关参数:
{ "name": "postgres-jdbc-sink-connector", "config": { // 保留原有其他配置 "value.converter": "io.confluent.connect.json.JsonSchemaConverter", "value.converter.schema.registry.url": "http://schema-registry:8081", // 替换为你的Schema Registry地址 "key.converter": "org.apache.kafka.connect.storage.StringConverter", "key.converter.schemas.enable": false } }
注意:若主题数据未关联Schema,需先向Schema Registry注册对应数据Schema,或让生产者发送数据时附带Schema信息。
方法2:使用Kafka Connect Transform处理无Schema JSON
若不想使用Schema Registry,可通过Transform功能将无Schema JSON包装成Struct,让JDBC Sink识别:
修改连接器配置,添加Transform参数:
{ "name": "postgres-jdbc-sink-connector", "config": { // 保留原有其他配置 "value.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter.schemas.enable": false, "key.converter": "org.apache.kafka.connect.storage.StringConverter", "key.converter.schemas.enable": false, // 添加Transform配置 "transforms": "HoistField", "transforms.HoistField.type": "org.apache.kafka.connect.transforms.HoistField$Value", "transforms.HoistField.field": "payload" } }
若需将Kafka记录的Key写入表中key字段,可额外添加InsertField Transform:
"transforms": "HoistField,InsertKey", "transforms.HoistField.type": "org.apache.kafka.connect.transforms.HoistField$Value", "transforms.HoistField.field": "payload", "transforms.InsertKey.type": "org.apache.kafka.connect.transforms.InsertField$Key", "transforms.InsertKey.field": "key"
额外配置调整
当前pk.mode设为record_key,但pk.fields指定name存在矛盾:record_key模式下主键字段对应Kafka记录的Key,而你的Key是字符串"A",表中name字段属于数据Value内容。若要以Value中的name作为主键,需将pk.mode改为record_value:
"pk.mode": "record_value", "pk.fields": "name"
内容的提问来源于stack exchange,提问作者Rathore
相关产品推荐
相关产品推荐

