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

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 {}类型。

问题分析与解决方法

核心问题

  1. SFTP Source未拆分数组数据:输入文件的payload是数组,但SFTPJsonSourceConnector默认将整个文件内容作为单条Kafka消息发送,导致JdbcSinkConnector接收到的消息值是一个数组而非单个对象。因此sink连接器在顶层结构中找不到rollno字段(该字段实际存在于数组的每个元素内)。
  2. 主键配置不匹配:错误信息显示配置的主键字段是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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.24 16:45:46