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

AWS MSK+S3 Sink问题:JSON消息以Struct格式存储而非标准JSON

排查S3 Sink输出Struct格式而非标准JSON的问题

核心问题在于S3 Sink连接器未配置正确的Value转换器,导致Kafka Connect将Debezium输出的Struct类型消息直接以toString的方式序列化,最终在S3中呈现为Struct{...}格式。以下是具体排查和解决方向:

1. 补充S3 Sink的Value转换器配置

当前你的S3 Sink配置仅指定了key.converter,缺少value.converter配置。根据Debezium的输出格式,需对应配置转换器:

场景1:Debezium输出纯JSON(无Schema)

如果Debezium通过unwrap转换后已输出纯数据记录,添加以下配置:

value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=false

schemas.enable=false表示不需要携带Schema信息,直接序列化纯JSON内容。

场景2:Debezium输出Avro格式(使用Schema Registry)

若你使用Confluent Schema Registry管理Avro Schema,需配置Avro转换器:

value.converter=io.confluent.connect.avro.AvroConverter
value.converter.schema.registry.url=http://your-schema-registry-endpoint:8081

修改后的完整S3 Sink配置示例(纯JSON场景):

connector.class=io.confluent.connect.s3.S3SinkConnector
s3.region=us-east-1
flush.size=65536
schema.compatibility=NONE
tasks.max=1
topics=reservation,reservation2
timezone=UTC
format.class=io.confluent.connect.s3.format.json.JsonFormat
partitioner.class=io.confluent.connect.storage.partitioner.DefaultPartitioner
schema.generator.class=io.confluent.connect.storage.hive.schema.DefaultSchemaGenerator
storage.class=io.confluent.connect.storage.S3Storage
key.converter=org.apache.kafka.connect.storage.StringConverter
# 新增Value转换器配置
value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=false
s3.bucket.name=bucket
rotate.schedule.interval.ms=60000

2. 验证Kafka主题中的消息格式

使用Kafka控制台消费者工具直接查看主题内容,确认Debezium的输出是否符合预期:

kafka-console-consumer.sh --bootstrap-server <your-msk-bootstrap-servers> --topic reservation --from-beginning \
  --property print.key=true \
  --property value.deserializer=org.apache.kafka.connect.json.JsonConverter \
  --property value.deserializer.schemas.enable=false
  • 如果输出是标准JSON,说明Debezium的unwrap转换生效,问题仅在S3 Sink的转换器配置;
  • 如果输出仍包含Envelope结构(如before/after/op字段),需检查Debezium的transforms.unwrap配置是否正确应用。

3. 检查MSK Worker的全局转换器配置

MSK的默认Worker配置(connect-distributed.properties)中若设置了全局key.converter和value.converter,会被连接器级别的配置覆盖,但如果连接器未指定,则会使用全局配置。确认全局配置是否与Debezium的输出格式匹配:

  • 若全局value.converter为Avro,但S3 Sink未配置对应的转换器,会导致序列化异常;
  • 若全局value.converter为JsonConverter但schemas.enable=true,也会导致输出带Schema的JSON,需在连接器级别覆盖该配置。

4. 确认Debezium的unwrap转换有效性

检查Debezium的transforms=unwrap,dropTopicPrefix配置顺序是否正确,unwrap应优先于其他转换执行,确保提取出变更后的纯数据记录。若顺序颠倒,可能导致unwrap未正确处理消息。

内容的提问来源于stack exchange,提问作者jyablonski

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 17:25:27