使用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:
- 修改
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
- 修改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);
- 确保数据库表的主键字段
id类型与Kafka消息的key(UUID字符串)匹配,PostgreSQL中可设为VARCHAR(36)类型。
额外细节:当前worker.properties同时配置了文件偏移存储和Kafka偏移存储,建议只保留一种避免冲突。生产环境推荐用Kafka主题存储偏移量,删除offset.storage.file.filename=/tmp/connect.offsets这一行。
内容的提问来源于stack exchange,提问作者RusskiT
相关产品推荐
相关产品推荐

