Apache Beam Java Kata 实战:用 TextIO.read() 从文本文件读取 PCollection
大数据批处理流处理数据工程【免费下载链接】beamApache 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)将一个或多个文本文件读入PCollectionString并通过一个读取 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); PCollectionString 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()是一个PTransformPBegin, PCollectionString它作用于管道起点PBegin输出一个有界bounded的PCollectionString其中每一行输入文件对应一个元素行尾换行符会被剥离。数据文件长什么样本任务的数据文件 countries.txt 内容为 10 行国家名Singapore United States Australia England France China Indonesia Mexico Germany Japan注意两点其一每行一个记录这正是TextIO按行读取的天然匹配其二United States含空格说明TextIO.read()不做分词整行原样成为一个元素。完整解题读取并转大写任务的完整解法在 Task.java 中通过一个可复用的applyTransform方法实现static PCollectionString applyTransform(PCollectionString 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 -PmainClassorg.apache.beam.learning.katas.io.textio.read.Task管道默认在 DirectRunner 上执行countries.txt使用相对路径因此请在learning/katas/java目录下运行或按实际环境调整路径。测试如何验证Katas 为每个任务都配有隐藏的单元测试本任务的 TaskTest.java 展示了 Beam 官方的验证方式Test public void textIO() { PCollectionString countries testPipeline.apply(TextIO.read().from(countries.txt)); PCollectionString 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 / ValueProviderString)指定文件路径或通配符模式不可为 null无必填否则expand时抛异常withCompression(Compression)指定压缩类型Compression.AUTO自动探测withDelimiter(byte[])自定义记录分隔符替代默认的\r、\n、\r\nnull使用默认换行withSkipHeaderLines(int)跳过文件头部指定行数0withHintMatchesManyFiles()提示 filepattern 匹配海量文件数万级以上falsewithEmptyMatchTreatment(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(ValueProviderString)——后者支持运行时才解析的值便于在 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()读取PCollectionFileIO.ReadableFile其默认 bundle 大小为 64MBDEFAULT_BUNDLE_SIZE_BYTES 64 * 1024 * 1024L见 TextIO.java#L190用于在打开文件的成本与单次 ProcessElement 输出上限之间取得平衡。PCollectionFileIO.ReadableFile matched pipeline.apply(FileIO.matchAll().withConfiguration(...)) .apply(FileIO.readMatches()); PCollectionString lines matched.apply(TextIO.readFiles());旧的TextIO.readAll()在源码中已被标记Deprecated官方建议用上述FileIO组合替代TextIO.java#L205-L227因为组合方式让执行语义更显式且ReadAll未来版本将被移除。常见问题与最佳实践文件路径找不到TextIO.read()默认EmptyMatchTreatment.DISALLOWfilepattern 匹配不到任何文件会直接失败。本地运行时建议使用相对于工作目录的路径或通过PipelineOptions参数化传入。每一行是一个元素TextIO.read()按行切分不做类型解析需要结构化数据时可在读取后用ParDo/MapElements自行解析如本任务的String::toUpperCase。想读压缩文件无需额外处理默认Compression.AUTO自动识别如确定文件未压缩可显式withCompression(Compression.UNCOMPRESSED)提升性能。文件有表头用withSkipHeaderLines(n)跳过无需在业务逻辑里手动过滤。大规模文件匹配数万级以上文件用withHintMatchesManyFiles()需要等待新文件到达时用watchForNewFiles(...)确认 Runner 支持可拆分 DoFn。测试优先参考 TaskTest.java 的TestPipelinePAssert模式将读取逻辑与转换逻辑分离如applyTransform让管道可在不启动完整作业的情况下被单元测试覆盖。延伸学习继续完成 Katas IO 章节 的其他任务巩固读写变换阅读 TextIO 完整源码其中Read、ReadAll、ReadFiles、Write、TypedWrite、sink()等完整展示了文本 I/O 的全景若需读取其他格式如 Avro、Parquet、JDBC可参考 Beam 内置 I/O 任务文档 的指引Beam SDK 为多种数据源提供了开箱即用的变换。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐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),仅供参考