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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 03:15:56