如何在Apache Flink批处理作业中并行读取多个Sequence文件
解决方案:实现Flink自定义多路径SequenceFile读取器
我来帮你搞定这个Flink读取S3上大量Sequence文件的难题,彻底解决你遇到的逗号分隔路径失效、循环union崩溃的问题。咱们直接对标Spark的便捷读取体验,实现一个自定义的读取器。
核心痛点拆解
你遇到的问题本质是:
- Flink原生
SequenceFileInputFormat对多S3路径的支持有限,逗号分隔的方式不生效 - 文件量太大时,循环读取+多次
union会带来巨大的调度开销,甚至触发稳定性问题 - 需要像Spark那样,能一次性批量处理多个路径的SequenceFile读取能力
步骤1:自定义支持多路径的SequenceFileInputFormat
咱们扩展Flink的Hadoop兼容层,实现一个能同时处理多个S3路径的InputFormat,内部统一处理所有文件的split,不用外部做union操作。
import org.apache.flink.api.common.io.DefaultInputSplitAssigner; import org.apache.flink.api.common.io.InputFormat; import org.apache.flink.api.common.io.InputSplit; import org.apache.flink.api.common.io.InputSplitAssigner; import org.apache.flink.core.io.InputSplitSource; import org.apache.flink.core.io.LocatableInputSplit; import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.Writable; import org.apache.hadoop.mapreduce.lib.input.SequenceFileInputFormat; import java.io.IOException; import java.util.ArrayList; import java.util.List; public class MultiPathSequenceFileInputFormat<K extends Writable, V extends Writable> implements InputFormat<K, V>, InputSplitSource<InputSplit> { private final SequenceFileInputFormat<K, V> hadoopInputFormat; private List<Path> paths; public MultiPathSequenceFileInputFormat() { this.hadoopInputFormat = new SequenceFileInputFormat<>(); } // 设置要读取的所有S3路径 public void setPaths(List<String> pathStrings) { this.paths = new ArrayList<>(); pathStrings.forEach(pathStr -> this.paths.add(new Path(pathStr))); } @Override public void configure(org.apache.flink.configuration.Configuration parameters) { try { // 利用Hadoop原生能力加载多路径 hadoopInputFormat.setInputPaths(new org.apache.hadoop.conf.Configuration(), paths.toArray(new Path[0])); } catch (IOException e) { throw new RuntimeException("Failed to initialize SequenceFile input paths", e); } } @Override public InputSplit[] createInputSplits(int minNumSplits) throws IOException { // 生成Hadoop的文件split,转成Flink的InputSplit格式 org.apache.hadoop.mapreduce.InputSplit[] hadoopSplits = hadoopInputFormat.getSplits(new org.apache.hadoop.mapreduce.JobContext() { @Override public org.apache.hadoop.conf.Configuration getConfiguration() { return new org.apache.hadoop.conf.Configuration(); } }); List<InputSplit> flinkSplits = new ArrayList<>(); for (int i = 0; i < hadoopSplits.length; i++) { flinkSplits.add(new LocatableInputSplit(i, hadoopSplits[i].getLocations())); } return flinkSplits.toArray(new InputSplit[0]); } @Override public InputSplitAssigner getInputSplitAssigner(InputSplit[] inputSplits) { return new DefaultInputSplitAssigner(inputSplits); } @Override public void open(InputSplit split) throws IOException { hadoopInputFormat.open(((LocatableInputSplit) split).getHadoopSplit()); } @Override public boolean reachedEnd() throws IOException { return hadoopInputFormat.reachedEnd(); } @Override public K nextRecord(K reuse) throws IOException { return hadoopInputFormat.nextKeyValue() ? hadoopInputFormat.getCurrentKey() : null; } @Override public void close() throws IOException { hadoopInputFormat.close(); } }
步骤2:封装成Spark风格的便捷工具类
为了用起来像Spark一样顺手,咱们写个工具类,把读取逻辑封装成简单的方法:
import org.apache.flink.api.java.ExecutionEnvironment; import org.apache.flink.api.java.DataSet; import org.apache.flink.api.java.tuple.Tuple2; import org.apache.hadoop.io.Writable; import java.util.List; public class SequenceFileReader { // 只读取Key的方法 public static <K extends Writable, V extends Writable> DataSet<K> readKeys(ExecutionEnvironment env, List<String> paths, Class<K> keyClass) { MultiPathSequenceFileInputFormat<K, V> inputFormat = new MultiPathSequenceFileInputFormat<>(); inputFormat.setPaths(paths); return env.createInput(inputFormat, keyClass); } // 只读取Value的方法 public static <K extends Writable, V extends Writable> DataSet<V> readValues(ExecutionEnvironment env, List<String> paths, Class<V> valueClass) { MultiPathSequenceFileInputFormat<K, V> inputFormat = new MultiPathSequenceFileInputFormat<>(); inputFormat.setPaths(paths); return env.createInput(inputFormat, valueClass).map(value -> value); } // 同时读取Key和Value,返回Tuple2 public static <K extends Writable, V extends Writable> DataSet<Tuple2<K, V>> read(ExecutionEnvironment env, List<String> paths, Class<K> keyClass, Class<V> valueClass) { MultiPathSequenceFileInputFormat<K, V> inputFormat = new MultiPathSequenceFileInputFormat<>(); inputFormat.setPaths(paths); return env.createInput(inputFormat, keyClass, valueClass); } }
步骤3:实际使用示例
现在你可以像用Spark那样,一行代码搞定大量S3 SequenceFile的读取:
import org.apache.flink.api.java.ExecutionEnvironment; import org.apache.flink.api.java.tuple.Tuple2; import org.apache.hadoop.io.LongWritable; import org.apache.hadoop.io.Text; import java.util.Arrays; import java.util.List; public class S3SequenceFileBatchJob { public static void main(String[] args) throws Exception { ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment(); // 这里可以传多个具体路径,也可以用通配符匹配整个目录下的文件 List<String> s3Paths = Arrays.asList( "s3://your-bucket/data/file1.seq", "s3://your-bucket/data/file2.seq", "s3://your-bucket/data/sub-dir/*" ); // 读取Key和Value,类型根据你的SequenceFile实际情况调整 DataSet<Tuple2<LongWritable, Text>> sequenceData = SequenceFileReader.read( env, s3Paths, LongWritable.class, Text.class ); // 后续的业务处理逻辑 sequenceData.print(); env.execute("Flink S3 SequenceFile Batch Job"); } }
关键优势说明
- 多路径/通配符支持:直接利用Hadoop的文件系统能力,支持多个具体路径、目录通配符(比如
*),不用手动枚举所有文件 - 彻底避免Union问题:一次性生成所有文件的InputSplit,不需要循环调用
union,消除了多次union带来的调度开销和稳定性风险 - 兼容原生逻辑:基于Flink的Hadoop兼容层实现,和原生SequenceFile读取逻辑完全一致,保证数据正确性
如果你的文件是按前缀组织的,直接用通配符路径会更方便,比如s3://your-bucket/year=2024/month=05/*,Hadoop会自动匹配所有符合条件的文件。
内容的提问来源于stack exchange,提问作者Abhinav Prakash
相关产品推荐
相关产品推荐

