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字段?
按以下逻辑配置即可,不需要修改现有生产者代码:
- 先解决Schema合法性问题:如果不需要解析字节数组的内部结构,就通过单消息转换(SMT)给原始字节套一层固定的Struct类型Schema,把原始Value作为Struct里的一个
bytes类型字段,保证顶层Schema是Record/Struct类型,满足Parquet写入的基本要求。 - 开启Header写入后,显式指定要落盘的Header列表,以及每个Header对应的数据类型(Header本身存储的是字节数组,连接器无法自动识别应该按字符串、数字还是其他类型解析)。
- 校验转换器和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配置指定字段和类型,可以换用HeaderToValueSMT,手动把需要的Header字段抽取到Value的Struct结构中,再写入Parquet即可。
内容的提问来源于stack exchange,提问作者code123
相关产品推荐
相关产品推荐

