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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 20:05:30