首页/新闻资讯/正文详情

Apache Beam Java Kata 实战:用 TextIO.read() 从文本文件读取 PCollection

发布时间:2026/9/27 4:21:34 来源:云帆数科 栏目:资讯中心
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),仅供参考

相关推荐

2026最新网站建设推广代运营避坑实录:从备案到SEO全链路拆解
2026最新网站建设推广代运营避坑实录:从备案到SEO全链路拆解

2026最新网站建设推广代运营避坑实录:从备案到SEO全链路拆解 做网站最怕什么?不是代码写不出来,也不是设计不好看,而是卡在“备案流程一头雾水”这一步,急得团团转。很多老板找外包做【网站建设推广代运营】,结果网站建好了,因为备案资料填错、… · 2026/9/27 4:21:28

STM32 流水灯三种写法:寄存器、标准外设库、HAL 库,以及一个让我意外的延时结论
STM32 流水灯三种写法:寄存器、标准外设库、HAL 库,以及一个让我意外的延时结论

这学期做 STM32 实验,同一个流水灯(4 只 LED 依次亮 1 秒)我用三种方式各写了一遍: 直接读写寄存器、ST 标准外设库、STM32Cube HAL 库 按键中断。 写完之后有两个收获挺意外的,记一下。 一、先说三种写法差在哪寄存器… · 2026/9/27 4:21:28

强电磁环境下以太网温湿度变送器的EMC设计实践
强电磁环境下以太网温湿度变送器的EMC设计实践

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/27 4:21:22

5G网络优化干扰排查全攻略:从分类到定位的实用指南
5G网络优化干扰排查全攻略:从分类到定位的实用指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/27 5:09:24

Vivado ILA抓不到数据?从时钟树到XDC约束的深度排查指南
Vivado ILA抓不到数据?从时钟树到XDC约束的深度排查指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/27 5:09:24

3个坑搞定一站式服务门户源码下载避坑
3个坑搞定一站式服务门户源码下载避坑

3个坑搞定一站式服务门户源码下载避坑 域名服务器搞不懂?别慌,这是西南中小老板建一站式服务门户时的头号噩梦。很多老板拿着预算,心里却打鼓:到底该买现成的还是找团队定制?… · 2026/9/27 5:09:24

股票查询网站模板wordpress新手入门避坑与加固
股票查询网站模板wordpress新手入门避坑与加固

股票查询网站模板wordpress新手入门避坑与加固 别再用那些一眼假、配色烂大街的模板了。你辛辛苦苦做的股票查询站,用户点进去第一反应不是“专业”,而是“这网站靠谱吗?会不会偷我钱?”这种不信任感,直接导致跳出率飙升,SEO排名也上不去。… · 2026/9/27 5:09:06

3个免费降AIGC网站,让你的论文彻底告别AI痕迹[必看]
3个免费降AIGC网站,让你的论文彻底告别AI痕迹[必看]

最近不少同学私信我,说论文明明是自己一个字一个字敲的,就用了AI帮忙理了理思路,结果学校AIGC检测直接飙到30%以上,整个人都懵了。这事儿不是个例,现在各大查重平台都加了AI检测功能,查重率能压到10%以下&a… · 2026/9/27 5:09:00

会议记录总翻车?2025实测:AI录音卡+智能转写,如何把1小时会议压缩成3分钟精华
会议记录总翻车?2025实测:AI录音卡+智能转写,如何把1小时会议压缩成3分钟精华

你有没有经历过这种崩溃瞬间——开了2小时的项目评审会,全程录音,会后对着长达3小时的音频文件发呆,从头听一遍?保守估计又要2小时。快进听?关键信息一不留神就滑过去了。自己手动整理会议纪要?写了开头就没… · 2026/9/27 5:09:00

MATLAB雷达信号脉冲压缩仿真:LFM线性调频、匹配滤波与距离分辨率实现
MATLAB雷达信号脉冲压缩仿真:LFM线性调频、匹配滤波与距离分辨率实现

简介:这套Matlab仿真工具完整呈现雷达信号脉冲压缩过程,从线性调频(LFM)信号生成、目标回波仿真到匹配滤波压缩处理均有可运行代码支撑,面向电子信息工程、计算机、数学等专业学生,适用于课程设计、期末大作… · 2026/9/27 0:00:01

汕头网站建设制作厂家避坑指南:5大注意事项救急
汕头网站建设制作厂家避坑指南:5大注意事项救急

汕头网站建设制作厂家避坑指南:5大注意事项救急 改个需求建站公司拖一周,这种憋屈事我见得太多了。 很多汕头老板找本地建站团队,签合同前看着方案挺美,一上线就变脸。 今天不聊虚的,直接拆解找 汕头网站建设制作厂家 时的5个核心 注意事项… · 2026/9/27 0:00:01

多模态虚假新闻检测实战:BERT+ResNet双塔与对比学习
多模态虚假新闻检测实战:BERT+ResNet双塔与对比学习

简介:基于PyTorch的多模态虚假新闻检测项目完整代码包,面向自然语言处理与计算机视觉交叉方向的开发者、科研人员及毕业设计选题者,解决社交媒体中文本与图像联合识别虚假新闻的问题。系统以BERT预训练模型提取文本语义特征,以Res… · 2026/9/27 0:00:01

MATLAB雷达信号脉冲压缩仿真:LFM线性调频、匹配滤波与距离分辨率实现
MATLAB雷达信号脉冲压缩仿真:LFM线性调频、匹配滤波与距离分辨率实现

简介:这套Matlab仿真工具完整呈现雷达信号脉冲压缩过程,从线性调频(LFM)信号生成、目标回波仿真到匹配滤波压缩处理均有可运行代码支撑,面向电子信息工程、计算机、数学等专业学生,适用于课程设计、期末大作… · 2026/9/27 0:00:01

汕头网站建设制作厂家避坑指南:5大注意事项救急
汕头网站建设制作厂家避坑指南:5大注意事项救急

汕头网站建设制作厂家避坑指南:5大注意事项救急 改个需求建站公司拖一周,这种憋屈事我见得太多了。 很多汕头老板找本地建站团队,签合同前看着方案挺美,一上线就变脸。 今天不聊虚的,直接拆解找 汕头网站建设制作厂家 时的5个核心 注意事项… · 2026/9/27 0:00:01

多模态虚假新闻检测实战:BERT+ResNet双塔与对比学习
多模态虚假新闻检测实战:BERT+ResNet双塔与对比学习

简介:基于PyTorch的多模态虚假新闻检测项目完整代码包,面向自然语言处理与计算机视觉交叉方向的开发者、科研人员及毕业设计选题者,解决社交媒体中文本与图像联合识别虚假新闻的问题。系统以BERT预训练模型提取文本语义特征,以Res… · 2026/9/27 0:00:01

了解更多?预约专属演示

我们的顾问将为您一对一讲解产品与方案

企业微信二维码