添加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
关键修改说明
- 替换value转换器:用
software.amazon.msk.connect.converters.AvroConverter替代StringConverter,该转换器会从Glue Schema Registry拉取你定义的UserDataschema,将Kafka中的Avro二进制消息解析为GenericRecord(即Struct类型),满足FieldPartitioner的要求。 - 修正拼写错误:将
value.convertor.schemaName改为value.converter.schemaName,确保配置项能被正确识别。 - 保留其他配置:原有的分区器、存储格式、S3相关配置均可保留,只需调整转换器部分。
验证步骤
- 停止当前故障的连接器任务,应用上述修正后的配置重启连接器。
- 检查MSK Connect任务日志,确认无
Value is not Struct type或PartitionException报错。 - 查看目标S3桶
clickstream/mskTrials/MSKTutorialTopic路径下,是否生成以name=xxx为分区目录的Avro文件。
内容的提问来源于stack exchange,提问作者Sreenath
相关产品推荐
相关产品推荐

