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

使用MSK Connector同步MSK至RDS PostgreSQL配置报错求助

问题描述

使用kafkajs向Kafka发送数据,每条消息的key设为UUID,value是序列化后的事件字符串:

// 生产者用TypeScript编写
const event = {
    eventtype: "event1",
    eventversion: "1.0.1",
    sourceurl: "https://some-url.com/source"
};
// 序列化字符串,因为kafkajs生产者只接受`string`或`Buffer`类型
const stringifiedEvent = JSON.stringify(event);

通过以下配置启动独立模式的JDBC Sink Connector:

connect-standalone.properties

name=local-jdbc-sink-connector
connector.class=io.confluent.connect.jdbc.JdbcSinkConnector
dialect.name=PostgreSqlDatabaseDialect
connection.url=jdbc:postgresql://postgres:5432/eventservice
connection.password=postgres
connection.user=postgres

auto.create=true
auto.evolve=true
topics=topic1
tasks.max=1
insert.mode=upsert
pk.mode=record_key
pk.fields=id

worker.properties

offset.storage.file.filename=/tmp/connect.offsets
offset.flush.interval.ms=10000

value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=false
value.converter.schema.registry.url=http://schema-registry:8081 

key.converter=org.apache.kafka.connect.storage.StringConverter
key.converter.schemas.enable=false

bootstrap.servers=localhost:9092
group.id=jdbc-sink-connector-worker
worker.id=jdbc-sink-worker-1

offset.storage.topic=connect-offsets
offset.storage.replication.factor=1
config.storage.topic=connect-configs
config.storage.replication.factor=1
status.storage.topic=connect-status
status.storage.replication.factor=1

启动连接器后能正常连接PostgreSQL,但生产消息时出现以下错误:

WorkerSinkTask{id=local-jdbc-sink-connector-0} Task threw an uncaught and unrecoverable exception. 
Task is being killed and will not recover until manually restarted. Error: Sink connector 'local-jdbc-sink-connector' is configured with 'delete.enabled=false' and 'pk.mode=record_key' and therefore requires records with a non-null Struct value and non-null Struct schema, but found record at (topic='topic1',partition=0,offset=0,timestamp=1676309784254) with a HashMap value and null value schema. 
(org.apache.kafka.connect.runtime.WorkerSinkTask:609)

堆栈信息:

org.apache.kafka.connect.errors.ConnectException: Sink connector 'local-jdbc-sink-connector' is configured with 
'delete.enabled=false' and 'pk.mode=record_key' and therefore requires records with a non-null Struct value and 
non-null Struct schema, but found record at (topic='txningestion2',partition=0,offset=0,timestamp=1676309784254) 
with a HashMap value and null value schema.
    at io.confluent.connect.jdbc.sink.RecordValidator.lambda$requiresValue$2(RecordValidator.java:86)
    at io.confluent.connect.jdbc.sink.RecordValidator.lambda$and$1(RecordValidator.java:41)
    at io.confluent.connect.jdbc.sink.BufferedRecords.add(BufferedRecords.java:81)
    at io.confluent.connect.jdbc.sink.JdbcDbWriter.write(JdbcDbWriter.java:74)
    at io.confluent.connect.jdbc.sink.JdbcSinkTask.put(JdbcSinkTask.java:85)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:581)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:333)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:234)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:203)
    at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:189)
    at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:244)
    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)

反复调整配置仍无法解决,求正确配置及解决方案。


解决方案

错误核心原因:当pk.mode=record_key且delete.enabled=false时,JDBC Sink要求消息value必须是带Schema的Struct类型,但当前配置中value.converter.schemas.enable=false,导致value被解析成无Schema的HashMap,不符合要求。

有两种可行的解决思路:

思路1:修改PK配置,从value字段提取主键

如果不需要用Kafka消息的key作为数据库主键,可调整JDBC Sink的主键配置:

  • 将pk.mode改为record_value
  • 确保pk.fields对应event中存在的字段(比如event里的唯一标识字段,或新增一个id字段到event中)

修改后的connect-standalone.properties相关配置:

insert.mode=upsert
pk.mode=record_value
pk.fields=eventtype  # 替换为你实际用来作为主键的字段,比如新增的id字段

这种方式不需要修改消息发送格式,也不需要启用Schema,适合不想引入Schema Registry的场景。

思路2:启用Schema,让value转为Struct类型

如果坚持要用Kafka消息的key作为数据库主键,需要让消息value带有Schema:

  1. 修改worker.properties中的value转换器配置,启用Schema并移除多余配置:
value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=true
# 移除该行:value.converter.schema.registry.url=http://schema-registry:8081
  1. 修改kafkajs生产者代码,发送带Schema的JSON消息:
const eventWithSchema = {
  schema: {
    type: "struct",
    fields: [
      { field: "eventtype", type: "string" },
      { field: "eventversion", type: "string" },
      { field: "sourceurl", type: "string" }
    ]
  },
  payload: {
    eventtype: "event1",
    eventversion: "1.0.1",
    sourceurl: "https://some-url.com/source"
  }
};
const stringifiedEvent = JSON.stringify(eventWithSchema);
  1. 确保数据库表的主键字段id类型与Kafka消息的key(UUID字符串)匹配,PostgreSQL中可设为VARCHAR(36)类型。

额外细节:当前worker.properties同时配置了文件偏移存储和Kafka偏移存储,建议只保留一种避免冲突。生产环境推荐用Kafka主题存储偏移量,删除offset.storage.file.filename=/tmp/connect.offsets这一行。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 16:00:59