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

如何控制Flink Streaming API FileSink输出Parquet文件的记录数?

问题描述

使用Flink 1.7.1(问题与版本无关),通过Streaming API的FileSink将Avro GenericRecord写入Parquet文件,功能正常但出现4条记录生成4个文件的情况(Checkpoint间隔10s,所有记录同时到达),需要控制每个文件中的记录数。

代码示例:

FileSink<GenericRecord> sink = FileSink
        .forBulkFormat(outputPath, AvroParquetWriters.forGenericRecord(avroSchema))
        .withRollingPolicy(
                OnCheckpointRollingPolicy
                        .build()
        )
        //.withBucketAssigner()
        .build();
解决方案

要控制Parquet文件的记录数,需结合批量写入配置和滚动策略调整,具体如下:

  • 配置批量写入的记录阈值
    针对AvroParquetWriter设置批量大小,当内存中缓存的记录数达到该阈值时,会在Checkpoint触发时写入文件。通过withBatchSize()方法指定单批记录数:

    AvroParquetWriters.forGenericRecord(avroSchema)
        .withBatchSize(1000); // 示例:每1000条记录生成一个文件
    
  • 替换滚动策略
    原代码使用的OnCheckpointRollingPolicy会在每次Checkpoint时强制滚动文件,哪怕当前文件只有少量记录。建议替换为DefaultRollingPolicy,它可以结合多个条件触发滚动(如文件大小、空闲时长、活跃时长),同时配合批量大小实现按记录数控文件:

    DefaultRollingPolicy<GenericRecord, String> rollingPolicy = DefaultRollingPolicy.builder()
        .withMaxPartSize(1024 * 1024 * 1024) // 可选:设置文件最大大小
        .withRolloverInterval(TimeUnit.MINUTES.toMillis(15)) // 可选:设置文件最长活跃时间
        .withInactivityInterval(TimeUnit.MINUTES.toMillis(5)) // 可选:设置空闲超时时间
        .build();
    
  • 整合后的完整代码

    FileSink<GenericRecord> sink = FileSink
            .forBulkFormat(outputPath, 
                AvroParquetWriters.forGenericRecord(avroSchema)
                    .withBatchSize(1000) // 控制单批记录数
            )
            .withRollingPolicy(rollingPolicy)
            //.withBucketAssigner() // 如需按规则分桶可启用
            .build();
    
  • 补充说明

    1. 批量大小的设置需要结合内存情况,避免过大导致OOM
    2. 如果并行度较高,每个并行子任务会独立生成文件,若想减少文件总数,可适当降低并行度,或使用BucketAssigner将数据路由到同一桶中
    3. BulkFormat的写入依赖Checkpoint,需确保Checkpoint机制正常运行

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 04:13:16