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

从S3读取压缩Avro文件时Flink报InputStream无法反向搜索错误

Flink读取S3上gzip压缩Avro文件报错解决方案

问题场景

学习Flink期间,尝试从S3读取Avro编码的gzip压缩文件并处理内容。未压缩的.avro文件可正常读取,但读取.gz文件时抛出如下异常:

2022-08-12 15:51:21,418 WARN org.apache.flink.runtime.taskmanager.Task [] - Split Reader: Custom File Source -> Map -> Sink: Print to Std. Out (1/1)#0 (bcd4507e1eac498d2c2fea1c4785679e) switched from RUNNING to FAILED with failure cause: java.lang.IllegalArgumentException: Wrapped InputStream: cannot search backwards.

使用代码如下:

public class StreamingJob {
    public static void main(String[] args) throws Exception {
        // String path = "s3a://my-test-bucket/tweets.avro"; // <-- this works fine
        String path = "s3a://my-test-bucket/tweets.gz";

        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        AvroInputFormat inputFormat = new AvroInputFormat<TwitterSchema>(
            new Path(path), 
            TwitterSchema.class
        );

        DataStreamSource<TwitterSchema> ds = env.readFile(
            inputFormat,
            path,
            FileProcessingMode.PROCESS_CONTINUOUSLY,
            10000l, // 10 seconds
            TypeInformation.of(TwitterSchema.class)
        );

        ds
            .map(t -> t.getTweet())
            .print();

        env.execute("s3 job example");
    }
}

版本信息:

  • Flink版本:1.14.5
  • Java版本:1.8

根据Flink文档「读取压缩文件」章节,预期Flink可自动解压文件,但实际报错。

问题原因与解决方法

异常核心原因是AvroInputFormat依赖可随机访问的输入流(支持seek操作),但gzip压缩流不支持反向搜索/seek。当使用PROCESS_CONTINUOUSLY模式时,Flink会尝试重复扫描文件,压缩流无法满足这种随机访问需求。

具体解决步骤:

  1. 修改文件处理模式
    将FileProcessingMode.PROCESS_CONTINUOUSLY改为FileProcessingMode.PROCESS_ONCE。该模式下Flink仅读取文件一次,无需对压缩流执行seek操作,适配gzip这类不支持随机访问的压缩格式。

修改后的代码片段:

DataStreamSource<TwitterSchema> ds = env.readFile(
    inputFormat,
    path,
    FileProcessingMode.PROCESS_ONCE, // 替换为PROCESS_ONCE
    10000l,
    TypeInformation.of(TwitterSchema.class)
);
  1. 规范文件命名
    Flink通过文件名后缀识别压缩格式,需将文件重命名为tweets.avro.gz(而非仅.gz),确保Flink正确识别这是gzip压缩的Avro文件,自动启用解压逻辑。同时更新代码中的路径:
String path = "s3a://my-test-bucket/tweets.avro.gz";
  1. 验证依赖配置
    Flink 1.14.x默认已包含gzip支持,若为自定义构建项目,需确认引入flink-compress依赖:
<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-compress</artifactId>
    <version>1.14.5</version>
</dependency>

额外说明

若业务必须使用PROCESS_CONTINUOUSLY模式监控文件变化,可选择:

  • 上传前先将gzip文件解压为Avro文件;
  • 改用支持随机访问的压缩格式(如Snappy,需对应调整文件后缀和依赖)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 08:48:21