- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
本篇技术指南以 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,核心是两步:
- 用
TextIO.read()实例化一个读取变换; - 用
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未来版本将被移除。
常见问题与最佳实践
- 文件路径找不到:
TextIO.read()默认EmptyMatchTreatment.DISALLOW,filepattern 匹配不到任何文件会直接失败。本地运行时建议使用相对于工作目录的路径,或通过PipelineOptions参数化传入。 - 每一行是一个元素:
TextIO.read()按行切分,不做类型解析,需要结构化数据时可在读取后用ParDo/MapElements自行解析(如本任务的String::toUpperCase)。 - 想读压缩文件:无需额外处理,默认
Compression.AUTO自动识别;如确定文件未压缩可显式withCompression(Compression.UNCOMPRESSED)提升性能。 - 文件有表头:用
withSkipHeaderLines(n)跳过,无需在业务逻辑里手动过滤。 - 大规模文件匹配:数万级以上文件用
withHintMatchesManyFiles();需要等待新文件到达时用watchForNewFiles(...)(确认 Runner 支持可拆分 DoFn)。 - 测试优先:参考 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.
相关推荐
Faker 实战指南:用 Ruby 生成假面骑士(Kamen Rider)假数据 —— Faker::JapaneseMedia::KamenRider 完整用法
Faker 实战指南:用 Ruby 生成假面骑士(Kamen Rider)假数据 —— Faker::JapaneseMedia::KamenRider 完整用
大数据批处理流处理数据工程Apache Beam Go SDK 文本 I/O 实战:用 textio.Read 读取文件并将 PCollection 转为大写
Apache Beam Go SDK 文本 I/O 实战:用 textio.Read 读取文件并将 PCollection 转为大写 Apache Beam 是
大数据批处理流处理数据工程Apache Beam Java Kata 实战:使用 Count 聚合变换统计 PCollection 元素个数
Apache Beam Java Kata 实战:使用 Count 聚合变换统计 PCollection 元素个数 本指南围绕 Apache Beam 官方 J
大数据批处理流处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考