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

Kafka Connect写入S3 Parquet报Avro schema must be a record错误排查

Kafka Connect S3 Sink写入Parquet异常排查与配置方案

场景复现

尝试通过Kafka Connect消费Kafka集群消息,以Parquet格式写入S3存储,编写了发送byte[]类型消息的简易生产者,相关代码、连接器配置、运行报错信息如下:

生产者代码

Properties propertiesAWS = new Properties();
    propertiesAWS.setProperty(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "myKafka:9092");
    propertiesAWS.setProperty(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, LongSerializer.class.getName());
    propertiesAWS.setProperty(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, ByteArraySerializer.class.getName());


    KafkaProducer<Long, byte[]> producer = new KafkaProducer<Long, byte[]>(propertiesAWS);
    Random rng = new Random();

    for (int i = 0; i < 100; i++) {
        try {
            Thread.sleep(1000);
            Headers headers = new RecordHeaders();
            headers.add(new RecordHeader("header1", "header1".getBytes()));
            headers.add(new RecordHeader("header2", "header2".getBytes()));
            ProducerRecord<Long, byte[]> recordOut = new ProducerRecord<Long, byte[]>
                    ("s3.test.topic", 1, rng.nextLong(), new byte[]{1, 2, 3}, headers);
            producer.send(recordOut);
        } catch (Exception e) {
            System.out.println(e);

        }
    }

S3 Sink连接器初始配置

{
"name": "test_2_s3",
"config": {
    "connector.class": "io.confluent.connect.s3.S3SinkConnector",
    "aws.access.key.id": "XXXXXXX",
    "aws.secret.access.key": "XXXXXXXX",
    "s3.region": "eu-central-1",
    "flush.size": "5",
    "rotate.schedule.interval.ms": "10000",
    "timezone": "UTC",
    "tasks.max": "1",
    "topics": "s3.test.topic",
    "parquet.codec": "gzip",
    "format.class": "io.confluent.connect.s3.format.parquet.ParquetFormat",
    "partitioner.class": "io.confluent.connect.storage.partitioner.DefaultPartitioner",
    "storage.class": "io.confluent.connect.s3.storage.S3Storage",
    "s3.bucket.name": "test-phase1",
    "key.converter": "org.apache.kafka.connect.converters.LongConverter",
    "value.converter": "org.apache.kafka.connect.converters.ByteArrayConverter",
    "behavior.on.null.values": "ignore",
    "store.kafka.headers": "true"
}

运行异常

Caused by: java.lang.IllegalArgumentException: Avro schema must be a record.
at org.apache.parquet.avro.AvroSchemaConverter.convert(AvroSchemaConverter.java:124)

问题解答

1. 当前配置的错误点

核心问题是Parquet作为列式存储格式,写入时必须拿到顶层为Record类型的结构化Schema,当前配置完全不满足这个要求:

  • Confluent S3 Sink自带的ParquetFormat默认通过Avro Schema转换生成Parquet文件元数据,你配置的ByteArrayConverter解析消息后,输出的Value是无结构的原始字节类型,对应的Schema是基础字节类型,不是Avro要求的顶层Record结构,AvroSchemaConverter处理不了这种非Record的顶层Schema,直接抛出参数非法异常。
  • 你只开启了store.kafka.headers=true,但没指定哪些Header要写入、每个Header对应什么数据类型,连接器拿到的Header全是字节数组,无法直接映射成Parquet的结构化列。

2. 若仅需要写入字节数组消息内容与指定Kafka消息头,是否必须使用Avro格式?

完全不需要。
Avro只是消息序列化方式的一种,Parquet要求的是明确的结构化Schema定义,不是必须用Avro序列化生产端消息。你完全可以继续发送纯字节的消息,只要给连接器传入符合要求的顶层Record结构Schema即可:比如把整个消息体作为一个bytes类型的列,要存储的Header分别对应成独立列,不需要生产者侧把消息序列化成Avro格式。
如果完全不想维护Schema,不要选择Parquet格式,换用纯二进制的字节文件格式即可。

3. 如何配置连接器,指定需要写入Parquet文件的Kafka header字段?

按以下逻辑配置即可,不需要修改现有生产者代码:

  1. 先解决Schema合法性问题:如果不需要解析字节数组的内部结构,就通过单消息转换(SMT)给原始字节套一层固定的Struct类型Schema,把原始Value作为Struct里的一个bytes类型字段,保证顶层Schema是Record/Struct类型,满足Parquet写入的基本要求。
  2. 开启Header写入后,显式指定要落盘的Header列表,以及每个Header对应的数据类型(Header本身存储的是字节数组,连接器无法自动识别应该按字符串、数字还是其他类型解析)。
  3. 校验转换器和SMT逻辑,确保最终传给Parquet写入器的结构是带明确字段名、字段类型的Struct,包含你要存储的Value列和Header列。

可直接运行的最小修正配置参考:

{
"name": "test_2_s3_fixed",
"config": {
    "connector.class": "io.confluent.connect.s3.S3SinkConnector",
    "aws.access.key.id": "XXXXXXX",
    "aws.secret.access.key": "XXXXXXXX",
    "s3.region": "eu-central-1",
    "flush.size": "5",
    "rotate.schedule.interval.ms": "10000",
    "timezone": "UTC",
    "tasks.max": "1",
    "topics": "s3.test.topic",
    "parquet.codec": "gzip",
    "format.class": "io.confluent.connect.s3.format.parquet.ParquetFormat",
    "partitioner.class": "io.confluent.connect.storage.partitioner.DefaultPartitioner",
    "storage.class": "io.confluent.connect.s3.storage.S3Storage",
    "s3.bucket.name": "test-phase1",
    "key.converter": "org.apache.kafka.connect.converters.LongConverter",
    "value.converter": "org.apache.kafka.connect.converters.ByteArrayConverter",
    "behavior.on.null.values": "ignore",
    // Header写入配置
    "store.kafka.headers": "true",
    "headers": "header1,header2",
    "header1.type": "string",
    "header2.type": "string",
    // SMT给原始字节套合法的Struct Schema
    "transforms": "wrapStruct",
    "transforms.wrapStruct.type": "org.apache.kafka.connect.transforms.SetSchemaMetadata$Value",
    "transforms.wrapStruct.schema.name": "s3.parquet.record",
    "transforms.wrapStruct.schema.type": "STRUCT",
    "transforms.wrapStruct.schema.fields": "msg_bytes:bytes"
}
}

如果你用的是5.x及更早版本的Confluent连接器,不支持直接通过headers配置指定字段和类型,可以换用HeaderToValue SMT,手动把需要的Header字段抽取到Value的Struct结构中,再写入Parquet即可。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 18:57:21