使用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.
问题原因
- 第一次错误原因:
ExtractNewRecordState转换器是专门处理Debezium CDC生成的带Envelope结构的消息(包含before/after/op等字段),你的消息是普通纯JSON,使用该转换器会导致Schema信息丢失,触发空指针异常。 - 第二次错误原因:移除配置后,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
相关产品推荐
相关产品推荐

