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

添加FieldPartitioner后AWS MSK S3 Sink Connector无法正常工作

解决MSK Connect S3 Sink添加FieldPartitioner后的"Value is not Struct type"错误

问题根源

报错Value is not Struct type的核心原因是转换器配置不匹配:
你当前使用value.converter=org.apache.kafka.connect.storage.StringConverter,这个转换器会将Kafka消息值转换成字符串类型,但FieldPartitioner需要从结构化的Struct类型数据中提取分区字段(比如你的name字段),字符串无法被分区器解析,导致分区编码失败。

另外配置里存在拼写错误:value.convertor.schemaName中的convertor应为converter,这个错误也会影响Schema Registry的正常调用。

解决方案

将value转换器替换为支持Glue Schema Registry的Avro转换器,确保消息被解析为结构化的GenericRecord(对应Struct类型),让FieldPartitioner可以正常提取name字段。

修正后的完整配置

connector.class=io.confluent.connect.s3.S3SinkConnector
format.class=io.confluent.connect.s3.format.avro.AvroFormat
flush.size=1
schema.compatibility=BACKWARD
tasks.max=2
topics=MSKTutorialTopic
storage.class=io.confluent.connect.s3.storage.S3Storage
topics.dir=mskTrials
s3.bucket.name=clickstream
s3.region=us-east-1

partitioner.class=io.confluent.connect.storage.partitioner.FieldPartitioner
partition.field.name=name

# 替换为Glue Schema Registry的Avro转换器
value.converter=software.amazon.msk.connect.converters.AvroConverter
value.converter.schemaAutoRegistrationEnabled=true
value.converter.registry.name=datalake-schema-registry
# 修正拼写错误:convertor -> converter
value.converter.schemaName=MSKTutorialTopic-value
value.converter.avroRecordType=GENERIC_RECORD
value.converter.region=us-east-1
value.converter.schemas.enable=true

key.converter=org.apache.kafka.connect.storage.StringConverter

关键修改说明

  1. 替换value转换器:用software.amazon.msk.connect.converters.AvroConverter替代StringConverter,该转换器会从Glue Schema Registry拉取你定义的UserData schema,将Kafka中的Avro二进制消息解析为GenericRecord(即Struct类型),满足FieldPartitioner的要求。
  2. 修正拼写错误:将value.convertor.schemaName改为value.converter.schemaName,确保配置项能被正确识别。
  3. 保留其他配置:原有的分区器、存储格式、S3相关配置均可保留,只需调整转换器部分。

验证步骤

  1. 停止当前故障的连接器任务,应用上述修正后的配置重启连接器。
  2. 检查MSK Connect任务日志,确认无Value is not Struct type或PartitionException报错。
  3. 查看目标S3桶clickstream/mskTrials/MSKTutorialTopic路径下,是否生成以name=xxx为分区目录的Avro文件。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 19:41:54