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

如何在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");
    }
}

关键优势说明

  1. 多路径/通配符支持:直接利用Hadoop的文件系统能力,支持多个具体路径、目录通配符(比如*),不用手动枚举所有文件
  2. 彻底避免Union问题:一次性生成所有文件的InputSplit,不需要循环调用union,消除了多次union带来的调度开销和稳定性风险
  3. 兼容原生逻辑:基于Flink的Hadoop兼容层实现,和原生SequenceFile读取逻辑完全一致,保证数据正确性

如果你的文件是按前缀组织的,直接用通配符路径会更方便,比如s3://your-bucket/year=2024/month=05/*,Hadoop会自动匹配所有符合条件的文件。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 10:02:16