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

如何将Apache Beam TextIO异常转为失败并分离失效文件列表?

解决Apache Beam TextIO读取文件时异常终止Pipeline的问题

当使用TextIO.readAll()读取不存在的文件时,会抛出如下异常并终止整个Pipeline:

org.apache.beam.sdk.Pipeline$PipelineExecutionException:
java.io.FileNotFoundException: No files matched spec:
src/test/resources/config/qqqqqqq

要实现Pipeline不终止,同时分离出有效文件内容和失效文件列表,核心做法是用Apache Beam的Try类型包装文件读取结果,将成功读取的内容和失败的文件路径分开处理。

完整解决方案代码

import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.PipelineOptions;
import org.apache.beam.sdk.PipelineOptionsFactory;
import org.apache.beam.sdk.io.TextIO;
import org.apache.beam.sdk.transforms.DoFn;
import org.apache.beam.sdk.transforms.FlatMapElements;
import org.apache.beam.sdk.transforms.MapElements;
import org.apache.beam.sdk.transforms.ParDo;
import org.apache.beam.sdk.transforms.partition.Partition;
import org.apache.beam.sdk.values.PCollection;
import org.apache.beam.sdk.values.PCollectionList;
import org.apache.beam.sdk.values.TypeDescriptors;
import org.apache.beam.vendor.guava.v26_0_jre.com.google.common.base.Try;

import java.io.FileNotFoundException;
import java.nio.file.Files;
import java.nio.file.Paths;
import java.util.List;

public class FileReadWithErrorHandling {

    public static void main(String[] args) {
        PipelineOptions options = PipelineOptionsFactory.create();
        Pipeline pipeline = Pipeline.create(options);

        // 读取文件路径列表
        PCollection<String> filePaths = pipeline
                .apply(TextIO.read().from("src/test/resources/config/W-PSFV-LOG-FILE-2022-05-16_23-59-59.txt"))
                .apply(MapElements.into(TypeDescriptors.strings())
                        .via(line -> "src/test/resources/config/" + line));

        // 包装文件读取结果为Try类型(区分成功/失败)
        PCollection<Try<List<String>>> readResults = filePaths.apply(ParDo.of(new ReadFileWithTryFn()));

        // 拆分成功和失败的结果分支
        PCollectionList<Try<List<String>>> partitioned = readResults.apply(Partition.of(2, (elem, numPartitions) -> {
            return elem.isSuccess() ? 0 : 1;
        }));

        // 处理成功分支:展开文件内容并输出
        PCollection<String> validContents = partitioned.get(0)
                .apply(FlatMapElements.into(TypeDescriptors.strings())
                        .via(tryResult -> tryResult.get()));

        // 处理失败分支:提取失效文件路径并输出
        PCollection<String> invalidFiles = partitioned.get(1)
                .apply(MapElements.into(TypeDescriptors.strings())
                        .via(tryResult -> {
                            Throwable throwable = tryResult.getFailure();
                            if (throwable instanceof FileNotFoundException) {
                                return "失效文件: " + ((FileNotFoundException) throwable).getMessage().split(": ")[1];
                            }
                            return "读取失败: " + throwable.getMessage();
                        }));

        // 输出结果到文件(也可替换为其他Sink)
        validContents.apply(TextIO.write().to("output/valid_contents").withSuffix(".txt"));
        invalidFiles.apply(TextIO.write().to("output/invalid_files").withSuffix(".txt"));

        pipeline.run().waitUntilFinish();
    }

    // 自定义DoFn:单个文件读取+异常捕获,返回Try包装的结果
    static class ReadFileWithTryFn extends DoFn<String, Try<List<String>>> {
        @ProcessElement
        public void processElement(@Element String filePath, OutputReceiver<Try<List<String>>> receiver) {
            try {
                List<String> lines = Files.readAllLines(Paths.get(filePath));
                receiver.output(Try.success(lines));
            } catch (FileNotFoundException e) {
                receiver.output(Try.failure(e));
            } catch (Exception e) {
                // 可扩展处理其他读取异常(如权限问题)
                receiver.output(Try.failure(e));
            }
        }
    }
}

关键逻辑说明

  1. ReadFileWithTryFn:负责单个文件的读取操作,捕获FileNotFoundException等异常,用Try.success返回成功读取的内容,Try.failure返回异常信息。单个文件读取失败不会终止整个Pipeline。
  2. Partition转换:将包含Try结果的数据集拆分为两个分支,索引0对应成功读取的结果,索引1对应读取失败的结果。
  3. 结果分流处理:成功分支展开文件内容列表,输出到有效内容文件;失败分支提取失效文件路径,输出到失效文件列表。

注意事项

  • 示例中用Files.readAllLines读取文件,适合小文件场景;如果处理大文件,建议改用TextIO.read()结合侧输出的方式,避免内存占用过高。
  • 可根据业务需求扩展异常处理逻辑,比如过滤特定异常、记录详细错误日志等。

内容的提问来源于stack exchange,提问作者veaceslav.kunitki

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.23 04:24:26