☰
Apache Beam Java Kata 实战:用 TextIO.read() 从文本文件读取 PCollection
2026/9/27 4:21:29 网站建设 项目流程
  • 大数据
  • 批处理
  • 流处理
  • 数据工程

【免费下载链接】beam

Apache Beam is a unified programming model for Batch and Streaming data processing.

项目地址:https://gitcode.com/gh_mirrors/beam4/beam
点击查看免费下载

本篇技术指南以 Apache Beam 官方学习课程(Katas)中 TextIO Read 任务 为核心,讲解如何使用TextIO.read()与TextIO.Read.from(String)将一个或多个文本文件读入PCollection<String>,并通过一个"读取 countries.txt 并将国家名转为大写"的完整 Kata 实战,带你掌握文本文件读取的标准写法、配置项与底层实现原理。读完本文,你将能独立完成 Beam Java 管道中最常见的"读文本文件 → 逐行转换 → 验证结果"全流程。

为什么管道需要 I/O 变换

在 Beam 中构建管道时,通常需要从外部数据源读取数据,例如文件或数据库;同样,你可能希望将管道结果输出到外部存储系统。Beam 为多种常见存储类型内置了 read / write 变换,文本文件是最基础也最常用的一种。正如 TextIO Read 任务文档 所指出的:如果内置变换不支持你所需的存储格式,你还可以自行实现 read / write 变换;但在绝大多数场景下,TextIO已经足够。

在 Katas 课程 的 IO 章节 中,TextIO Read是第一个动手练习(该 lesson 仅包含这一项任务),它的目标非常明确:

Kata:读取countries.txt文件,并将每个国家名转换为大写。

TextIO.read() 基本用法

要从一个或多个文本文件读取PCollection,核心是两步:

  1. 用TextIO.read()实例化一个读取变换;
  2. 用TextIO.Read.from(String)指定要读取的文件或文件模式(filepattern)路径。

在 Katas 的 Task.java 中,读取部分正是这样完成的:

PipelineOptions options = PipelineOptionsFactory.fromArgs(args).create(); Pipeline pipeline = Pipeline.create(options); PCollection<String> countries = pipeline.apply("Read Countries", TextIO.read().from(FILE_PATH));

其中FILE_PATH是相对仓库根目录的路径:

private static final String FILE_PATH = "IO/TextIO/TextIO Read/countries.txt";

TextIO.read()是一个PTransform<PBegin, PCollection<String>>:它作用于管道起点PBegin,输出一个有界(bounded)的PCollection<String>,其中每一行输入文件对应一个元素(行尾换行符会被剥离)。

数据文件长什么样

本任务的数据文件 countries.txt 内容为 10 行国家名:

Singapore United States Australia England France China Indonesia Mexico Germany Japan

注意两点:其一,每行一个记录,这正是TextIO按行读取的天然匹配;其二,"United States"含空格,说明TextIO.read()不做分词,整行原样成为一个元素。

完整解题:读取并转大写

任务的完整解法在 Task.java 中通过一个可复用的applyTransform方法实现:

static PCollection<String> applyTransform(PCollection<String> input) { return input.apply(MapElements.into(strings()).via(String::toUpperCase)); }

这里使用了MapElements与TypeDescriptors.strings()(通过静态导入),把每个字符串元素映射为大写形式。applyTransform被设计为独立的静态方法,便于测试直接调用——这是 Katas 课程的标准模式。

整个管道的主流程为:

pipeline.apply("Read Countries", TextIO.read().from(FILE_PATH)); // 读取 applyTransform(countries); // 转换为大写 output.apply(Log.ofElements()); // 打印结果 pipeline.run();

Log.ofElements()来自 learning/katas/java/util 工具包,负责将PCollection的每个元素打印出来,便于本地观察运行结果。

如何运行

本 Kata 位于 learning/katas/java 模块,使用该目录下的gradlew即可运行:

cd learning/katas/java ./gradlew run -PmainClass=org.apache.beam.learning.katas.io.textio.read.Task

管道默认在 DirectRunner 上执行,countries.txt使用相对路径,因此请在learning/katas/java目录下运行(或按实际环境调整路径)。

测试如何验证

Katas 为每个任务都配有隐藏的单元测试,本任务的 TaskTest.java 展示了 Beam 官方的验证方式:

@Test public void textIO() { PCollection<String> countries = testPipeline.apply(TextIO.read().from("countries.txt")); PCollection<String> results = Task.applyTransform(countries); PAssert.that(results) .containsInAnyOrder( "AUSTRALIA", "CHINA", "ENGLAND", "FRANCE", "GERMANY", "INDONESIA", "JAPAN", "MEXICO", "SINGAPORE", "UNITED STATES"); testPipeline.run().waitUntilFinish(); }

要点解读:

  • TestPipeline.create()是 Beam 官方的测试管道(@Rule),自动处理pipeline.run()与断言时机;
  • TextIO.read().from("countries.txt")读取测试工作目录下的数据文件;
  • PAssert.that(results).containsInAnyOrder(...)断言结果集合与顺序无关地包含全部大写国家名——这正是分布式PCollection无序特性的体现;
  • 测试先调用Task.applyTransform,再断言结果,保证被测逻辑与管道构建解耦。

任务配置 task-info.yaml 中定义了两个TODO()占位符(分别对应读取与转换两个待补全位置),学习者需要自行补全后运行测试通过,即完成 Kata。

TextIO.read() 的完整配置项

除了最基础的from(String),Beam 的 TextIO.Read 还提供了丰富的链式配置方法,全部从源码 TextIO.java 中可直接确认:

配置方法作用默认值(来自read()源码)
from(String / ValueProvider<String>)指定文件路径或通配符模式,不可为 null无(必填,否则expand时抛异常)
withCompression(Compression)指定压缩类型Compression.AUTO(自动探测)
withDelimiter(byte[])自定义记录分隔符,替代默认的'\r'、'\n'、'\r\n'null(使用默认换行)
withSkipHeaderLines(int)跳过文件头部指定行数0
withHintMatchesManyFiles()提示 filepattern 匹配海量文件(数万级以上)false
withEmptyMatchTreatment(EmptyMatchTreatment)设置无文件匹配时的处理策略EmptyMatchTreatment.DISALLOW(不允许空匹配)
watchForNewFiles(Duration, TerminationCondition, boolean)周期性轮询等待新文件出现(需支持可拆分 DoFn 的 Runner)不启用

路径与通配符

from(String)中的 filepattern 可以是:

  • 本地路径(本地运行时),如countries.txt、/local/path/to/files/*;
  • 云存储路径(配合远程执行服务),如gs://<bucket>/<filepath>;
  • 支持标准 Java Filesystem glob 模式:*、?、[...]。

从源码 TextIO.java 可以看到,from(String)内部先做checkArgument(filepattern != null)校验,再包装为StaticValueProvider委托给from(ValueProvider<String>)——后者支持运行时才解析的值,便于在 Dataflow 等场景中延迟绑定参数。

压缩、分隔符与表头

读取压缩文件时,withCompression(Compression)支持AUTO/GZIP/BZIP2/DEFLATE/UNCOMPRESSED,默认AUTO会根据文件扩展名或魔数自动解压。自定义分隔符withDelimiter(byte[])则可用于读取非换行分隔的记录(如以\t或特定字节序列分隔),源码还专门校验分隔符不能"自重叠"(self-overlapping),避免边界解析歧义(见 TextIO.java#L408-L427)。

空匹配与流式监听

  • withEmptyMatchTreatment控制 filepattern 一个文件都匹配不到时的行为:默认DISALLOW直接失败,可改为ALLOW(返回空集合)或ALLOW_IF_WILDCARD(仅当模式本身含通配符时才允许)。
  • watchForNewFiles(pollInterval, terminationCondition, matchUpdatedFiles)让TextIO.read()具备流式文件监听能力(仅支持可拆分 DoFn 的 Runner,如 Dataflow 与 Flink)。类注释中的示例展示了每分钟轮询一次、一小时无新文件则停止的写法(见 TextIO.java#L119-L130)。

源码级原理:read() 内部如何工作

深入 TextIO.java 可以看清TextIO.read()的底层机制。

默认参数如何构建

read()静态工厂方法(TextIO.java#L196-L203)通过 AutoValue Builder 构造Read实例,默认配置为:

.setCompression(Compression.AUTO) .setHintMatchesManyFiles(false) .setSkipHeaderLines(0) .setMatchConfiguration(MatchConfiguration.create(EmptyMatchTreatment.DISALLOW))

这些默认值决定了不调用任何额外配置时TextIO.read().from(path)的行为:自动解压、不跳过表头、空匹配报错。

expand() 的分发逻辑

Read.expand()(TextIO.java#L429-L448)是核心分发点:

if (getMatchConfiguration().getWatchInterval() == null && !getHintMatchesManyFiles()) { return input.apply("Read", org.apache.beam.sdk.io.Read.from(getSource())); } // 其余情况走 FileIO + ReadFiles 组合 return input .apply("Create filepattern", Create.ofProvider(getFilepattern(), StringUtf8Coder.of())) .apply("Match All", FileIO.matchAll().withConfiguration(getMatchConfiguration())) .apply("Read Matches", FileIO.readMatches()...) .apply("Via ReadFiles", readFiles()...);

也就是说:

  • 常规静态读取(不监听新文件、不设海量文件提示)走Read.from(CompressedSource)的经典FileBasedSource路径,按 bundle 并行分片读取;
  • 一旦启用了watchForNewFiles或withHintMatchesManyFiles,则改写为FileIO.matchAll()+FileIO.readMatches()+readFiles()的组合,以获得流式监听与更高的文件级并行度。

getSource()(TextIO.java#L451-L459)则把TextSource(承载 filepattern、空匹配策略、分隔符与跳表头行数)包进CompressedSource,按指定压缩策略读取。

读取海量文件的性能提示

若 filepattern 会匹配非常多的文件(至少数万个),应使用withHintMatchesManyFiles()。源码注释明确说明:该提示可能让 Runner 以不同方式执行以提升性能;但如果实际只匹配少量文件,在支持动态工作再平衡的 Runner 上可能反而变慢(TextIO.java#L390-L401)。因此它是一把需要按场景谨慎使用的双刃剑。

进阶:readFiles() 与 FileIO 的组合

对于更复杂的读取场景,Beam 推荐显式组合FileIO与TextIO.readFiles()(TextIO.java#L233-L241),例如"先按目录匹配 → 过滤 → 再按文件读"。readFiles()读取PCollection<FileIO.ReadableFile>,其默认 bundle 大小为 64MB(DEFAULT_BUNDLE_SIZE_BYTES = 64 * 1024 * 1024L,见 TextIO.java#L190),用于在打开文件的成本与单次 ProcessElement 输出上限之间取得平衡。

PCollection<FileIO.ReadableFile> matched = pipeline.apply(FileIO.matchAll().withConfiguration(...)) .apply(FileIO.readMatches()); PCollection<String> lines = matched.apply(TextIO.readFiles());

旧的TextIO.readAll()在源码中已被标记@Deprecated,官方建议用上述FileIO组合替代(TextIO.java#L205-L227),因为组合方式让执行语义更显式,且ReadAll未来版本将被移除。

常见问题与最佳实践

  1. 文件路径找不到:TextIO.read()默认EmptyMatchTreatment.DISALLOW,filepattern 匹配不到任何文件会直接失败。本地运行时建议使用相对于工作目录的路径,或通过PipelineOptions参数化传入。
  2. 每一行是一个元素:TextIO.read()按行切分,不做类型解析,需要结构化数据时可在读取后用ParDo/MapElements自行解析(如本任务的String::toUpperCase)。
  3. 想读压缩文件:无需额外处理,默认Compression.AUTO自动识别;如确定文件未压缩可显式withCompression(Compression.UNCOMPRESSED)提升性能。
  4. 文件有表头:用withSkipHeaderLines(n)跳过,无需在业务逻辑里手动过滤。
  5. 大规模文件匹配:数万级以上文件用withHintMatchesManyFiles();需要等待新文件到达时用watchForNewFiles(...)(确认 Runner 支持可拆分 DoFn)。
  6. 测试优先:参考 TaskTest.java 的TestPipeline+PAssert模式,将读取逻辑与转换逻辑分离(如applyTransform),让管道可在不启动完整作业的情况下被单元测试覆盖。

延伸学习

  • 继续完成 Katas IO 章节 的其他任务,巩固读写变换;
  • 阅读 TextIO 完整源码,其中Read、ReadAll、ReadFiles、Write、TypedWrite、sink()等完整展示了文本 I/O 的全景;
  • 若需读取其他格式(如 Avro、Parquet、JDBC),可参考 Beam 内置 I/O 任务文档 的指引,Beam SDK 为多种数据源提供了开箱即用的变换。
  • 大数据
  • 批处理
  • 流处理
  • 数据工程

【免费下载链接】beam

Apache Beam is a unified programming model for Batch and Streaming data processing.

项目地址:https://gitcode.com/gh_mirrors/beam4/beam
点击查看免费下载

相关推荐

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询