S3SinkConnector报「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

