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

S3SinkConnector报「Avro schema must be a record」异常求助

问题排查:S3 Sink Connector抛出「Avro schema must be a record」异常

环境与配置

使用版本:S3SinkConnector 10.3.0 + Kafka Connect 7.0.1
连接器配置如下:

{"connector.class": "io.confluent.connect.s3.S3SinkConnector","format.class": "io.confluent.connect.s3.format.parquet.ParquetFormat","flush.size": 1,"s3.bucket.name": "*****","s3.object.tagging": "true","s3.region": "us-east-2","aws.access.key.id": "*****","aws.secret.access.key": "*****","s3.part.retries": 5,"s3.retry.backoff.ms": 1000,"behavior.on.null.values": "ignore","keys.format.class": "io.confluent.connect.s3.format.json.JsonFormat","headers.format.class": "io.confluent.connect.s3.format.json.JsonFormat","store.kafka.headers": "true","store.kafka.keys": "true","topics": "***","storage.class": "io.confluent.connect.s3.storage.S3Storage","topics.dir": "kafka-backup","value.converter": "io.confluent.connect.json.JsonSchemaConverter","value.converter.schema.registry.url": "http://schema-registry:8081","key.converter": "org.apache.kafka.connect.storage.StringConverter","partitioner.class": "io.confluent.connect.storage.partitioner.HourlyPartitioner","locale": "en-US","timezone": "UTC","timestamp.extractor": "Record"}

问题现象

Kafka中的记录通过io.confluent.connect.json.JsonSchemaConverter以JSON格式存储,所有消息均有严格Schema,但Sink Connector读取记录时抛出如下异常:

org.apache.kafka.connect.errors.ConnectException: Exiting WorkerSinkTask due to unrecoverable exception.
    at org.apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:631)
    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:188)
    at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:243)
    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:829)
Caused by: java.lang.IllegalArgumentException: Avro schema must be a record.
    at org.apache.parquet.avro.AvroSchemaConverter.convert(AvroSchemaConverter.java:124)
    at org.apache.parquet.avro.AvroParquetWriter.writeSupport(AvroParquetWriter.java:150)
    at org.apache.parquet.avro.AvroParquetWriter.access$200(AvroParquetWriter.java:36)
    at org.apache.parquet.avro.AvroParquetWriter$Builder.getWriteSupport(AvroParquetWriter.java:182)
    at org.apache.parquet.hadoop.ParquetWriter$Builder.build(ParquetWriter.java:563)
    at io.confluent.connect.s3.format.parquet.ParquetRecordWriterProvider$1.write(ParquetRecordWriterProvider.java:102)
    at io.confluent.connect.s3.format.S3RetriableRecordWriter.write(S3RetriableRecordWriter.java:46)
    at io.confluent.connect.s3.format.KeyValueHeaderRecordWriterProvider$1.write(KeyValueHeaderRecordWriterProvider.java:107)
    at io.confluent.connect.s3.TopicPartitionWriter.writeRecord(TopicPartitionWriter.java:562)
    at io.confluent.connect.s3.TopicPartitionWriter.checkRotationOrAppend(TopicPartitionWriter.java:311)
    at io.confluent.connect.s3.TopicPartitionWriter.executeState(TopicPartitionWriter.java:254)
    at io.confluent.connect.s3.TopicPartitionWriter.write(TopicPartitionWriter.java:205)
    at io.confluent.connect.s3.S3SinkTask.put(S3SinkTask.java:234)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:601)

原因分析

从堆栈信息可见,Confluent官方的Parquet格式Sink内部依赖Avro工具链生成Parquet文件:ParquetRecordWriterProvider使用AvroParquetWriter写入数据,而AvroSchemaConverter要求输入的Avro Schema必须是record类型。

问题根源是:你的消息JSON Schema并非对象结构(比如是primitive类型、数组或map),经过JsonSchemaConverter转换后得到的Avro Schema对应的数据类型不是record,不符合AvroParquetWriter的要求,因此触发「Avro schema must be a record」异常。

解决方法

方案1:调整消息为对象结构

修改Kafka中的消息格式,确保消息是对象类型(而非单一primitive、数组或map)。例如:

  • 原始消息:"test-string"(string类型)
  • 修改后:{"content": "test-string"}(对象类型)
    同时更新对应的JSON Schema为对象结构,确保转换后的Avro Schema是record类型。

方案2:启用消息自动包装

在连接器配置中添加value.converter.wrap.message参数,让JsonSchemaConverter自动将非对象类型的消息包装为对象结构:

"value.converter.wrap.message": "true"

启用后,转换器会将原始消息包装为{"payload": 原始消息内容}的对象,转换后的Avro Schema即为record类型,满足Parquet Writer的要求。

方案3:使用非Avro依赖的Parquet格式实现

如果不想修改消息结构或包装消息,可以替换官方的Parquet格式实现,使用支持直接基于JSON Schema生成Parquet Schema的第三方组件,比如io.streamthoughts.kafka.connect.filepulse.format.parquet.ParquetFormat(来自Kafka Connect File Pulse)。注意需要确保组件版本与你的Kafka Connect版本兼容。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 05:20:25