尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Apache Beam Java 实战:使用 PubsubIO 将数据写入 Google Pub/Sub Topic

Apache Beam Java 实战:使用 PubsubIO 将数据写入 Google Pub/Sub Topic 大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载本文以仓库 learning/prompts/code-generation/java/02_io_pubsub.md 中的代码生成示例为骨架完整讲解如何使用 Apache Beam 的PubsubIO连接器把数据写入 Google Pub/Sub Topic。你将掌握一段可直接复用的WritePubSubTopic管道代码、PipelineOptions 命令行参数模式的写法以及PubsubIO.Write在源码层面的实现原理、Topic 命名规范与批量写入、时间戳、消息 ID、错误处理等高级配置从而具备在生产管道中安全落地的实战能力。一、核心思路用 PubsubIO 连接器写 Pub/SubGoogle Cloud Pub/Sub 是 Google Cloud 上的消息发布/订阅服务Apache Beam 通过org.apache.beam.sdk.io.gcp.pubsub.PubsubIO连接器为它提供读写PTransform。原文档给出的答案是在管道中先用Create.of构造一个内存中的PCollection再通过PubsubIO.writeStrings().to(topic)把每条字符串作为 UTF-8 消息发布到指定 Topic。写 Pub/Sub 会产生无界数据流因此该写入端天然适配流式Streaming管道也同样可以在批量Batch管道中作为输出端使用。二、完整示例WritePubSubTopic原文档提供了如下可直接运行的完整代码这里原样继承并补充了类说明注释package pubsub; import org.apache.beam.sdk.Pipeline; import org.apache.beam.sdk.io.gcp.pubsub.PubsubIO; import org.apache.beam.sdk.options.Description; import org.apache.beam.sdk.options.PipelineOptions; import org.apache.beam.sdk.options.PipelineOptionsFactory; import org.apache.beam.sdk.transforms.Create; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.util.Arrays; import java.util.List; // Pipeline to write data to a Google Pub/Sub topic public class WritePubSubTopic { private static final Logger LOG LoggerFactory.getLogger(WritePubSubTopic.class); // Pipeline options to configure the pipeline public interface WritePubSubTopicOptions extends PipelineOptions { Description(PubSub topic name to write to) String getTopicName(); void setTopicName(String value); } public static void main(String[] args) { // Parse the pipeline options from the command line WritePubSubTopicOptions options PipelineOptionsFactory.fromArgs(args).withValidation().as(WritePubSubTopicOptions.class); // Create a pipeline Pipeline p Pipeline.create(options); // Sample messages to write to the Pub/Sub topic final ListString messages Arrays.asList( PubSub message 1, PubSub message 2, PubSub message 3 ); p // Create a list of messages to write to the Pub/Sub topic .apply(Create.of(messages)) // Write the messages to the Pub/Sub topic .apply(PubsubIO.writeStrings().to(options.getTopicName())); // Execute the pipeline p.run(); } }这段代码体现了原文档强调的要点利用 Pipeline Options 模式从命令行解析参数。运行时通过--topicName...即可指定目标 Topic无需改动代码。三、代码逐段拆解3.1 用 PipelineOptions 定义命令行参数WritePubSubTopicOptions继承自PipelineOptions其中Description(PubSub topic name to write to)注解用于描述该选项的用途getTopicName()/setTopicName()的 JavaBean 命名约定决定了命令行参数名是--topicName。PipelineOptionsFactory.fromArgs(args).withValidation().as(WritePubSubTopicOptions.class)完成三件事fromArgs(args)解析命令行传入的--keyvalue参数withValidation()开启选项校验缺失必填项或参数非法时会在启动阶段报错.as(...)将解析结果转换为自定义选项类型。3.2 创建管道Pipeline.create(options)使用上一步解析出的选项创建管道实例。执行器Runner同样通过命令行传入例如--runnerDirectRunner或--runnerDataflowRunner未指定时使用本地 Direct Runner。3.3 构造输入数据final ListString messages Arrays.asList( PubSub message 1, PubSub message 2, PubSub message 3 ); p.apply(Create.of(messages))Create.of将内存中的ListString变成一个有界PCollectionString用于模拟管道上游产生数据的来源。在生产场景中这里通常替换为KafkaIO、PubsubIO.readStrings()或批处理文件源等真实数据源。3.4 写入 Pub/Sub Topic.apply(PubsubIO.writeStrings().to(options.getTopicName()));PubsubIO.writeStrings()返回WriteString写入端.to(topic)指定目标 Topic。由于 Topic 名称来自命令行参数运行时可以灵活切换目标。3.5 启动管道p.run();p.run()提交管道执行。对写 Pub/Sub 这类持续写入场景流式管道会长期运行并持续发布消息。四、源码视角writeStrings() 背后发生了什么仓库中 PubsubIO.java 是连接器的核心实现。writeStrings()的源码位于 PubsubIO.java#L794-L801public static WriteString writeStrings() { return Write.newBuilder( (ValueInSingleWindowString stringAndWindow) - new PubsubMessage( stringAndWindow.getValue().getBytes(StandardCharsets.UTF_8), ImmutableMap.of())) .setDynamicDestinations(false) .build(); }从中可以看出两点实现事实编码方式每条字符串被getBytes(StandardCharsets.UTF_8)编码为字节数组作为PubsubMessage的 payload属性为空写入时的消息属性attributes为空 map即writeStrings()不带自定义属性。PubsubIO类中还定义了默认传输客户端 PubsubIO.java#L196private static final PubsubClient.PubsubClientFactory FACTORY PubsubJsonClient.FACTORY;即默认使用PubsubJsonClientJSON over HTTP 实现与 Pub/Sub 服务通信如需切换到 gRPC 客户端可调用Write.withClientFactory(PubsubGrpcClientFactory)见 PubsubIO.java#L1645-L1647。五、Topic 路径格式与命名规则.to()接收的字符串必须是完整的 Pub/Sub 资源路径。源码通过正则校验路径格式见 PubsubIO.java#L209private static final Pattern TOPIC_REGEXP Pattern.compile(projects/([^/])/topics/(.));即标准格式为projects/project_id/topics/topic_name例如projects/my-project/topics/my-topic。此外源码对 project 与 topic 名称还有额外约束项目 ID 需匹配[a-z][-a-z0-9:.]{4,61}[a-z0-9]PubsubIO.java#L203-L204名称长度必须在 3255 字符之间且不能以goog开头PubsubIO.java#L238-L258 的validatePubsubName。这些校验逻辑同样被单元测试覆盖例如 PubsubIOTest.java#L993-L998 验证了合法 Topic 名含大写、数字、-_.~%等字符可以通过而 PubsubIOTest.java#L1007 验证了含*的非法名称会被拒绝。注意.to(String)在管道构建阶段就调用validateTopic做格式校验PubsubIO.java#L1603-L1607所以路径写错会在本地立即抛IllegalArgumentException而不是等到远端执行时才失败。六、不止字符串更多写入 API除了writeStrings()连接器还提供多种写入入口均定义在 PubsubIO.javaAPI源码位置说明writeMessages()L768-L774写入PubsubMessage可携带自定义属性writeMessagesDynamic()L782-L788每条消息按PubsubMessage.withTopic(...)写入不同 TopicwriteStrings()L794-L801写入 UTF-8 字符串writeProtos(Class)L807-L813写入 protobuf 二进制消息writeAvros(Class)L833-L839写入 Avro 二进制消息6.1 携带属性的写入当需要携带 key-value 属性时使用PubsubMessagep.apply(Create.of(event-1, event-2)) .apply(MapElements.via((String s) - new PubsubMessage( s.getBytes(StandardCharsets.UTF_8), ImmutableMap.of(type, user_event, source, mobile)))) .apply(PubsubIO.writeMessages().to(options.getTopicName()));6.2 动态 Topic 路由如果消息需要按内容路由到不同 Topic有两种方式源码注释见 PubsubIO.java#L161-L181方式一为每个元素动态计算目标 TopicPubsubIO.java#L1631-L1637events.apply(PubsubIO.writeAvros(MyType.class) .to((ValueInSingleWindowEvent quote) - { String country quote.getCountry(); return projects/myproject/topics/events_ country; }));方式二把 Topic 直接写在PubsubMessage上再使用writeMessagesDynamic()events.apply(MapElements.into(new TypeDescriptorPubsubMessage() {}) .via(e - new PubsubMessage( e.toByteString(), Collections.emptyMap()).withTopic(e.getCountry()))) .apply(PubsubIO.writeMessagesDynamic());七、Write 端高级配置PubsubIO.Write提供链式配置方法均可在.to(...)之后继续调用7.1 时间戳属性withTimestampAttributePubsubIO.java#L1699-L1701.writeStrings().to(topic) .withTimestampAttribute(ts)把每条记录的时间戳以毫秒数Unix epoch写入名为ts的属性下游若用PubsubIO.readStrings().withTimestampAttribute(ts)读取读端实现见 PubsubIO.java#L1106-L1108就能恢复业务时间戳而不是默认的 Pub/Sub 发布时间。7.2 消息 ID 属性withIdAttributePubsubIO.java#L1711-L1713.writeStrings().to(topic).withIdAttribute(id)为每条消息附加不透明唯一 ID下游配合读端的withIdAttribute(id)PubsubIO.java#L1119-L1121可用于去重——因为 Pub/Sub 无法保证不重复投递。7.3 批量写入控制写入端默认是批量的相关方法源码位于 PubsubIO.java#L1660-L1670withMaxBatchSize(int)攒够指定条数再发送withMaxBatchBytesSize(int)按字节数控制批量大小。之所以需要控制批量大小是因为 Pub/Sub 单次请求有 10MB 上限源码常量 PubsubIO.java#L218 明确写为PUBSUB_MESSAGE_MAX_TOTAL_SIZE 10_000_000。7.4 消息排序键withOrderingKeyPubsubIO.java#L1685-L1687 的withOrderingKey()用于启用消息排序键发布。注意源码同时说明 Beam 向服务端发布记录的顺序仍是不确定的排序键只保证同一键下消息按服务端接收顺序排列。7.5 序列化失败处理withErrorHandlerPubsubIO.java#L1724 提供withErrorHandler(ErrorHandlerBadRecord, ?)将序列化失败的记录路由到自定义错误处理逻辑schema 错误不在其管辖范围仍按 runner 默认行为处理。7.6 本地模拟器与自定义客户端源码类注释PubsubIO.java#L112-L115说明如需使用本地 Pub/Sub 模拟器可通过PubsubOptions#setPubsubRootUrl(String)设置模拟器的 host 与 port写入端也有对应的withPubsubRootUrl(String)PubsubIO.java#L1715-L1717。需要 gRPC 传输时可改用withClientFactory(new PubsubGrpcClientFactory())。八、运行方式与验证8.1 运行参数编译打包后通过命令行传入参数运行Topic 名与 Runner 都是可配置项java -cp classpath pubsub.WritePubSubTopic \ --topicNameprojects/my-project/topics/my-topic \ --runnerDirectRunner其中--topicName目标 Topic 完整路径对应WritePubSubTopicOptions的getTopicName--runner执行器如DirectRunner本地调试、DataflowRunnerGoogle Cloud Dataflow等。8.2 权限要求源码注释PubsubIO.java#L117-L121明确指出权限要求取决于运行管道的PipelineRunner。例如在 Dataflow Runner 上运行时需要为 Dataflow 服务账号授予相应 Pub/Sub 权限本地 Direct Runner 则使用本机凭据。8.3 测试佐证仓库中的 PubsubIOTest.java 覆盖了写入端的多种行为testWriteDisplayDataL288-L302验证 Topic、timestampAttribute、idAttribute会正确登记到管道 DisplayData便于运维观测若干用例L993-L1013验证 Topic 名称格式校验另有针对writeMessages().to(...)指向不存在 Topic 的写入场景测试。九、注意事项小结Topic 路径必须完整使用projects/project_id/topics/topic_name格式否则本地构建阶段即报错单批请求有 10MB 上限超大消息需注意withMaxBatchSize/withMaxBatchBytesSize的配合字符串消息编码为 UTF-8writeStrings()不携带属性需要属性请改用writeMessages()时间戳与 ID 传递跨管道传递业务时间戳与去重 ID 时写入端与读取端的withTimestampAttribute/withIdAttribute必须配套使用凭据与权限按所选 Runner 提前配置 Pub/Sub 访问权限本地开发可借助 Pub/Sub 模拟器通过PubsubOptions#setPubsubRootUrl指向本地服务在无云环境的情况下联调管道。通过上述内容你已经掌握了从内存数据 →PubsubIO.writeStrings()→ 指定 Topic的完整写入链路以及连接器在源码层面的编码、校验、批量与配置机制可直接在此基础上构建自己的 Pub/Sub 数据出口管道。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐Apache Beam PubsubIO 实战从 Google Cloud Pub/Sub 流式读取数据的完整代码解读与源码剖析Apache Beam PubsubIO 实战从 Google Cloud Pub/Sub 流式读取数据的完整代码解读与源码剖析 本指南以仓库 learnin大数据批处理流处理数据工程Apache Beam KafkaToPubsub 示例实战从 Apache Kafka 到 Google Cloud Pub/Sub 的流式数据接入Apache Beam KafkaToPubsub 示例实战从 Apache Kafka 到 Google Cloud Pub/Sub 的流式数据接入 本文围大数据批处理流处理数据工程Apache Beam 实战使用 BigQueryIO 向 Google BigQuery 写入数据的 Java 指南Apache Beam 实战使用 BigQueryIO 向 Google BigQuery 写入数据的 Java 指南 导读 本文以 Apache Beam大数据批处理流处理数据工程上一篇PDFdir基于正则表达式的PDF智能书签生成解决方案下一篇智能农业灌溉系统用Arduino-ESP32打造物联网节水神器创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表