Kafka Connect PostgreSQL Sink报错Value schema must be of type Struct如何解决?
问题描述
将嵌套JSON数据发送至PostgreSQL Sink的Kafka消费者,构建Sink Connector时无法修改源端数据,希望直接发送原始数据不做转换,但Kafka Connect抛出如下错误:
[2023-01-04 22:58:15,227] ERROR WorkerSinkTask{id=Kafkapgsink-0} Task threw an uncaught and unrecoverable exception. Task is being killed and will not recover until manually restarted. Error: Value schema must be of type Struct (org.apache.kafka.connect.runtime.WorkerSinkTask:609) org.apache.kafka.connect.errors.ConnectException: Value schema must be of type Struct at io.confluent.connect.jdbc.sink.metadata.FieldsMetadata.extract(FieldsMetadata.java:86) at io.confluent.connect.jdbc.sink.metadata.FieldsMetadata.extract(FieldsMetadata.java:67) at io.confluent.connect.jdbc.sink.BufferedRecords.add(BufferedRecords.java:115) 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:834)
当前全局连接器配置
bootstrap.servers=localhost:9092 key.converter=org.apache.kafka.connect.storage.StringConverter value.converter=org.apache.kafka.connect.storage.StringConverter key.converter.schemas.enable=false value.converter.schemas.enable=false offset.storage.file.filename=/tmp/connect.offsets offset.flush.interval.ms=10000
Sink Connector配置
name=Kafkapgsink connector.class=io.confluent.connect.jdbc.JdbcSinkConnector task.max=100 connection.url=jdbc:postgresql://localhost:5432/fileintegrity connection.user=postgres connection.password=09900 insert.mode=insert auto.create=true auto.evolve=true table.name.format=oi pk.mode=record_key delete.enabled=true
解决方案
错误原因
JDBC Sink Connector需要将消息数据映射到数据库表的列,因此要求消息值为Struct类型(对应Kafka Connect的结构化Schema)。但当前配置使用StringConverter,将嵌套JSON当作纯字符串处理,未解析为Struct,导致不符合Connector的类型要求。
方案1:用JSON Converter自动解析结构化数据
修改全局连接器的value.converter配置,使用JSON Converter解析嵌套JSON为Struct:
# 替换原value.converter相关配置 value.converter=org.apache.kafka.connect.json.JsonConverter value.converter.schemas.enable=false # 源数据无Schema,保持关闭
- 此方案会自动将JSON的顶级字段映射到数据库表的列,嵌套JSON部分会被转换为PostgreSQL的
jsonb类型(需PostgreSQL 9.4+版本支持)。 - 结合
auto.create=true和auto.evolve=true,Connector会自动创建或更新表结构以匹配JSON字段。
方案2:将原始JSON存入单个字段
如果无需解析JSON内部字段,只想把整个嵌套JSON作为单一值存入数据库,可按以下步骤操作:
- 手动在PostgreSQL中创建目标表
oi,需包含主键字段(对应消息的key,与pk.mode=record_key匹配)和存储JSON的字段(例如data jsonb)。 - 在Sink Connector配置中添加Single Message Transform(SMT),将字符串消息值包装为Struct:
# 在原有Sink配置基础上添加以下内容 transforms=unwrap transforms.unwrap.type=org.apache.kafka.connect.transforms.HoistField$Value transforms.unwrap.field=data
- 该SMT会将整个原始JSON字符串包装为仅含
data字段的Struct,Connector即可将其写入数据库的data列。
注意事项
- 方案1中,若JSON结构复杂或频繁变化,可能导致数据库表结构频繁变更,需评估是否符合业务稳定性要求。
pk.mode=record_key要求消息的key类型与数据库表的主键字段类型一致,否则会出现类型不匹配错误。
内容的提问来源于stack exchange,提问作者user20822279
相关产品推荐
相关产品推荐

