如何将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)); } } } }
关键逻辑说明
- ReadFileWithTryFn:负责单个文件的读取操作,捕获
FileNotFoundException等异常,用Try.success返回成功读取的内容,Try.failure返回异常信息。单个文件读取失败不会终止整个Pipeline。 - Partition转换:将包含Try结果的数据集拆分为两个分支,索引0对应成功读取的结果,索引1对应读取失败的结果。
- 结果分流处理:成功分支展开文件内容列表,输出到有效内容文件;失败分支提取失效文件路径,输出到失效文件列表。
注意事项
- 示例中用
Files.readAllLines读取文件,适合小文件场景;如果处理大文件,建议改用TextIO.read()结合侧输出的方式,避免内存占用过高。 - 可根据业务需求扩展异常处理逻辑,比如过滤特定异常、记录详细错误日志等。
内容的提问来源于stack exchange,提问作者veaceslav.kunitki
相关产品推荐
相关产品推荐

