读取文本文件的 Java / Python / Go 三语言指南)
大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载Apache Beam 的TextIO是读写文本文件最常用的 I/O 变换Transform它天然支持以gs://前缀寻址的 Google Cloud StorageGCS对象。本篇以本仓库 Tour of Beam 教学单元 text-io-gcs-read 单元说明 为主线逐一讲解 Java、Python、Go 三种 SDK 从 GCS 读取文本文件的最小可运行示例、底层实现原理与进阶配置并给出可动手验证的 Playground 练习方案。读完本文你将能在自己的 Beam 流水线中无缝读取任意 GCS 文本对象并理解每一行数据在 SDK 内部是如何被切分、解码与发射的。一、核心思路一行文本 PCollection 中的一个元素TextIO在 Apache Beam 中的作用是“按行”读写文本文件读取时它会打开指定文件把每行文本作为一个独立的元素发射到下游PCollectionString写出时则把每个字符串元素写成一行。要从 GCS 读取文本文件只需三步用gs://前缀指定 GCS 对象路径形如gs://bucket-name/object-name调用各语言 SDK 的文本读取 API对返回的PCollection施加下游变换如逐行打印或分词。例如读取名为mybucket存储桶中的myfile.txt路径即为gs://mybucket/myfile.txt。仓库中的三个官方示例Go 示例、Java 示例、Python 示例均读取公开样例对象gs://apache-beam-samples/shakespeare/kinglear.txt《李尔王》全文演示了从 GCS 读取并按单词拆分的完整流程。二、三种 SDK 的最小读取示例2.1 JavaTextIO.read().from(...)Pipeline pipeline Pipeline.create(); pipeline.apply(TextIO.read().from(gs://mybucket/myfile.txt)) .apply(ParDo.of(new DoFnString, Void() { ProcessElement public void processElement(ProcessContext c) { System.out.println(c.element()); } })); pipeline.run();TextIO.read()返回一个Read变换对象.from(gs://...)指定读取来源。仓库的 Task.java 在真实示例中通过PipelineOptionsFactory.create()创建PipelineOptions再以Pipeline.create(options)构建流水线并在ParDo内用c.element().split( )将每行拆成单词逐一输出。TextIO.read()静态工厂方法定义于 sdks/java/core/src/main/java/org/apache/beam/sdk/io/TextIO.java它位于 Beam Java SDK 的核心包中任何 Java 流水线都可直接使用无需额外依赖 GCS 专用 I/O 模块。2.2 Pythonbeam.io.ReadFromText(...)import apache_beam as beam p beam.Pipeline() p | beam.io.ReadFromText(gs://mybucket/myfile.txt) | beam.Map(print) p.run()Python 侧的读取变换名为ReadFromText位于 sdks/python/apache_beam/io/textio.py。仓库的 task.py 给出了带命名标签的等价写法p beam.Pipeline() input p | ReadMyFile beam.io.ReadFromText(gs://apache-beam-samples/shakespeare/kinglear.txt) input | Print words beam.Map(print_words) p.run()print_words对每行调用line.split()后逐个打印单词直观展示了“按行读取 → 元素级处理”的数据流。2.3 Gotextio.Read(...)p, s : beam.NewPipelineWithRoot() lines : textio.Read(p, gs://mybucket/myfile.txt) beam.ParDo(p, func(line string) { fmt.Println(line) }, lines) if err : p.Run(); err ! nil { fmt.Printf(Failed to execute job: %v, err) }Go SDK 的textio.Read接收一个beam.Scope与 glob 路径返回PCollectionstring。仓库 main.go 展示了完整的工程化写法通过空导入_ github.com/apache/beam/sdks/v2/go/pkg/beam/io/filesystem/gcs和_ github.com/apache/beam/sdks/v2/go/pkg/beam/io/filesystem/local完成 GCS 与本地文件系统的注册再调用beamx.Run(context.Background(), p)执行流水线并用log.Exitf处理失败。值得注意的是Go 示例在applyTransform中定义了一个正则wordRE[a-zA-Z]([a-z])?用于从每一行中提取英文单词——这正是 Playground 练习中“分词输出”的基础。三、底层实现原理数据如何从 GCS 流入 PCollection3.1 Python 侧_TextSource文件源Python 的ReadFromText最终实例化_TextSource定义于 textio.py它继承自filebasedsource.FileBasedSource负责把文本文件解析为按换行符分隔的元素。从源码可见_TextSource.__init__接受了compression_type、coder、buffer_size、validate、skip_header_lines、delimiter、escapechar等一整套参数读取记录时还会通过_skip_lines跳过指定数量的表头行。这意味着 Python 的文本读取不止是“读行”还内置了压缩识别、表头跳过、自定义分隔符等文件解析能力。3.2 Go 侧文件系统抽象与压缩选项Go 的textio.Read定义于 sdks/go/pkg/beam/io/textio/textio.go签名如下func Read(s beam.Scope, glob string, opts ...ReadOptionFn) beam.PCollection实现上它会先调用filesystem.ValidateScheme(glob)校验路径协议再把 glob 通过beam.Create注入读取 DoFn。ReadOptionFn是可选的函数式选项源码中提供了三种开箱即用的配置ReadAutoCompression()按文件扩展名自动探测压缩类型默认行为ReadGzip()强制按 gzip 解压ReadUncompressed()强制按未压缩处理。例如显式声明 gzip 压缩的写法为lines : textio.Read(s, gs://mybucket/logs.txt.gz, textio.ReadGzip())这种“默认按扩展名自动判断、可显式覆盖”的设计让同一份代码既能读普通文本也能直接读压缩日志。3.3 Java 侧统一的TextIO.Read变换Java 的TextIO.read().from(...)返回的Read变换同样支持压缩配置withCompression(Compression.UNCOMPRESSED)等并可通过withoutValidation()跳过构建期的路径校验、以withHintMatchesManyFiles()提示框架该 glob 可能匹配大量文件。核心工厂方法read()位于 TextIO.java同一文件还定义了配套的TextIO.write()二者构成 Java 侧文本文件读写的完整入口。四、进阶配置按需定制读取行为4.1 PythonReadFromText完整参数结合 textio.py 源码ReadFromText支持以下常用参数括号内为源码默认值参数默认值作用min_bundle_size0每个读取分片的最小字节数影响并行度compression_typeauto压缩类型auto/gzip/bzip2/deflate/zstd/uncompressedstrip_trailing_newlinesTrue读取时是否去除行尾换行符validateFalse流水线构建期是否校验文件存在文件较多时可能变慢coderStrUtf8Coder()用于解码每行的 Coderskip_header_lines0跳过每个源文件开头的表头行数源码会对负值抛ValueError并对大于 10 的值输出性能告警delimiterNone自定义记录分隔符字节串替代默认换行符escapecharNone转义字符用于处理带转义的记录with_filenameFalse是否同时输出文件名此时元素变为(file_name, record)二元组例如读取带 1 行表头的 CSV 文本对象p | beam.io.ReadFromText( gs://mybucket/events.csv, skip_header_lines1, strip_trailing_newlinesTrue, )4.2 Go 压缩选项与 Go 读取示例的完整骨架Go 侧除压缩选项外textio包还提供ReadAll从一个PCollectionstring中按多个 glob 批量展开读取。若想完整复现仓库示例的工程骨架可参考 main.gofunc main() { beam.Init() p : beam.NewPipeline() s : p.Root() input : textio.Read(s, gs://apache-beam-samples/shakespeare/kinglear.txt) beam.ParDo0(s, func(line string) { for _, word : range wordRE.FindAllString(line, -1) { fmt.Println(word) } }, input) err : beamx.Run(context.Background(), p) if err ! nil { log.Exitf(context.Background(), Failed to execute job: %v, err) } }注意其中的beam.Init()与空导入 GCS 文件系统包二者是 Go 流水线访问 GCS 的必要前提前者初始化 Beam 运行时后者把gs://协议绑定到对应的文件系统实现。五、Playground 实战练习把单词反转后写出本教学单元在 description.md 末尾给出了一个 Playground 练习默认示例会读取文本文件并逐行输出其中的单词请修改示例把这些单词以“反转形式”写入另一个文件。参考解题思路如下以 Python 为例将单词反转并写出import apache_beam as beam def reverse_word(word): return word[::-1] p beam.Pipeline() lines p | ReadFile beam.io.ReadFromText(gs://apache-beam-samples/shakespeare/kinglear.txt) words lines | SplitWords beam.FlatMap(lambda line: line.split()) rev words | ReverseWords beam.Map(reverse_word) rev | WriteFile beam.io.WriteToText(gs://mybucket/reversed_words.txt) p.run()练习要点读取阶段用ReadFromText得到“行”的PCollection用FlatMap或 Go 的beam.ParDo0配合正则、Java 的ParDosplit( )把行摊平成单词对单词逐个做反转变换写出阶段改用WriteToTextGo 为textio.WriteJava 为TextIO.write().to(...)路径同样使用gs://前缀。你可以把输出路径替换为自己的 GCS 存储桶即可在云端观察到反转结果文件。六、注意事项与最佳实践路径前缀读取 GCS 对象必须使用gs://bucket/object形式本地文件则用普通路径Go 中需额外导入 local 文件系统包。仓库示例对apache-beam-samples公共存储桶的读取无需任何鉴权配置即可运行。压缩处理默认按扩展名自动识别压缩格式.gz后缀的文件会自动解压在 Python 中若文件扩展名不规范可显式指定compression_typegzip。性能与并行TextIO读取天然支持分片并行Python 的desired_bundle_size控制目标分片大小文件越大越能体现 Beam 的分布式读取优势Java 读取大量小文件时可考虑withHintMatchesManyFiles()优化提示。校验开关validatePython与withoutValidation()Java控制是否在流水线构建阶段预先检查文件存在性——对运行时才会生成的文件应关闭构建期校验。表头与分隔符结构化文本如 CSV可借助skip_header_lines、delimiter等参数在读取阶段就完成初步解析减少下游变换的负担。通过本文的示例与源码剖析你现在可以在 Java、Python、Go 三种语言中自由读取 GCS 上的文本数据并依据业务需求定制压缩、表头、分隔符等读取行为。更深一步可继续阅读仓库中 textio.py 与 textio.go 的完整源码探索WriteToText、ReadAll等配套变换构建完整的文本读写流水线。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐从 PyTorch 模型到多设备本地推理OpenVINO 快速上手指南从 PyTorch 模型到多设备本地推理OpenVINO 快速上手指南 OpenVINO 是 Intel 开源的 AI 推理优化工具箱负责把你用 PyTor大数据批处理流处理数据工程PyTorch Lightning 在本地On-Prem集群上使用 TorchRun 启动多节点分布式训练PyTorch Lightning 在本地On Prem集群上使用 TorchRun 启动多节点分布式训练 本篇指南以 PyTorch Lightning大数据批处理流处理数据工程Apache Beam Kotlin Kata 实战使用 TextIO 从文本文件读取 PCollectionApache Beam Kotlin Kata 实战使用 TextIO 从文本文件读取 PCollection 创建 Beam 管道时最常见的第一步就是从文大数据批处理流处理数据工程上一篇Ramsey 数下界研究仓库导读用邻接矩阵与 igraph 验证 R(r, s) 新下界下一篇react-admin 软删除系列useSoftDeleteMany 钩子实现批量归档与恢复创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考