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),仅供参考

相关新闻

IGBT选型实战指南:从工况分析到参数计算与验证

IGBT选型实战指南:从工况分析到参数计算与验证

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/24 13:59:32 阅读更多 →
PaddleSpeech WaveFlow 声码器波形合成指南:synthesize.py 全流程解析与实战

PaddleSpeech WaveFlow 声码器波形合成指南:synthesize.py 全流程解析与实战

人工智能语音音频NLP媒体生成 【免费下载链接】PaddleSpeech Easy-to-use Speech Toolkit including Self-Supervised Learning model, SOTA/Streaming ASR with punctuation, Streaming TTS with text frontend, Speaker Verification System, End-to-End Speech Translation …

2026/9/24 13:59:32 阅读更多 →
ESP32-S3-N16R8实战指南:存储配置、引脚复用与选型避坑

ESP32-S3-N16R8实战指南:存储配置、引脚复用与选型避坑

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/24 13:59:32 阅读更多 →

最新新闻

Claude Code嵌入式AI编程实战:SPI/I2C/ADC驱动开发指南

Claude Code嵌入式AI编程实战:SPI/I2C/ADC驱动开发指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/24 14:42:58 阅读更多 →
python代码保存到Gitee完整步骤

python代码保存到Gitee完整步骤

一、环境准备 1.本地安装Git Git - Install for Windows 2.创建Gitee账号 3.创建本地Python-study文件夹 二、完整步骤操作 1.在python-study文件夹里右键选择Open Git Bash here开Git终端 2.输入Gitee中的代码 3.弹出登录界面 登录成功后 点击刷新 弹出以下界面 4.点击“第…

2026/9/24 14:42:58 阅读更多 →
RGBD相机视觉检测:MV-EB435i深度对齐与点云融合实践

RGBD相机视觉检测:MV-EB435i深度对齐与点云融合实践

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/24 14:42:58 阅读更多 →
现在口碑好的AI论文写作工具有哪些品牌?聊聊真实使用体验

现在口碑好的AI论文写作工具有哪些品牌?聊聊真实使用体验

每到期末、毕业答辩、课题申报阶段,很多学生都会陷入论文写作的困境:选题毫无头绪、大纲搭建逻辑混乱、正文撰写耗时长、参考文献格式出错、查重重复率偏高、AIGC检测告警、本校论文排版标准复杂。依靠纯人工从零开始撰写、一遍遍修改格式和降重&#xf…

2026/9/24 14:42:58 阅读更多 →
DevilutionX GDB 调试增强:pretty-printer 的加载方式、配置实战与实现原理

DevilutionX GDB 调试增强:pretty-printer 的加载方式、配置实战与实现原理

游戏开发 【免费下载链接】DevilutionX Diablo build for modern operating systems 项目地址: https://gitcode.com/gh_mirrors/de/DevilutionX 点击查看 免费下载 导读 DevilutionX(暗黑破坏神 1 的现代操作系统移植版)在仓库中内置了一套…

2026/9/24 14:42:58 阅读更多 →
StoryDiffusion 快速上手:角色前后一致的多格漫画生成

StoryDiffusion 快速上手:角色前后一致的多格漫画生成

StoryDiffusion 快速上手:角色前后一致的多格漫画生成 【免费下载链接】StoryDiffusion Accepted as [NeurIPS 2024] Spotlight Presentation Paper 项目地址: https://gitcode.com/GitHub_Trending/st/StoryDiffusion StoryDiffusion 是一个基于 Consistent…

2026/9/24 14:41:58 阅读更多 →

日新闻

基于YOLOv8的渔船作业监控系统:从环境搭建到边缘部署全流程

基于YOLOv8的渔船作业监控系统:从环境搭建到边缘部署全流程

简介:这是一套面向计算机、人工智能、自动化等专业学生与教师的毕业设计级项目资源,围绕YOLOv8实现渔船作业监控系统,可用于毕设、课程设计、大作业或项目立项演示。压缩包共97个文件,约24.21MB,以70个Python源码文件为…

2026/9/24 0:00:19 阅读更多 →
单细胞注释实战:基于Scanpy的标记基因与参考映射流程解析

单细胞注释实战:基于Scanpy的标记基因与参考映射流程解析

简介:一份基于单细胞RNA测序数据的细胞类型注释算法研究Python毕业设计源码,针对计算机相关专业正在做毕设或需要项目实战的学习者,可用于课程设计与期末大作业。项目代码完整、经导师指导评审通过,可直接运行,覆盖数据…

2026/9/24 0:00:19 阅读更多 →
C#源生成器实战:用增量生成器替代反射,告别AOT崩溃

C#源生成器实战:用增量生成器替代反射,告别AOT崩溃

第一次在项目里被反射卡住,是在一个老旧的WinForms模块里:几十个类依赖PropertyChanged通知,运行时反射读属性、发通知,每次启动慢半拍不说,一上.NET Native/AOT裁剪模式几乎全面崩盘。后来我把这段逻辑全部改成C#源生…

2026/9/24 0:00:19 阅读更多 →

周新闻

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

直接铺开项目本身吧。这几个月我一直在折腾一件事:用Flutter给OpenHarmony做一款游戏集合类的App,说白了就是把若干小游戏塞进一个壳里,用统一入口分发。这个方向本身不算新鲜,真正让我花了不少心思的,是首页那堆游戏卡…

2026/9/24 14:34:13 阅读更多 →
Word表格编号全攻略:从列表编号到题注交叉引用

Word表格编号全攻略:从列表编号到题注交叉引用

写Word文档,最让人头疼的往往是那些“看起来不起眼”的小问题。比如表格编号这事:今天在表后面多加了两个空白行,明天给客户交稿前发现整个章节的编号全部错位,光是挨个改序号就能耗掉大半个下午。我前阵子帮人整理一份上百页的技…

2026/9/24 9:10:42 阅读更多 →
从第一个站到第二个站:独立开发者的静态网站选型与落地实践

从第一个站到第二个站:独立开发者的静态网站选型与落地实践

1. 项目概述1.1 核心需求解析做独立开发者这几年,说实话,第一个网站上线的那天晚上我兴奋得没睡着。但等它跑了半年,流量惨淡、功能臃肿、代码自己都懒得看第二遍之后,我才慢慢琢磨明白一个道理:第一个网站是练手&…

2026/9/24 14:33:56 阅读更多 →

月新闻

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能分类:[AI/大模型]细分主题:AI 增强型 CI/CD 流水线自动化与 GitOps 实践:Agent 工作流、工具调用与任务拆解:从原型到生产的验收清单很多团队在尝试用大…

2026/9/24 12:50:34 阅读更多 →
容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场分类:[工程技术]细分主题:Kubernetes 生产环境运维与排障实战:可复制的项目复盘模板与决策记录大部分团队的事故复盘报告,最后都变成了躺在 Confluence 或钉…

2026/9/24 14:33:48 阅读更多 →
容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步分类:[工程技术]细分主题:Docker 容器化技术与镜像安全管理:核心链路的逐步实现与关键代码取舍面对一个积累了五六年历史包袱的单体架构应用(包含 Web 接口、后台…

2026/9/24 12:49:17 阅读更多 →