恒美微站
首页
关于我们
建站服务
主题模板
案例展示
资讯中心
联系我们
Apache Beam 使用指南:基于 Google Cloud Storage 文件系统(gs://)读写数据
首页
资讯中心
/
Apache Beam 使用指南:基于 Google Cloud Storage 文件系统(gs://)读写数据
Apache Beam 使用指南:基于 Google Cloud Storage 文件系统(gs://)读写数据
发布时间:2026/10/10 11:30:41
批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载Apache Beam 提供了对 Google Cloud StorageGCS对象存储的完整支持既能通过内置 I/O 连接器从 GCS 桶中读取数据、写入数据也能让管道直接与 GCS 文件系统交互如检查文件是否存在、列出目录文件、删除文件。本文基于当前仓库中learning/prompts/documentation-lookup-nolinks/23_io_gcs.md的说明展开并结合 Java、Python、Go 三个 SDK 的源码实现进行纵深解析读完你将对gs://路径、通配符匹配、文件系统 API 以及底层实现原理有完整的实战认知。Google Cloud Storage 与 Apache Beam 的集成概览Google Cloud Storage 是 Google Cloud 提供的对象存储服务用于在云端存储和访问数据。Apache Beam 从两个层面支持 GCS文件系统层面Beam 将gs://作为一种受支持的分布式文件系统接入Java SDK 中对应GcsFileSystemPython 中对应GcsIOGo 中对应pkg/beam/io/filesystem/gcs。这使得任何基于 Beam 文件系统抽象FileSystem接口的读写操作都可以透明地作用在 GCS 对象上。I/O 连接器层面Beam 内置的TextIO连接器原生支持读写 GCS 桶AvroIO、XMLIO、TFRecordIO、ParquetIO等格式连接器同样支持在 GCS 桶内读写各自的文件格式。因此无论是文本日志、Avro 记录、TFRecord 样本还是 Parquet 列存数据只要指定gs://路径Beam 管道即可直接读写无需额外编写定制连接器。gs:// 路径格式与命名规范GCS 对象的访问路径统一使用以下格式gs://bucket/path其中bucket是 GCS 存储桶名称path是桶内对象对象名可视为虚拟目录文件名。例如gs://my-bucket/my-file.txt在 Apache Beam 中所有接受文件路径的参数如TextIO.read().from(...)、TextIO.write().to(...)都可以直接传入gs://路径。以 Java 源码实现为例GcsFileSystem.java 通过getScheme()返回gs来声明自己负责gs://前缀的 URI路径在内部被解析为GcsPathGcsPath.fromUri进而映射到GcsResourceId参与后续的文件系统操作。该文件系统通过 GcsFileSystemRegistrar.java 以AutoService(FileSystemRegistrar.class)方式注册Beam 在解析gs://路径时会自动加载它。提示GCS 中并不存在真正的目录gs://my-bucket/dir/这种带尾部斜杠的路径在 Java 的GcsFileSystem中被视为目录型资源标识源码中matchNewResource会在目录路径后自动补/而文件路径以/结尾则会被拒绝见 GcsFileSystem.java。用内置 I/O 连接器读写 GCSTextIO文本数据读写Beam 内置的TextIO连接器是读写 GCS 文本文件的最直接方式其实现位于 TextIO.java。Java 示例// 从 GCS 读取文本 PCollectionString lines p.apply( ReadFromGcs, TextIO.read().from(gs://my-bucket/my-file.txt)); // 写入 GCS生成 gs://my-bucket/output-00000-of-00001 等分片文件 p.apply(WriteToGcs, TextIO.write().to(gs://my-bucket/output));Python 示例import apache_beam as beam with beam.Pipeline() as p: # 读取 lines p | ReadFromGcs beam.io.ReadFromText(gs://my-bucket/my-file.txt) # 写入 lines | WriteToGcs beam.io.WriteToText(gs://my-bucket/output)Go 示例import ( github.com/apache/beam/sdks/v2/go/pkg/beam github.com/apache/beam/sdks/v2/go/pkg/beam/io/textio ) func main() { beam.Init() p : beam.NewPipeline() s : p.Root() lines : textio.Read(s, gs://my-bucket/my-file.txt) textio.Write(s, gs://my-bucket/output, lines) }读取与写入的具体字节流操作最终落到各 SDK 的文件系统实现上Java 端由GcsFileSystem.open(...)/create(...)调用GcsUtil完成通道创建GcsFileSystem.javaGo 端由gcs包实现OpenRead/OpenWrite见 gcs.goPython 端由GcsIO.open(...)提供流式读写见 gcsio.py。其他格式连接器AvroIO、XMLIO、TFRecordIO、ParquetIO除了TextIOBeam 的AvroIO、XMLIO、TFRecordIO、ParquetIO连接器同样支持从 GCS 桶读取数据、向 GCS 桶写入不同文件格式的数据连接器适用格式典型场景AvroIOAvro 序列化记录大数据批处理、跨系统数据交换XMLIOXML 文档解析 XML 事件流TFRecordIOTFRecord 样本机器学习训练样本读写ParquetIOParquet 列式存储分析型查询、与数据仓库互操作这些连接器的使用方式与TextIO一致只需将文件路径替换为gs://形式例如// 读取 GCS 中的 Parquet 文件 PCollectionGenericRecord records p.apply( ReadParquet, ParquetIO.read(ParquetIO.ReadFiles.class) .from(gs://my-bucket/records/*.parquet));由于它们均构建在 Beam 统一的文件系统抽象之上只要路径以gs://开头就会被自动路由到 GCS 文件系统实现因此在 GCS 桶内读写不同文件格式与在本地/其他文件系统上的代码完全一致仅路径前缀不同。通配符Wildcard读写多个文件Beam 的读写变换支持在gs://路径中使用通配符从而一次处理多个文件读或一次写入多个分片写。例如gs://my-bucket/my-files-*.txt可以匹配该桶内所有文件名符合该模式的对象用于读取或写入读取TextIO.read().from(gs://my-bucket/my-files-*.txt)会展开匹配到的全部文件并合并为一条PCollection写入TextIO.write().to(gs://my-bucket/my-files-*.txt)时*会被替换为实际生成的分片编号如my-files-00000-of-00002.txt。在 Java 的GcsFileSystem中通配符展开由match(ListString specs)负责路径先按GcsUtil.isWildcard(path)区分为 glob 与非 glob 两类GcsUtil.javaglob 路径通过expand(gcsPattern)处理先由GcsUtil.getNonWildcardPrefix提取通配符前的静态前缀GcsUtil.java调用 GCS List API 分页列出该前缀下的对象再用正则wildcardToRegexp(...)对每个对象名做精确过滤最终返回匹配集合GcsFileSystem.java。Go SDK 采用类似策略gcs包的List(ctx, glob)先用前缀列出候选对象再在内存中做 glob 匹配见 gcs.go。值得注意的是Beam 文件系统对通配符的支持以尽可能将*之前的字符作为前缀以缩小 List 范围为优化原则因此在设计对象命名时将可枚举的静态部分放在通配符之前如gs://bucket/date2024-*/data-*.avro可以获得更好的性能。直接操作 GCS 文件系统除了通过 I/O 连接器读写数据Apache Beam 还允许管道直接与 GCS 文件系统交互例如验证某个文件是否存在获取目录下的文件列表删除文件。Java 中可通过FileSystems工具类获得文件系统实例并执行上述操作底层即GcsFileSystem。从源码可以看出GcsFileSystem完整实现了 BeamFileSystemGcsResourceId抽象的所有核心能力匹配/列举match返回MatchResult含文件大小、MD5 校验和、最后修改时间等元数据GcsFileSystem.java打开/创建open/create后者还支持设置 MIME 类型、预期文件不存在校验以及上传缓冲区大小GcsCreateOptionsGcsFileSystem.java复制/重命名copy/rename并可通过GcsOptions.getGcsPerformanceMetrics()开启复制/重命名次数与耗时的指标统计GcsFileSystem.java删除delete批量删除对象GcsFileSystem.java。Python 的GcsIOgcsio.py同样提供了一组等价的能力包括open读/写流、delete、copy、rename、exists判断对象是否存在与size获取对象大小可直接用于管道内的文件操作。Go SDK 的 GCS 文件系统实现位于 gcs.go通过init()中的filesystem.Register(gs, New)注册gs://schemegcs.go并实现了List、OpenRead、OpenWrite、Size、LastModified、Remove、Copy等接口。其客户端初始化优先使用应用默认凭据storage.ScopeReadWrite失败时回退为匿名访问gcs.go。语言支持与适用前提GCS 文件系统在 Java、Python、Go 三个 Beam SDK 中均受支持与本仓库中的实现一一对应Javasdks/java/extensions/google-cloud-platform-core模块中的GcsFileSystem及其注册器Pythonsdks/python/apache_beam/io/gcp/gcsio.py中的GcsIOGosdks/go/pkg/beam/io/filesystem/gcs/gcs.go中的 gcs 文件系统。需要说明的前提与限制依赖与凭据使用 GCS 需要安装相应的 GCP 依赖Java 依赖google-cloud-platform-corePython 需安装apache-beam[gcp]或对应的 GCS 客户端库并配置可用的 Google Cloud 凭据如服务账号、应用默认凭据Go 在无法取得凭据时会尝试匿名访问公开对象。路径语法所有读写路径必须符合gs://bucket/path格式本地路径、其他 scheme 不会路由到 GCS 文件系统。分片写入语义write().to(gs://...)通常生成多个分片文件输出文件名的具体形态取决于连接器与 Runner 的实现通配符仅用于指定输出前缀/模式。小结Apache Beam 对 Google Cloud Storage 的支持是文件系统抽象 格式连接器双层架构的自然结果gs://被注册为受支持的 schemeTextIO、AvroIO、XMLIO、TFRecordIO、ParquetIO等连接器在此基础上提供格式化的读写能力通配符机制让多文件批处理与分片写入变得透明而 Java、Python、Go 三端对齐的文件系统 API匹配、列举、存在性检查、删除、复制、重命名则保证了在管道内直接管理 GCS 对象的可操作性。设计数据管道时只需遵循gs://bucket/path路径规范并把静态前缀放在通配符之前即可高效、可靠地在云端对象存储上运行 Beam 作业。赞分享批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载相关推荐Apache Beam 与 Google Cloud StorageGCS文件系统集成实战指南Apache Beam 与 Google Cloud StorageGCS文件系统集成实战指南 Apache Beam 提供统一的批处理和流处理编程模型GApache Beam 中 Google Cloud StorageGCS文件系统与对象存储 I/O 实战指南Apache Beam 中 Google Cloud StorageGCS文件系统与对象存储 I/O 实战指南 Apache Beam 是一套面向批处理与流大数据批处理流处理数据工程Apache Beam TextIO 实战从 Google Cloud StorageGCS读取文本文件的 Java / Python / Go 三语言指南Apache Beam TextIO 实战从 Google Cloud StorageGCS读取文本文件的 Java / Python / Go 三语言指大数据批处理流处理数据工程上一篇wrk与消息队列集成异步系统性能测试方案下一篇10个顶级Android UI库提升开发效率的终极指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考