Apache Beam PubsubIO 实战:从 Google Cloud Pub/Sub 流式读取数据的完整代码解读与源码剖析 大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载本指南以仓库learning/prompts/code-explanation/java/02_io_pubsub.md中讲解的ReadPubSubTopic示例代码为核心逐行拆解 Apache Beam 中通过PubsubIO.readStrings().fromTopic(...)从 Google Cloud Pub/Sub 主题读取流式数据的完整链路。读完本文你将掌握如何通过自定义PipelineOptions接收命令行参数、PubsubIO各读取方法的底层实现与适用场景、fromTopic与fromSubscription的区别以及时间戳属性、去重属性、死信主题等生产级配置项的真实用法。一、示例代码概览一段最小可用的 Pub/Sub 流式读取管道关联文档给出的核心示例是一段结构完整的最小管道它包含了一个 Beam 流式作业的四个典型组成部分选项接口定义输入参数→ 管道创建 → 数据源读取Source→ 逐元素处理Transform→ 提交执行。public class ReadPubSubTopic { private static final Logger LOG LoggerFactory.getLogger(ReadPubSubTopic.class); public interface ReadPubSubTopicOptions extends PipelineOptions { Description(Pub/Sub Topic to read from) Default.String(projects/pubsub-public-data/topics/taxirides-realtime) String getTopicName(); void setTopicName(String value); } public static void main(String[] args) { ReadPubSubTopicOptions options PipelineOptionsFactory.fromArgs(args).withValidation().as(ReadPubSubTopicOptions.class); Pipeline p Pipeline.create(options); p .apply(Read from Pub/Sub, PubsubIO.readStrings().fromTopic(options.getTopicName())) .apply(Process elements, ParDo.of(new DoFnString, String() { ProcessElement public void processElement(ProcessContext c) { c.output(c.element()); } }) ); p.run(); } }这段代码使用 Apache Beam 的PubsubIO从 Pub/Sub 主题读取数据返回的PCollectionString是一个无界数据集unbounded PCollection因为 Pub/Sub 是一个持续到达的数据流读取 transform 会不间断地消费新消息。二、PipelineOptions用接口声明管道运行参数示例中的ReadPubSubTopicOptions接口继承了PipelineOptions用于声明这条管道运行时可以接收的配置项。其关键设计如下public interface ReadPubSubTopicOptions extends PipelineOptions { Description(Pub/Sub Topic to read from) Default.String(projects/pubsub-public-data/topics/taxirides-realtime) String getTopicName(); void setTopicName(String value); }getTopicName()/setTopicName(String)Beam 的PipelineOptions使用 JavaBean 风格的方法对getter/setter定义选项接口编译时由PipelineOptionsFactory自动生成实现类。Description为该选项提供人类可读的描述会在--help输出中展示方便命令行使用者理解参数含义。Default.String声明该选项的默认值。此处默认指向 Google 公开的纽约出租车实时数据流主题projects/pubsub-public-data/topics/taxirides-realtime这意味着不传参数直接运行也能工作。命令行参数运行管道时通过--topicNameprojects/xxx/topics/yyy即可覆盖该默认值。这是由PipelineOptionsFactory.fromArgs(args)完成的参数解析机制。ReadPubSubTopicOptions options PipelineOptionsFactory.fromArgs(args).withValidation().as(ReadPubSubTopicOptions.class); Pipeline p Pipeline.create(options);PipelineOptionsFactory.fromArgs(args)解析 main 方法接收的命令行参数包括--runner、--topicName等。.withValidation()开启选项校验例如检查必填参数是否缺失。.as(ReadPubSubTopicOptions.class)将解析结果转换为自定义选项接口的实例。Pipeline.create(options)基于这些选项创建管道对象p。从仓库源码看选项体系中还有更丰富的配套接口例如PubsubOptions提供setPubsubRootUrl(String)方法可用于指向本地 Pub/Sub 模拟器的 host 与端口见 PubsubIO.java 源码注释便于离线开发与测试。三、数据源PubsubIO.readStrings().fromTopic(...)的底层实现3.1readStrings()做了什么示例使用PubsubIO.readStrings()读取消息。从源码实现PubsubIO.java可以看到它的本质/** * Returns A {link PTransform} that continuously reads UTF-8 encoded strings from a Google Cloud * Pub/Sub stream. */ public static ReadString readStrings() { return Read.newBuilder( (PubsubMessage message) - new String(message.getPayload(), StandardCharsets.UTF_8)) .setCoder(StringUtf8Coder.of()) .build(); }也就是说readStrings()会为每条 Pub/Sub 消息做两件事把消息的二进制payload按UTF-8解码为 JavaString为输出PCollectionString绑定StringUtf8Coder保证元素在 runner 之间传输时的序列化一致性。3.2fromTopic(...)与主题路径校验fromTopic有两种重载直接传字符串或传ValueProviderString后者支持运行时求值。其核心逻辑PubsubIO.javapublic ReadT fromTopic(String topic) { return fromTopic(StaticValueProvider.of(topic)); } public ReadT fromTopic(ValueProviderString topic) { validateTopic(topic); return toBuilder() .setTopicProvider(NestedValueProvider.of(topic, PubsubTopic::fromPath)) .build(); }其中PubsubTopic::fromPath负责把形如projects/my-project/topics/my-topic的字符串解析并校验为合法的主题路径。测试用例也印证了主题名校验规则例如合法的projects/my-project/topics/AbC-DeF、projects/my-project/topics/AbC-1234-_.~%-_.~%-_.~%-abc可通过而带*通配符的非法名称会被拒绝见 PubsubIOTest.java。一个重要语义源码注释明确说明fromTopic读取时runner 只会消费管道启动之后发布到该主题的数据管道启动前积压的历史消息不会被读取PubsubIO.java。这决定了fromTopic天然适用于实时流处理场景若需要消费历史存量数据应改用订阅subscription方式。3.3fromTopic与fromSubscription如何选择PubsubIO.Read同时支持fromTopic(...)与fromSubscription(...)PubsubIO.java两者互斥public ReadT fromSubscription(String subscription) { return fromSubscription(StaticValueProvider.of(subscription)); }源码注释给出了两者差异PubsubIO.java从fromTopic读取时管道启动时会自动创建订阅读取自启动时刻起的实时数据从fromSubscription读取时多个 reader 共享同一订阅会各自分到任意一部分数据因此多个消费者通常应使用各自独立的订阅。在expand校验阶段如果既没设置 topic 也没设置 subscription或两者同时设置会直接抛出IllegalStateException见 PubsubIO.java。四、逐元素处理ParDo与DoFn的匿名实现读取得到的PCollectionString接着被ParDo变换逐元素处理.apply(Process elements, ParDo.of(new DoFnString, String() { ProcessElement public void processElement(ProcessContext c) { c.output(c.element()); } }) );ParDo.of(...)将自定义的DoFn包装成并行处理变换ProcessElement注解标记的方法会在每个元素到达时被调用ProcessContext c携带当前元素此处c.output(c.element())是一个原样透传的恒等处理——把输入元素直接输出实际业务中通常会把c.element()替换为解析、清洗、聚合等真实逻辑。ParDo是 Beam 中最核心的通用处理原语本示例用它演示了读取 → 处理 → 输出的标准管道形态。五、提交执行p.run()p.run();run()将管道提交给指定的 runner 执行。示例默认没有显式指定 runner实际运行时通过--runnerDirectRunner本地测试或--runnerDataflowRunnerGoogle Cloud Dataflow 云端执行等参数选择。由于 Pub/Sub 读取产生的是无界PCollection这条管道会持续运行直到被显式取消。六、从示例走向生产PubsubIO.Read的完整配置面关联文档只覆盖了readStrings().fromTopic(...)这一条最简路径。仓库中PubsubIO.Read还提供了大量生产级配置理解它们能让你在真实场景中直接复用6.1 按消息类型选择读取方法PubsubIO为不同负载类型预置了多种读取入口PubsubIO.java方法返回元素类型适用场景readStrings()StringUTF-8 文本消息本示例readMessages()PubsubMessage仅 payload二进制负载关注原始消息readMessagesWithMessageId()PubsubMessage messageId需要消息 IDreadMessagesWithAttributes()PubsubMessagepayload attributes需要自定义属性做路由或业务字段readMessagesWithAttributesAndMessageId()payload attributes messageId同时需要属性与消息 IDreadMessagesWithAttributesAndMessageIdAndOrderingKey()全字段使用消息排序键ordering key的场景readProtos(ClassT)protobuf 消息负载为 protobuf 编码readProtoDynamicMessages(ProtoDomain, String)DynamicMessage编译期未知消息类型的 protobufreadAvros(ClassT)Avro 记录负载为 Avro 编码6.2 高级读取配置项以下方法均定义于PubsubIO.ReadPubsubIO.javawithTimestampAttribute(String)指定消息属性中承载事件时间的字段。属性值可以是 Unix 毫秒数或 RFC 3339 格式字符串如2015-10-29T23:41:41.123Z。不设置时默认使用 Pub/Sub 的发布时间作为事件时间所有窗口计算都基于该时间戳且按到达时间分配时间戳时系统可保证不会出现迟到数据PubsubIO.java。仓库示例 GameStats.java 就使用了withTimestampAttribute(GameConstants.TIMESTAMP_ATTRIBUTE)配合fromTopic(...)读取游戏事件流。withIdAttribute(String)指定承载唯一记录 ID 的属性。Pub/Sub 不保证不重复投递设置该属性后 Beam 可据此对重复消息做尽力去重不设置则完全无法保证去重PubsubIO.java。withDeadLetterTopic(String)把读取阶段解析失败的消息写入死信主题死信消息会附带exceptionClassName、exceptionMessage、pubsubMessageId三个属性便于排查与withErrorHandler互斥。注意读取成功之后业务处理阶段产生的错误不在此列需要自行配置Write或依赖 Pub/Sub 侧的死信设置PubsubIO.java。withErrorHandler(ErrorHandlerBadRecord, ?)为解析失败的消息注册 Beam 的错误处理机制与死信主题二选一。withClientFactory(PubsubClient.PubsubClientFactory)默认使用PubsubJsonClient可替换为PubsubGrpcClientFactory切换底层通信客户端PubsubIO.java。withValidation()显式开启读取端的参数校验。组合示例可替换示例中的读取行p.apply(Read from Pub/Sub, PubsubIO.readStrings() .withTimestampAttribute(event_time) .withIdAttribute(record_id) .fromTopic(options.getTopicName()));七、运行与验证7.1 本地运行命令以 DirectRunner 在本地跑通这条管道并指定自定义主题mvn compile exec:java -Dexec.mainClassReadPubSubTopic \ -Dexec.args--runnerDirectRunner --topicNameprojects/my-project/topics/my-topic其中--topicName就是ReadPubSubTopicOptions.getTopicName()对应的命令行参数名不传该参数时会使用Default.String指定的纽约出租车公开主题。7.2 前提与限制凭证访问 Pub/Sub 需要有效的 Google Cloud 凭证实际权限要求取决于所选 runnerPubsubIO.java 权限说明。本地模拟器可通过PubsubOptions#setPubsubRootUrl指向本地模拟器离线开发时无需真实 GCP 资源PubsubIO.java。7.3 源码级测试佐证仓库测试 PubsubIOTest.java 覆盖了本节涉及的核心行为readStrings().fromTopic(projects/myproject/topics/mytopic)的合法主题解析、StaticValueProvider包装的 topic/subscription 读取、以及与withIdAttribute(myId)组合使用的场景可作为你编写集成测试时的参考模板。八、小结ReadPubSubTopic示例虽短却完整演示了 Beam 流式管道的骨架PipelineOptions声明参数 →PipelineOptionsFactory解析命令行 →Pipeline.create建管道 →PubsubIO.readStrings().fromTopic建立无界数据源 →ParDo逐元素处理 →p.run()提交。在真实项目中你可以在此基础上按需叠加withTimestampAttribute事件时间与窗口、withIdAttribute去重、withDeadLetterTopic/withErrorHandler失败治理等配置或切换到readMessagesWithAttributes()、readProtos()等读取入口以满足不同的消息格式与业务诉求。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐Apache Beam KafkaToPubsub 示例实战从 Apache Kafka 到 Google Cloud Pub/Sub 的流式数据接入Apache Beam KafkaToPubsub 示例实战从 Apache Kafka 到 Google Cloud Pub/Sub 的流式数据接入 本文围大数据批处理流处理数据工程Apache Beam 中使用 BigQueryIO 读取 BigQuery 表ReadBigQueryTable 完整实战解析Apache Beam 中使用 BigQueryIO 读取 BigQuery 表ReadBigQueryTable 完整实战解析 Apache Beam 的大数据批处理流处理数据工程shadPS4 安装与配置完整指南Windows / Linux / macOS 三平台跑通 PS4 模拟器shadPS4 安装与配置完整指南Windows / Linux / macOS 三平台跑通 PS4 模拟器 shadPS4 是一个用 C 写的 Play大数据批处理流处理数据工程上一篇AI 高管人才战争全景解读OpenAI 权力更迭、天价薪酬与人才流动图谱附 RAG 知识库实战下一篇I-wanna-clean-keyboard 使用教程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考