尧图网络科技YAOTU DIGITAL 获取报价
获取报价
首页 / 资讯中心 / 文章详情

Apache Beam KafkaToPubsub 示例深入解析:从 Apache Kafka 到 Google Cloud Pub/Sub 的流式数据管道

发布时间:2026/9/29 2:47:33

资讯中心
01
ARTICLE

Apache Beam KafkaToPubsub 示例深入解析:从 Apache Kafka 到 Google Cloud Pub/Sub 的流式数据管道

Apache Beam KafkaToPubsub 示例深入解析:从 Apache Kafka 到 Google Cloud Pub/Sub 的流式数据管道
【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载本文基于 Apache Beam 官方示例 kafkatopubsub 展开它构建了一条从 Kafka 主题读取消息、写入指定 Pub/Sub 主题的流式管道并完整支持 SASL/SCRAM 认证、HashiCorp Vault 凭据托管、SSL 加密连接以及 PUBSUB / AVRO 两种输出格式。读完本文你将掌握该示例的完整运行方式、全部命令行参数的含义与底层实现并能够在此基础上定制自己的 Schema 与序列化逻辑搭建一套可落地的 Kafka 到 Pub/Sub 数据接入方案。示例定位与能力总览KafkaToPubsub是 Apache Beam 提供的一个完整的流式管道示例其核心职责是从一个或多个 Kafka 主题持续读取消息并把消息写入 Google Cloud Pub/Sub 的单个输出主题。它位于 examples/java/src/main/java/org/apache/beam/examples/complete/kafkatopubsub 目录下可作为真实生产管道的参考模板也可以直接扩展复用。该示例支持的数据格式可序列化的明文格式例如 JSONPubSubMessageGoogle Cloud Pub/Sub 的消息对象。支持的输入源配置单个或多个 Kafka bootstrap server基于明文或 SSL 加密连接的 Kafka SASL/SCRAM 认证通过密钥托管服务 HashiCorp Vault 读取 Kafka 凭据。支持的输出配置单个 Google Cloud Pub/Sub topic。在简单场景下示例会创建一个读取指定 Kafka 源主题的 Beam 管道并将文本消息流式写入目标 Pub/Sub 主题在复杂场景下Kafka 侧可能需要 SASL/SCRAM 认证认证既可以走明文连接也可以走 SSL 加密连接。示例支持使用单个 Kafka 用户账号对提供的所有源服务器与主题完成认证若需要基于 SSL 的 SASL 认证则必须提供 SSL 证书位置并配置可访问的密钥保管库当前支持 HashiCorp Vault以获取 Kafka 用户名与密码。管道整体架构与执行流程主类 KafkaToPubsub.java 的入口逻辑非常清晰其运行流程分为三个阶段读取配置与认证信息解析命令行参数为KafkaToPubsubOptions若有 Vault 地址与 Token则从 Vault 拉取 Kafka 凭据并构造认证配置若指定了 SSL 相关参数则构造 SSL 配置读取 Kafka 消息根据输出格式选择FormatTransform.readFromKafka字符串或FormatTransform.readAvrosFromKafkaAVRO写出到 Pub/Sub先用Values.create()丢弃 Kafka 记录的 Key、只保留 Value再写入 Pub/Sub 输出主题。源码中的核心调用链如下KafkaToPubsub.java/* * Steps: * 1) Read messages in from Kafka * 2) Extract values only * 3) Write successful records to PubSub */ if (options.getOutputFormat() FormatTransform.FORMAT.AVRO) { pipeline .apply(readAvrosFromKafka, FormatTransform.readAvrosFromKafka( options.getBootstrapServers(), topicsList, kafkaConfig, sslConfig)) .apply(createValues, Values.create()) .apply(writeAvrosToPubSub, PubsubIO.writeAvros(AvroDataClass.class)); } else { pipeline .apply(readFromKafka, FormatTransform.readFromKafka( options.getBootstrapServers(), topicsList, kafkaConfig, sslConfig)) .apply(createValues, Values.create()) .apply(writeToPubSub, new FormatTransform.FormatOutput(options)); }从源码结构可以推断Kafka 侧的读取统一封装在 FormatTransform.java它基于 Beam 官方的KafkaIO构建消费者通过withKeyDeserializerAndCoder/withValueDeserializerAndCoder指定反序列化器与编码器通过withConsumerConfigUpdates注入认证/SSL 配置通过withConsumerFactoryFn(new SslConsumerFactoryFn(sslConfig))注入自定义消费者工厂最后withoutMetadata()去除 Kafka 元数据保持输出干净。运行前置要求要跑通该示例环境上需要满足以下条件与主类 Javadoc 中列出的 Pipeline Requirements 一致Java 8Kafka Bootstrap Server(s) 处于运行状态已存在 Kafka 源主题source topic已存在 Pub/Sub 目标输出主题destination output topic可选已存在可用的 HashiCorp Vault 实例可选为 Kafka 配置了安全的 SSL 连接。Gradle 准备工作示例通过 Gradle 的JavaExec任务驱动管道运行。需要在build.gradle中添加如下任务task execute (type:JavaExec) { mainClass System.getProperty(mainClass) classpath sourceSets.main.runtimeClasspath systemProperties System.getProperties() args System.getProperty(exec.args, ).split() }该任务通过 JVM 系统属性传递主类名与执行参数随后即可用下面的命令启动管道gradle clean execute -DmainClassorg.apache.beam.examples.complete.kafkatopubsub.KafkaToPubsub \ -Dexec.args--argumentvalue --argumentvalue这里的-DmainClass指定管道入口类-Dexec.args中传入的全部参数会以空格分隔后作为管道选项。运行管道核心参数详解执行管道时需要指定四类核心参数Kafka Bootstrap servers、Kafka 输入主题、Pub/Sub 输出主题、输出格式命令格式如下--bootstrapServershost:port \ --inputTopicsyour-input-topic \ --outputTopicprojects/your-project-id/topics/your-topic-pame \ --outputFormatAVRO|PUBSUB这些参数在 KafkaToPubsubOptions.java 中均有定义其中前四个带有Validation.Required注解即缺失会导致参数校验失败。各参数含义与约束如下表参数必填含义与取值说明--bootstrapServers是逗号分隔的 Kafka Bootstrap Server 列表例如server1:[port],server2:[port]--inputTopics是逗号分隔的待读取 Kafka 主题列表例如topic1,topic2源码中按,拆分后校验不能为空--outputTopic是Pub/Sub 输出主题格式必须是projects/project-id/topics/topic-name--outputFormat是输出格式取值为枚举FORMAT.AVRO或FORMAT.PUBSUB定义见 FormatTransform.java--secretStoreUrl否Vault 中 Kafka 凭据的 URL 地址--vaultToken否访问 Vault 所需的 Token--truststorePath否truststore 文件路径本地路径或gs://开头的 GCS 路径--keystorePath否keystore 文件路径本地路径或gs://开头的 GCS 路径--truststorePassword否truststore 密码--keystorePassword否keystore 密码--keyPassword否keystore 中私钥的密码--kafkaConsumerConfig否附加的 Kafka Consumer 配置格式为key1value1;key2value2会合并进消费者配置关于最后一项--kafkaConsumerConfig虽然示例 README 未展开说明但它同样定义在选项接口中KafkaToPubsubOptions.java并且在主类run方法开头就通过parseKafkaConsumerConfig解析并入 Kafka 配置KafkaToPubsub.java。其底层解析逻辑位于 Utils.java按;拆分键值对再按拆分出 key 与 value可用于传递auto.offset.reset、max.poll.records等任意原生 Kafka Consumer 属性。切换 Runner默认情况下管道使用DirectRunner在本地执行。要切换到其他执行引擎只需追加--runnerYOUR_SELECTED_RUNNER例如在 Google Cloud Dataflow 上运行时可指定--runnerDataflowRunner并配合相应的 GCP 参数项目、区域、临时目录等。关于如何准备和运行各种 Runner可参考 examples/java 的 README 中的说明。SASL/SCRAM 认证与 HashiCorp Vault 凭据管理当 Kafka 集群启用了 SASL/SCRAM 认证时需要提供 Kafka 用户名与密码。示例的推荐做法是把凭据托管在 HashiCorp Vault 中运行时通过--secretStoreUrl与--vaultToken两个参数动态拉取--secretStoreUrlhttp(s)://host:port/path/to/credentials --vaultTokenyour-token底层实现位于 Utils.java 的getKafkaCredentialsFromVault方法它通过 Apache HttpClient 发起 GET 请求并在请求头携带X-Vault-Token令牌随后解析 Vault 返回的 JSON。Vault 响应具有固定嵌套结构真实数据位于data.data节点之下{ request_id: 6a0bb14b-ef24-256c-3edf-cfd52ad1d60d, lease_id: , renewable: false, lease_duration: 0, data: { data: { bucket: kafka_to_pubsub_test, key_password: secret, keystore_password: secret, keystore_path: ssl_cert/kafka.keystore.jks, password: admin-secret, truststore_password: secret, truststore_path: ssl_cert/kafka.truststore.jks, username: admin }, metadata: { created_time: 2020-10-20T11:43:11.109186969Z, deletion_time: , destroyed: false, version: 8 } }, wrap_info: null, warnings: null, auth: null }凭据字段名在 KafkaPubsubConstants.java 中定义为常量username、password、kafka凭据分组名。拿到用户名/密码后configureKafka方法Utils.java会构造 Kafka 原生配置sasl.mechanism设置为SCRAM-SHA-512sasl.jaas.config设置为ScramLoginModule格式的 JAAS 字符串内嵌用户名与密码。如果未提供--secretStoreUrl与--vaultToken主类会打印警告日志并尝试发起未授权连接KafkaToPubsub.java若 Vault 中缺少username/password字段同样会回退到未授权连接。SSL 加密连接配置如果要启用 Kafka 与 Beam 管道之间的安全 SSL 连接需要提供以下参数truststore 文件路径支持本地路径或gs://开头的 GCS 路径keystore 文件路径支持本地路径或gs://开头的 GCS 路径truststore 密码keystore 密码key 密码keystore 中私钥的密码。--truststorePathpath/to/kafka.truststore.jks --keystorePathpath/to/kafka.keystore.jks --truststorePasswordyour-truststore-password --keystorePasswordyour-keystore-password --keyPasswordyour-key-password这些值首先由configureSslUtils.java映射为 Kafka 的SslConfigs属性ssl.truststore.location、ssl.keystore.location、ssl.truststore.password、ssl.keystore.password、ssl.key.password。判断是否启用 SSL 的条件是isSslSpecifiedUtils.java只要 truststore/keystore 路径、truststore 密码或 key 密码任一非空即认为需要 SSL 配置。真正创建带 SSL 的 Kafka 消费者的是 SslConsumerFactoryFn.java它实现了 Beam 的SerializableFunctionMapString, Object, Consumerbyte[], byte[]作为自定义消费者工厂注入KafkaIO。其关键逻辑若 truststore/keystore 路径缺失则回退为普通KafkaConsumer并告警若路径以gs://开头通过getGcsFileAsLocal把 GCS 文件下载到本地临时路径/tmp/kafka.truststore.jks、/tmp/kafka.keystore.jks否则检查本地文件是否存在最终在消费者配置中写入security.protocolSASL_SSL以及全套 SSL 位置/密码属性SslConsumerFactoryFn.java。特别提醒Kafka 到 Pub/Sub 任务在使用分布式 Runner 执行时SSL 证书只能放在 GCSgs://路径。本地文件仅适用于本地 Runner——源码checkFileExists的日志也明确提示 Local files dont support when in using distribute runnerSslConsumerFactoryFn.java。支持的输出格式PUBSUB 与 AVRO管道支持两种输出格式由--outputFormat参数决定PUBSUB或AVRO。PubSubMessage 格式开箱即用选择PUBSUB时无需任何额外改动。写入逻辑封装在FormatTransform.FormatOutput中FormatTransform.java先用MapElements把每条 JSON 字符串按 UTF-8 编码包装成带空属性集的PubsubMessage再通过PubsubIO.writeMessages().to(outputTopic)写入目标主题。input .apply(convertMessagesToPubsubMessages, MapElements.into(TypeDescriptor.of(PubsubMessage.class)) .via((String json) - new PubsubMessage(json.getBytes(Charsets.UTF_8), ImmutableMap.of()))) .apply(writePubsubMessagesToPubSub, PubsubIO.writeMessages().to(options.getOutputTopic()));AVRO 格式需定制 SchemaAVRO 输出需要做四步定制示例已提供全套参考实现目录结构如下avro/AvroDataClass.java定义 AVRO Schema 的 POJO通过DefaultCoder(AvroCoder.class)注解配合 Beam 的AvroCoder进行序列化avro/AvroDataClassKafkaAvroDeserializer.java基于 ConfluentAbstractKafkaAvroDeserializer的自定义 Kafka 反序列化器transforms/FormatTransform.javaKafka 读取与 Pub/Sub 写入的组装点。具体步骤如下创建描述 AVRO Schema 的类。参照AvroDataClass仅需定义所需字段及其 getter/setterDefaultCoder(AvroCoder.class) public class AvroDataClass { String field1; Float field2; Float field3; // 构造器、getter、setter 省略 }创建自己的 Kafka Avro 反序列化器。参照AvroDataClassKafkaAvroDeserializer重命名类并把类型参数替换为你的 Schema 类即可。修改FormatTransform.readAvrosFromKafka将你的 Schema 类与反序列化器放入对应参数return KafkaIO.String, AvroDataClassread() ... .withValueDeserializerAndCoder( AvroDataClassKafkaAvroDeserializer.class, AvroCoder.of(AvroDataClass.class)) // put your classes here ...可选添加 Beam Transform如果你的业务需要在读写之间做加工例如过滤、清洗、聚合可在读 Kafka 与写 Pub/Sub 之间插入自定义 Beam Transform。修改 KafkaToPubsub.java 中的写入步骤把 Schema 类放入writeAvrosToPubSubif(options.getOutputFormat()FormatTransform.FORMAT.AVRO){ ... .apply(writeAvrosToPubSub,PubsubIO.writeAvros(AvroDataClass.class)); // put your SCHEMA class here }注意如果数据在 Transform 过程中 Schema 发生了改变应使用变化后的类定义。作为 Google Dataflow 模板运行该示例同样以 Google Dataflow 模板的形式存在可基于 Google Cloud Platform 构建并运行对应 DataflowTemplates 仓库中的kafka-to-pubsub模板。如果你需要直接以模板方式部署请参考该模板自身附带的说明文档README.md。端到端测试示例 README 中明确标注端到端测试E2E tests状态为TBD待定。当前仓库内该示例目录下未提供对应的集成测试代码读者在自行验证时可通过本地 DirectRunner 结合本地 Kafka 与本地模拟 Pub/Sub 或测试环境完成链路验证。使用注意事项小结单个用户账号示例仅支持使用单个 Kafka 用户账号对全部源服务器与主题进行认证明文与 SSL 二选一依据提供 Vault 凭据则启用 SASL/SCRAM提供 SSL 参数则启用 SASL_SSL 协议两者都未提供时回退为未授权明文连接分布式 Runner 的证书位置SSL 证书必须放在 GCSgs://本地文件仅限本地运行参数校验bootstrapServers、inputTopics、outputTopic、outputFormat四项为必填缺失会直接校验失败多主题支持inputTopics支持逗号分隔的多个主题bootstrapServers也支持多个服务器输出侧则固定为单个 Pub/Sub 主题。赞分享【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载相关推荐Apache Beam KafkaToPubsub 示例实战从 Apache Kafka 到 Google Cloud Pub/Sub 的流式数据接入Apache Beam KafkaToPubsub 示例实战从 Apache Kafka 到 Google Cloud Pub/Sub 的流式数据接入 本文围大数据批处理流处理数据工程Apache Beam PubsubIO 实战从 Google Cloud Pub/Sub 流式读取数据的完整代码解读与源码剖析Apache Beam PubsubIO 实战从 Google Cloud Pub/Sub 流式读取数据的完整代码解读与源码剖析 本指南以仓库 learnin大数据批处理流处理数据工程Scriptographer完全指南JavaScript自动化Adobe Illustrator的终极解决方案Scriptographer完全指南JavaScript自动化Adobe Illustrator的终极解决方案 Adobe Illustrator作为专业设计上一篇Velite与Next.js完美集成构建现代内容网站下一篇Automa多浏览器支持Chrome和Firefox的兼容性指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
02
RELATED NEWS

相关资讯

更多网站建设与数字化升级内容

03
WHY YAOTU

想打造同款高转化官网?

懂行业、懂生意,从建站到增长一站式陪跑

◈

场景化定制

不做模板站,围绕你的业务场景量身设计,小众不撞款。

◐

营销型架构

以转化目标组织内容与路径,让官网真正带来询盘。

▲

全周期服务

设计、开发、运营、运维一体,上线只是开始。

免费获取你的建站方案

留下需求,专属顾问 24 小时内为你输出方案建议。