从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会尝试重复扫描文件,压缩流无法满足这种随机访问需求。
具体解决步骤:
- 修改文件处理模式
将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) );
- 规范文件命名
Flink通过文件名后缀识别压缩格式,需将文件重命名为tweets.avro.gz(而非仅.gz),确保Flink正确识别这是gzip压缩的Avro文件,自动启用解压逻辑。同时更新代码中的路径:
String path = "s3a://my-test-bucket/tweets.avro.gz";
- 验证依赖配置
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
相关产品推荐
相关产品推荐

