Kafka Connect同步SFTP JSON到MySQL报错:record_value缺指定PK字段
问题描述
尝试通过SFTPJsonSourceConnector将SFTP中的JSON数据同步至MySQL(使用JdbcSinkConnector)时,出现以下错误:
[2022-07-31 00:03:20,239] ERROR [local-mysql-snik|task-2] WorkerSinkTask{id=local-mysql-snik-2} Task threw an uncaught and unrecoverable exception. Task is being killed and will not recover until manually restarted (org.apache.kafka.connect.runtime.WorkerTask:207)org.apache.kafka.connect.errors.ConnectException: Exiting WorkerSinkTask due to unrecoverable exception.at apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:618) at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:334) at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:235) at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:204) at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:200) at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:255) 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:834) Caused by: org.apache.kafka.connect.errors.ConnectException: PK mode for table 'SAMPLE' is RECORD_VALUE with configured PK fields [rollno], but record value schema does not contain field: rollno at io.confluent.connect.jdbc.sink.metadata.FieldsMetadata.extractRecordValuePk(FieldsMetadata.java:279) at io.confluent.connect.jdbc.sink.metadata.FieldsMetadata.extract(FieldsMetadata.java:104) at io.confluent.connect.jdbc.sink.metadata.FieldsMetadata.extract(FieldsMetadata.java:66) at io.confluent.connect.jdbc.sink.BufferedRecords.add(BufferedRecords.java:116) at io.confluent.connect.jdbc.sink.JdbcDbWriter.write(JdbcDbWriter.java:66) at io.confluent.connect.jdbc.sink.JdbcSinkTask.put(JdbcSinkTask.java:74) at org.apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:584)
相关文件与配置
输入JSON文件
{"schema": {"type": "struct","fields": [{"type": "int32", "field": "rollno"}, {"type": "string","field": "first_name"}],"name": "simple"},"payload": [{"rollno": 2,"first_name": "Sai"},{"rollno": 3,"first_name": "kumar"}]
connect-standalone.properties
key.converter=org.apache.kafka.connect.storage.StringConverter value.converter=org.apache.kafka.connect.json.JsonConverter schema.generation.enabled=true key.converter.schemas.enable=false value.converter.schemas.enable=true internal.key.converter=org.apache.kafka.connect.json.JsonConverter internal.value.converter=org.apache.kafka.connect.json.JsonConverter internal.key.converter.schemas.enable=false internal.value.converter.schemas.enable=false offset.flush.interval.ms=10000 plugin.path=/usr/share/java,/home/local/confluent-7.2.1/share/confluent-hub-components
sftp-source.properties
name=local-JsonSftp tasks.max=3 connector.class=io.confluent.connect.sftp.SftpJsonSourceConnector input.path=/home/local/confluent-7.2.1/path/to/data error.path=/home/local/confluent-7.2.1/path/to/data finished.path=/home/local/confluent-7.2.1/path/to/data cleanup.policy=NONE input.file.pattern=simple-test.json behavior.on.error=IGNORE sftp.username=myuser sftp.password=user@123 sftp.host=localhost sftp.port=22 kafka.topic=SAMPLE schema.generation.enabled=true value.converter.schema.registry.url=http://localhost:8081 key.converter=org.apache.kafka.connect.storage.StringConverter value.converter=org.apache.kafka.connect.json.JsonConverter key.converter.schemas.enable=false value.converter.schemas.enable=true
mysql-sink.properties
name=local-mysql-snik connector.class=io.confluent.connect.jdbc.JdbcSinkConnector tasks.max=3 topics=SAMPLE connection.url=jdbc:mysql://localhost:3306/empdetails?user=root&password=user@321 connection.user=user1 connection.password=user@321 schema.generation.enabled=true value.converter.schema.registry.url=http://localhost:8081 key.converter=org.apache.kafka.connect.storage.StringConverter value.converter=org.apache.kafka.connect.json.JsonConverter key.converter.schemas.enable=false value.converter.schemas.enable=true auto.create=true auto.evolve=true insert.mode=upsert pk.mode=record_value pk.fields=id
注意:使用pk.mode=record_key时,创建并插入的表为struct {}类型。
问题分析与解决方法
核心问题
- SFTP Source未拆分数组数据:输入文件的
payload是数组,但SFTPJsonSourceConnector默认将整个文件内容作为单条Kafka消息发送,导致JdbcSinkConnector接收到的消息值是一个数组而非单个对象。因此sink连接器在顶层结构中找不到rollno字段(该字段实际存在于数组的每个元素内)。 - 主键配置不匹配:错误信息显示配置的主键字段是
rollno,但mysql-sink.properties中写的是pk.fields=id,存在配置不一致(可能是配置未更新或输入错误)。
解决步骤
步骤1:配置SFTP Source拆分数组
修改sftp-source.properties,添加以下配置,让连接器将payload数组拆分为多条独立的Kafka消息:
# 指定处理JSON数组,将每个元素作为单独记录 mode=JSON # 配置从payload字段读取数组内容 json.schema.location=INLINE json.payload.field=payload
步骤2:修正主键配置
根据输入数据的实际字段,调整mysql-sink.properties的主键配置:
pk.mode=record_value pk.fields=rollno
输入数据中唯一的标识字段是rollno,而非id(输入数据中不存在id字段)。
步骤3:检查转换器配置一致性
确保connect-standalone.properties、source和sink的转换器配置保持一致:
- 由于使用带schema的JSON格式,
value.converter.schemas.enable=true需统一设置 - 若未使用Schema Registry,可移除
value.converter.schema.registry.url配置
步骤4:重启Connect服务
修改配置后,重启Kafka Connect standalone服务,确保新配置生效。
内容的提问来源于stack exchange,提问作者user2045848
相关产品推荐
相关产品推荐

