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

使用Debezium JDBC Sink Connector无法同步Kafka Topic至PostgreSQL

问题:Debezium JDBC Sink Connector同步无Schema JSON到PostgreSQL失败

问题背景

应用向Kafka Topic写入的消息为纯JSON(无Schema结构):

{
    "sub_id": "574",
    "sub_name": "john"
}

尝试使用以下Debezium JDBC Sink Connector配置将数据写入PostgreSQL:

config:
    class: io.debezium.connector.jdbc.JdbcSinkConnector

    connection.url: jdbc:postgresql://10.10.10.10:26257/db_dev
    connection.username: "******"
    connection.password: "********"
    topics: "kafka-crdb" 
    insert.mode: "upsert"
    primary.key.mode: "none"
    primary.key.fields: "sub_id"
    value.converter: "org.apache.kafka.connect.json.JsonConverter"
    value.converter.schemas.enable: "false"
    transforms: "unwrap"
    transforms.unwrap.type: "io.debezium.transforms.ExtractNewRecordState"

第一次执行错误

运行后抛出空指针异常:

2024-09-27 09:27:03,379 ERROR [kafka-sink-connector|task-0] Failed to process record: Failed to process a sink record (io.debezium.connector.jdbc.JdbcSinkConnectorTask) [task-thread-kafka-sink-connector-crdb-0]
org.apache.kafka.connect.errors.ConnectException: Failed to process a sink record
        at io.debezium.connector.jdbc.JdbcChangeEventSink.buildRecordSinkDescriptor(JdbcChangeEventSink.java:200)
        at io.debezium.connector.jdbc.JdbcChangeEventSink.execute(JdbcChangeEventSink.java:85)
        at io.debezium.connector.jdbc.JdbcSinkConnectorTask.put(JdbcSinkConnectorTask.java:103)
        at org.apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:601)
        at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:350)
        at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:250)
        at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:219)
        at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:204)
        at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:259)
        at org.apache.kafka.connect.runtime.isolation.Plugins.lambda$withClassLoader$1(Plugins.java:237)
        at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:539)
        at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264)
        at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136)
        at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635)
        at java.base/java.lang.Thread.run(Thread.java:840)
Caused by: java.lang.NullPointerException: Cannot invoke "org.apache.kafka.connect.data.Schema.name()" because the return value of "org.apache.kafka.connect.sink.SinkRecord.valueSchema()" is null
        at io.debezium.connector.jdbc.SinkRecordDescriptor$Builder.isFlattened(SinkRecordDescriptor.java:321)
        at io.debezium.connector.jdbc.SinkRecordDescriptor$Builder.build(SinkRecordDescriptor.java:310)
        at io.debezium.connector.jdbc.JdbcChangeEventSink.buildRecordSinkDescriptor(JdbcChangeEventSink.java:197)
        ... 14 more

修改配置后的错误

移除value.converter、schemas.enable和transforms相关配置后,出现以下错误:

Caused by: org.apache.kafka.connect.errors.DataException: 
JsonConverter with schemas.enable requires "schema" and "payload" fields and may not contain additional fields. 
If you are trying to deserialize plain JSON data, set schemas.enable=false in your converter configuration.

问题原因

  1. 第一次错误原因:ExtractNewRecordState转换器是专门处理Debezium CDC生成的带Envelope结构的消息(包含before/after/op等字段),你的消息是普通纯JSON,使用该转换器会导致Schema信息丢失,触发空指针异常。
  2. 第二次错误原因:移除配置后,Kafka Connect默认JSON转换器开启schemas.enable=true,要求消息必须包含schema和payload字段,但你的消息是纯JSON结构,不符合格式要求。

解决方案

使用以下修正后的配置,适配纯JSON消息的upsert同步:

config:
    class: io.debezium.connector.jdbc.JdbcSinkConnector

    connection.url: jdbc:postgresql://10.10.10.10:26257/db_dev
    connection.username: "******"
    connection.password: "********"
    topics: "kafka-crdb" 
    insert.mode: "upsert"
    primary.key.mode: "record_value"
    primary.key.fields: "sub_id"
    value.converter: "org.apache.kafka.connect.json.JsonConverter"
    value.converter.schemas.enable: "false"

关键配置说明

  • value.converter.schemas.enable: "false":明确告知转换器处理纯JSON消息,无需schema和payload结构
  • 移除transforms相关配置:消息不是Debezium CDC格式,无需使用ExtractNewRecordState进行unwrap处理
  • primary.key.mode: "record_value":指定从消息value中读取主键字段(对应primary.key.fields: "sub_id"),配合upsert模式实现插入/更新逻辑
  • insert.mode: "upsert":保持原有需求,根据主键判断操作类型

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 20:22:11