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

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作为单一值存入数据库,可按以下步骤操作:

  1. 手动在PostgreSQL中创建目标表oi,需包含主键字段(对应消息的key,与pk.mode=record_key匹配)和存储JSON的字段(例如data jsonb)。
  2. 在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 19:45:32