Apache Pulsar 与 Spark Streaming 集成实战:基于 SparkStreamingPulsarReceiver 构建实时流处理应用
消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载导读本文围绕 Apache Pulsar 官方文档中关于 Spark Streaming 适配器的核心内容系统讲解如何通过pulsar-spark库中提供的SparkStreamingPulsarReceiver自定义 Receiver让 Spark Streaming 直接消费 Pulsar 中的原始消息并以 RDDResilient Distributed Dataset形式进行批式处理。读完本文你将掌握在 Maven/Gradle 中正确引入pulsar-spark依赖、基于JavaStreamingContext.receiverStream接入 Pulsar 的完整代码范式以及如何切换AuthenticationDisabled与AuthenticationToken等不同认证方式。Spark Streaming 与 Pulsar 的集成方式Apache Pulsar 提供了一组语言无关的协议适配器adaptor其中 Spark Streaming 适配器以自定义 Receiver的形式存在。它的定位非常明确让 Spark Streaming 能够“接收来自 Pulsar 的原始数据raw data”而不是把 Pulsar 当作普通 Socket 或 Kafka 来处理。从官方文档site2/website-next/versioned_docs/version-2.2.0/adaptors-spark.md的表述看其工作机制可以概括为Spark Streaming 应用通过该 Receiver 从 Pulsar 订阅 topic持续拉取消息拉取到的数据被组织为 RDD从而可以借助 Spark 生态的各类算子map、filter、reduce、窗口计算等进行灵活处理。这种模式的本质是把 Pulsar 当作 Spark Streaming 的输入源input source两者之间通过 Pulsar 客户端协议而非 Kafka 协议通信。需要说明pulsar-spark适配器代码维护在独立的apache/pulsar-adapters仓库中当前 Pulsar 仓库内保存的是其使用文档因此本文以文档为主线结合本仓库中的 Pulsar 客户端 API 源码如ConsumerConfigurationData、Authentication系列类补充实现层面的细节。前置条件在构建配置中引入 pulsar-spark 依赖使用该 Receiver 前需要先在 Java 工程中声明pulsar-spark库的依赖。文档分别给出了 Maven 与 Gradle 两种配置方式。Maven 配置在pom.xml的properties与dependencies两个块中分别添加版本属性和依赖项!-- in your properties block -- pulsar.versionpulsar:version/pulsar.version !-- in your dependencies block -- dependency groupIdorg.apache.pulsar/groupId artifactIdpulsar-spark/artifactId version${pulsar.version}/version /dependency其中pulsar:version是文档站点构建时自动替换的占位符实际使用时应替换为具体版本号例如2.2.0。pulsar-spark的 groupId 为org.apache.pulsar与 Pulsar 其他 Java 组件保持一致。Gradle 配置在build.gradle中对应添加def pulsarVersion pulsar:version dependencies { compile group: org.apache.pulsar, name: pulsar-spark, version: pulsarVersion }提示Gradle 示例中的compile是 Gradle 3.x 及更早版本中的经典写法若使用 Gradle 4.x通常建议改用implementation配置。同时注意版本号应与本机 Pulsar 服务端版本匹配避免协议不兼容。核心用法将 Receiver 接入 JavaStreamingContext文档给出的核心使用模式非常简洁构造一个SparkStreamingPulsarReceiver实例然后把它传给JavaStreamingContext.receiverStream(...)方法得到一个JavaReceiverInputDStreambyte[]。以 version-2.2.0 文档中的原始示例为基础SparkConf conf new SparkConf().setMaster(local[*]).setAppName(pulsar-spark); JavaStreamingContext jssc new JavaStreamingContext(conf, Durations.seconds(5)); ClientConfiguration clientConf new ClientConfiguration(); ConsumerConfiguration consConf new ConsumerConfiguration(); String url pulsar://localhost:6650/; String topic persistent://public/default/topic1; String subs sub1; JavaReceiverInputDStreambyte[] msgs jssc .receiverStream(new SparkStreamingPulsarReceiver(clientConf, consConf, url, topic, subs));关键点拆解SparkConf配置了运行模式local[*]本地多线程与应用名称JavaStreamingContext的批处理间隔设为 5 秒Durations.seconds(5)clientConf/consConf分别对应 Pulsar 客户端与消费者的配置对象早期 API 形态url指向 Pulsar broker 服务地址pulsar://localhost:6650/为单机默认topic是完整 topic 名称subs是订阅名称最终msgs是一个JavaReceiverInputDStreambyte[]每条消息以byte[]形式进入 Spark 处理管道可继续调用 Spark Streaming 算子处理。演进后的推荐用法ConsumerConfigurationData 与认证注入随着 Pulsar Java 客户端 API 演进官方文档参见 site2/docs/adaptors-spark.md 与 site2/website-next/docs/adaptors-spark.md中的推荐写法改为使用ConsumerConfigurationDatabyte[]统一描述订阅配置并把认证对象作为构造参数的第三个入参。这一形态与本仓库中 pulsar-client-api 下的客户端抽象完全对应。基础用法禁用认证String serviceUrl pulsar://localhost:6650/; String topic persistent://public/default/test_src; String subs test_sub; SparkConf sparkConf new SparkConf().setMaster(local[*]).setAppName(Pulsar Spark Example); JavaStreamingContext jsc new JavaStreamingContext(sparkConf, Durations.seconds(60)); ConsumerConfigurationDatabyte[] pulsarConf new ConsumerConfigurationData(); SetString set new HashSet(); set.add(topic); pulsarConf.setTopicNames(set); pulsarConf.setSubscriptionName(subs); SparkStreamingPulsarReceiver pulsarReceiver new SparkStreamingPulsarReceiver( serviceUrl, pulsarConf, new AuthenticationDisabled()); JavaReceiverInputDStreambyte[] lineDStream jsc.receiverStream(pulsarReceiver);与旧版 API 相比这里的变化体现在订阅配置集中管理通过ConsumerConfigurationDatabyte[]的setTopicNames(SetString)设置 topic 集合注意是Set天然支持一次订阅多个 topic通过setSubscriptionName(String)设置订阅名认证显式化构造函数第三个参数传入认证实现。AuthenticationDisabled表示不启用认证其实现位于 pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDisabled.java是未配置认证时的默认行为返回值类型不变receiverStream返回的仍是JavaReceiverInputDStreambyte[]。使用 JWT Token 认证如果 Pulsar 集群开启了认证可以替换认证参数。文档给出的 Token 认证示例SparkStreamingPulsarReceiver pulsarReceiver new SparkStreamingPulsarReceiver( serviceUrl, pulsarConf, new AuthenticationToken(token:secret-JWT-token));AuthenticationToken使用 JWT Token 完成认证token:secret-JWT-token即认证凭证字符串。从源码结构看Pulsar 的认证体系以 pulsar-client-api/src/main/java/org/apache/pulsar/client/api/Authentication.java 为统一抽象接口AuthenticationDisabled、AuthenticationToken等均为该接口的具体实现因此理论上其他认证实现如 TLS 客户端证书等同样可以按此模式传入。完整示例统计包含 Pulsar 的消息数官方文档给出的配套示例位于独立的pulsar-adapters仓库examples/spark模块其业务逻辑是在接收到的消息流中统计包含字符串Pulsar的消息数量。基于上文的基础用法一个完整的 Spark Streaming 应用骨架如下JavaReceiverInputDStreambyte[] lineDStream jsc.receiverStream(pulsarReceiver); // 将字节流转换为字符串过滤出包含 Pulsar 的消息并计数 JavaDStreamString lines lineDStream.map(bytes - new String(bytes, StandardCharsets.UTF_8)); JavaPairDStreamString, Long counts lines .filter(line - line.contains(Pulsar)) .mapToPair(line - new Tuple2(line, 1L)) .reduceByKey(Long::sum); counts.print(); jsc.start(); jsc.awaitTermination();该示例展示了从 Pulsar 消费 → RDD 转换 → 业务过滤 → 聚合统计的完整链路可作为自定义流处理逻辑的起点。使用要点与注意事项结合文档与 Pulsar 客户端实现以下几点值得在实际开发中关注Topic 命名规范示例中的persistent://public/default/topic1是带完整 domain 的 topic 全名tenant/namespace/topic 三段式Pulsar 中 topic 须以persistent://或non-persistent://前缀标识存储类型订阅模式Receiver 内部是标准的 Pulsar Consumer 行为subs订阅名决定了消费位置与消息确认ack的归属不同批处理间隔下Receiver 会持续拉取数据并交给 Spark 按批封装为 RDD认证一致性认证对象由 Pulsar 客户端 API 抽象Authentication接口统一管理若集群启用了认证务必在构造 Receiver 时传入匹配的认证实现否则消费会因认证失败而中断字节流语义JavaReceiverInputDStreambyte[]说明消息以原始字节到达schema 解析与反序列化需要由 Spark 侧完成这与 Pulsar Java Client 的泛型消费模型一致。总结通过pulsar-spark适配器Apache Pulsar 可以无缝接入 Spark Streaming 生态以自定义 Receiver 消费原始消息以 RDD 形式交给 Spark 做分布式处理。本文覆盖了依赖引入、旧/新两代 API 用法、认证切换以及端到端计数示例足以支撑读者搭建第一个 Pulsar Spark Streaming 的实时数据处理管道。更多历史版本的文档变体可参考 site2/website/versioned_docs/version-2.2.0/adaptors-spark.md 以及仓库 site2/website-next/versioned_docs 目录下各版本的对应文档。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐Apache Airflow与流处理平台集成Flink、Spark Streaming实战Apache Airflow与流处理平台集成Flink、Spark Streaming实战 引言为什么需要流处理与工作流调度集成 在现代数据架构中实时数后端任务调度工作流自动化数据编排批处理数据工程流程编排Angular 2与Apache Spark Streaming集成构建实时数据分析应用Angular 2与Apache Spark Streaming集成构建实时数据分析应用 你是否还在为实时数据处理与前端展示的割裂而困扰本文将带你通过 gh文档Apache Pulsar 与 Spark 集成指南使用 Spark Streaming Receiver 消费 Pulsar 消息Apache Pulsar 与 Spark 集成指南使用 Spark Streaming Receiver 消费 Pulsar 消息 本指南讲解 Apache消息队列后端流处理上一篇如何构建完整的响应式设计系统inuitcss与Sass-MQ集成终极指南下一篇CGrep 高级搜索技巧正则表达式、语义匹配与测试代码过滤全解析创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考