Apache Beam从Kafka读取数据后无法写入Parquet文件求助
问题描述
我需要将Kafka主题中固定数量(约15k条)的数据通过批处理读取并写入Parquet文件到S3。目前数据读取、解析逻辑均正常,Row转GenericRecord的转换也已通过日志验证,但数据始终无法写入Parquet文件,且无任何错误日志输出。我的需求是:不希望基于偏移量进行窗口操作,数据按Kafka分区生成多个文件存储在S3中,且无需等待所有数据读取完成再开始处理。
以下是我的代码:
PCollection<KafkaRecord<byte[], byte[]>> messages = pipeline.apply("ReadFromKafka", KafkaIO.<byte[], byte[]>read() .withBootstrapServers(CONSUMER_CONFIGS.getKafkaServer()) .withConsumerFactoryFn(new consumerFactory()) ); PCollection<KafkaRecord<byte[], byte[]>> windowedMessages = messages.apply(Window.<KafkaRecord<byte[], byte[]>>into(new GlobalWindows()) .triggering(AfterWatermark.pastEndOfWindow()) // Trigger once after watermark passes the end of the window .discardingFiredPanes()); PCollection<Row> A= (ConvertToRowFn) A.setCoder(RowCoder.of(schema2)); String outputPath = configs.getFolderPath() ; PCollection<GenericRecord> records = A.apply("Convert Rows to GenericRecord", MapElements.into(TypeDescriptor.of(GenericRecord.class)) .via(this::convertRowToGenericRecord)); records.setCoder(AvroCoder.of(GenericRecord.class, AvroSchema.getSchema())); records.apply("Write Parquet", FileIO.<GenericRecord> write() .via(ParquetIO.sink(AvroSchema.getSchema())) .to(outputPath) .withSuffix(PARQUET)); pipeline.run().waitUntilFinish();
问题分析与解决方案
核心问题
你使用的GlobalWindow+AfterWatermark.pastEndOfWindow()组合在批处理场景下完全无效:GlobalWindow的结束时间是无穷大,批处理中的水印永远无法推进到这个时间点,导致窗口永远不会触发,数据一直滞留在窗口中,无法进入后续的写入环节,自然不会生成任何Parquet文件。
修复步骤
1. 移除无效的全局窗口配置
直接删除windowedMessages相关的窗口操作代码,批处理模式下不需要这种全局窗口设置,数据会被按分区并行处理。
2. 按Kafka分区生成独立文件
要实现按Kafka分区输出文件,需要保留KafkaRecord的分区信息,并用动态写入功能拆分输出:
- 从KafkaRecord中提取分区号作为分组Key
- 使用
FileIO.writeDynamic根据分区号生成独立文件路径
3. 修复代码语法错误
代码中PCollection<Row> A= (ConvertToRowFn)存在语法错误,正确写法应为:
PCollection<Row> rowCollection = messages.apply("Convert to Row", ParDo.of(new ConvertToRowFn()));
修正后完整代码示例
import org.apache.beam.sdk.io.FileIO; import org.apache.beam.sdk.io.kafka.KafkaIO; import org.apache.beam.sdk.io.kafka.KafkaRecord; import org.apache.beam.sdk.io.parquet.ParquetIO; import org.apache.beam.sdk.transforms.DoFn; import org.apache.beam.sdk.transforms.MapElements; import org.apache.beam.sdk.transforms.ParDo; import org.apache.beam.sdk.values.KV; import org.apache.beam.sdk.values.PCollection; import org.apache.beam.sdk.values.TypeDescriptor; import org.apache.avro.generic.GenericRecord; import org.apache.beam.sdk.coders.BigEndianIntegerCoder; import org.apache.beam.sdk.coders.KvCoder; import java.time.Duration; import java.util.Collections; // 读取Kafka全量数据(批处理模式) PCollection<KafkaRecord<byte[], byte[]>> messages = pipeline.apply("ReadFromKafka", KafkaIO.<byte[], byte[]>read() .withBootstrapServers(CONSUMER_CONFIGS.getKafkaServer()) .withConsumerFactoryFn(new consumerFactory()) .withTopics(Collections.singletonList("your-kafka-topic")) // 替换为你的主题名 .withReadCommitted() .withStartReadTime(Duration.ZERO) // 从最早偏移量开始读 .withEndReadTime(Duration.ofMillis(System.currentTimeMillis()))); // 读到当前最新偏移量 // 提取Kafka分区号,并转换为Row PCollection<KV<Integer, Row>> partitionedRows = messages.apply("Extract Partition & Convert to Row", ParDo.of(new DoFn<KafkaRecord<byte[], byte[]>, KV<Integer, Row>>() { @ProcessElement public void processElement(ProcessContext c) { KafkaRecord<byte[], byte[]> record = c.element(); int partition = record.getPartition(); // 调用你的转换逻辑生成Row Row row = new ConvertToRowFn().convert(record); // 根据你的ConvertToRowFn实现调整 c.output(KV.of(partition, row)); } })); partitionedRows.setCoder(KvCoder.of(BigEndianIntegerCoder.of(), RowCoder.of(schema2))); // 转换为GenericRecord,保留分区Key PCollection<KV<Integer, GenericRecord>> partitionedRecords = partitionedRows.apply("Convert to GenericRecord", MapElements.into(TypeDescriptor.of(KV.class)) .via(kv -> KV.of(kv.getKey(), convertRowToGenericRecord(kv.getValue())))); partitionedRecords.setCoder(KvCoder.of(BigEndianIntegerCoder.of(), AvroCoder.of(GenericRecord.class, AvroSchema.getSchema()))); // 按分区动态写入S3 Parquet文件 String outputPath = configs.getFolderPath(); partitionedRecords.apply("Write Parquet by Partition", FileIO.<Integer, GenericRecord>writeDynamic() .by(KV::getKey) // 按分区号作为分组Key .withDestinationCoder(BigEndianIntegerCoder.of()) .via(Contextful.fn(record -> ParquetIO.sink(AvroSchema.getSchema())), ParquetIO.sink(AvroSchema.getSchema())) .to(outputPath) .withNaming(partition -> FileIO.Write.defaultNaming("partition-" + partition, ".parquet"))); // 每个分区生成独立文件 pipeline.run().waitUntilFinish();
额外说明
- 批处理模式下,KafkaIO需要明确读取范围,通过
withStartReadTime和withEndReadTime指定读取全量数据 FileIO.writeDynamic会自动按分区并行写入文件,无需等待所有数据读取完成- 确保
ConvertToRowFn的转换逻辑能正确处理KafkaRecord的原始数据,并保留必要的元数据
内容的提问来源于stack exchange,提问作者IndiePump
相关产品推荐
相关产品推荐

