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

Flink DataStream API中Parquet Sink输出文件如何配置压缩?

当然可以在DataStream API中配置Parquet输出文件的压缩,只是配置方式和Table API不同,需要直接通过ParquetOutputFormat或自定义ParquetWriterFactory来指定压缩格式:

方法一:通过ParquetOutputFormat的内置配置

如果使用Flink提供的ParquetOutputFormat构建批量写入的FileSink,可以直接调用withCompressionCodec方法指定压缩编码:

import org.apache.flink.connector.file.sink.FileSink;
import org.apache.flink.formats.parquet.avro.ParquetOutputFormat;
import org.apache.parquet.hadoop.metadata.CompressionCodecName;
import org.apache.flink.core.fs.Path;

// 假设你的数据类型为自定义POJO/Avro类 MyEnhancedData
ParquetOutputFormat<MyEnhancedData> parquetFormat = ParquetOutputFormat
    .buildParquetOutputFormat()
    .withSchema(AvroSchemaConverter.convertToSchema(MyEnhancedData.class)) // 关联数据的Avro Schema
    .withCompressionCodec(CompressionCodecName.GZIP) // 设置压缩格式,比如GZIP/SNAPPY等
    .build();

// 创建写入S3的FileSink
FileSink<MyEnhancedData> s3Sink = FileSink
    .forBulkFormat(new Path("s3://your-target-bucket/output-path"), parquetFormat)
    .build();

// 将处理后的数据流写入S3
enhancedDataStream.sinkTo(s3Sink);

方法二:自定义ParquetWriterFactory

如果需要更灵活的配置,可以通过自定义ParquetWriterFactory,直接在AvroParquetWriter构造时指定压缩选项:

import org.apache.flink.connector.file.sink.FileSink;
import org.apache.flink.formats.parquet.ParquetWriterFactory;
import org.apache.parquet.avro.AvroParquetWriter;
import org.apache.parquet.hadoop.metadata.CompressionCodecName;
import org.apache.flink.core.fs.Path;

ParquetWriterFactory<MyEnhancedData> customWriterFactory = new ParquetWriterFactory<>(
    outputStream -> AvroParquetWriter.<MyEnhancedData>builder(outputStream)
        .withSchema(AvroSchemaConverter.convertToSchema(MyEnhancedData.class))
        .withCompressionCodec(CompressionCodecName.SNAPPY) // 选择压缩格式
        .build()
);

FileSink<MyEnhancedData> s3Sink = FileSink
    .forBulkFormat(new Path("s3://your-target-bucket/output-path"), customWriterFactory)
    .build();

enhancedDataStream.sinkTo(s3Sink);

可用压缩编码

压缩格式基于Parquet官方的CompressionCodecName枚举,常用选项包括:

  • GZIP: 高压缩比,适合存储场景
  • SNAPPY: 低延迟高吞吐,适合实时/准实时场景
  • BROTLI: 更高压缩比,压缩耗时略长
  • LZO: 平衡压缩比和性能

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 22:48:42