消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载Apache Pulsar 官方提供了面向 Apache Spark Streaming 的自定义 Receiver接收器——SparkStreamingPulsarReceiver它位于pulsar-spark适配库中允许 Spark Streaming 应用直接从 Pulsar 消费原始数据并以弹性分布式数据集RDD的形式交给下游处理。本文基于仓库中 version-2.3.2 的适配器文档 展开完整覆盖依赖引入、Receiver 初始化、订阅配置、认证方式与统计示例并结合本仓库的 Pulsar 客户端源码说明底层配置项的默认值与作用帮助你快速搭建Pulsar 消息 → Spark 实时处理的数据管道。Spark Streaming 与 Pulsar 的集成方式Spark Streaming 是 Apache Spark 提供的实时流处理框架。Pulsar 为其提供的是一个自定义 Receiver应用程序通过该 Receiver 从 Pulsar 中接收原始raw消息数据数据以RDDResilient Distributed Dataset的形式进入 Spark Streaming 的 DStream 数据流之后便可以使用 Spark 提供的各类转换算子如map、filter、reduceByKey等对其进行灵活处理。这种模式的价值在于Pulsar 承担消息的生产、存储与分发Spark Streaming 专注流式计算与窗口聚合两者通过一个标准化的 Receiver 桥接无需在应用侧自行实现 Pulsar Consumer 与 Spark 数据源的对接逻辑。前置条件引入 pulsar-spark 依赖使用该 Receiver 前需要在 Java 工程中引入pulsar-spark库的依赖它由org.apache.pulsar组织发布。本仓库的适配器文档提供了 Maven 与 Gradle 两种构建工具的配置方式。Maven在pom.xml中添加如下配置其中pulsar.version应替换为你实际使用的 Pulsar 版本文档中原样保留了构建期的版本占位符pulsar:version发布时会被替换为具体版本号!-- 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 /dependencyGradle在build.gradle中添加def pulsarVersion pulsar:version dependencies { compile group: org.apache.pulsar, name: pulsar-spark, version: pulsarVersion }注意本文所依据的 version-2.3.2 文档示例中compile是较老的 Gradle 配置写法新版 Gradle 中可等价替换为implementation或api。pulsar-spark适配器与pulsar-client属于同一org.apache.pulsar坐标体系使用相同版本号能避免客户端与服务端协议不兼容的问题。使用方式创建 Receiver 并接入 JavaStreamingContext核心用法非常简洁构造一个SparkStreamingPulsarReceiver实例将其作为参数传给JavaStreamingContext的receiverStream方法即可得到JavaReceiverInputDStreambyte[]随后就能对接收到的消息字节流做各种转换。文档给出的完整示例代码如下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);关键参数逐个拆解serviceUrlPulsar Broker 的服务地址。示例使用pulsar://localhost:6650/即默认运行在 6650 端口的二进制协议入口生产环境通常替换为集群地址可包含多个 Broker如pulsar://host1:6650,host2:6650。topic要订阅的主题。示例使用完整主题名persistent://public/default/test_src其中persistent表示持久化主题public/default是租户public下的默认命名空间。subs订阅名称subscription name示例为test_sub。订阅名用于标记消费组进度同一个订阅名下的多个消费者会按订阅类型共享消息消费进度。SparkConf 与 JavaStreamingContextsetMaster(local[*])表示本地多线程运行模式Durations.seconds(60)设定批处理间隔为 60 秒即每 60 秒生成一个 RDD 批次交给下游计算。ConsumerConfigurationDatabyte[]Pulsar 客户端的消费配置载体这里以byte[]作为消息类型表示接收原始字节数据。通过setTopicNames指定主题集合、setSubscriptionName指定订阅名。AuthenticationDisabled认证策略对象表示不启用认证用于本地开发或未开启认证的集群。ConsumerConfigurationData 的底层实现与可选配置ConsumerConfigurationDataT是pulsar-client中真正被 Consumer 使用的配置类位于 pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ConsumerConfigurationData.java。从源码看它实现了Serializable, Cloneable因此可以随 Spark 任务序列化分发到 Executor。该类的字段及其默认值直接决定了 Receiver 的消费行为常用项包括配置字段默认值说明topicNames空 TreeSet订阅的主题集合文档示例即通过setTopicNames设置subscriptionNamenull订阅名称须显式设置subscriptionTypeExclusive独占消费模式可改为Shared、Failover、Key_Shared以支持并行消费subscriptionModeDurable订阅是否持久化持久订阅重启后可从上次消费位点继续receiverQueueSize1000消费者本地接收队列大小影响背压与吞吐ackTimeoutMillis0消息确认超时时间毫秒0 表示不启用consumerNamenull消费者名称便于监控与追踪priorityLevel0消费者优先级用于 Shared 订阅下的消息分发以receiverQueueSize默认 1000为例该值决定了 Pulsar 客户端在本地预取的消息数量上限在 Spark Streaming 场景下它直接影响每个批次能够缓冲的数据量进而影响吞吐与内存占用。上述默认值均为当前仓库pulsar-client源码中的真实取值供你在调优时参考——在使用SparkStreamingPulsarReceiver前可通过对应的 setter 方法按需修改这些配置。认证配置从关闭认证到 JWT TokenSparkStreamingPulsarReceiver的第三个构造参数是 Pulsar 的认证策略对象。文档示例使用的是AuthenticationDisabled适用于未开启认证的测试环境。该类的实现位于 pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationDisabled.java。当集群开启认证后需要替换为其他认证实现。例如使用 JWT Token 认证时可构造AuthenticationToken并传入密钥SparkStreamingPulsarReceiver pulsarReceiver new SparkStreamingPulsarReceiver( serviceUrl, pulsarConf, new AuthenticationToken(token:secret-JWT-token));AuthenticationToken是pulsar-client提供的 Token 认证实现位于 pulsar-client/src/main/java/org/apache/pulsar/client/impl/auth/AuthenticationToken.java其行为在 AuthenticationTokenTest 中有对应的单元测试覆盖。除 Token 外Pulsar 还支持 TLS、Athenz 等认证插件均可按同样的方式作为构造参数传入。完整示例统计包含 Pulsar 的消息条数仓库文档中提到一个完整的可运行示例位于 Spark 适配器仓库的 examples 目录SparkStreamingPulsarReceiverExample.java其业务逻辑是统计接收到的消息中包含字符串 Pulsar 的消息条数。虽然该示例源码托管在独立的适配器仓库中但核心逻辑可以基于本文的接入代码延伸理解通过receiverStream(pulsarReceiver)得到JavaReceiverInputDStreambyte[]即原始消息流对消息字节流应用filter或map转换判断每条消息内容是否包含Pulsar使用count算子对每个批次进行计数并输出即可得到包含 Pulsar 关键词的消息数量这一实时统计指标。如果你希望在本仓库中进一步核对 Receiver 消费侧的行为可对照查看 ConsumerConfigurationData.java 中各配置字段的默认值以及 conf/broker.conf 中 Broker 侧的认证与服务端口配置从而端到端理解从Spark 任务创建 Receiver到Pulsar Broker 推送消息的完整链路。版本与适用范围说明本文依据的文档版本为 Pulsar 2.3.2见 version-2.3.2 的 adaptors-spark 文档同一份适配器文档在仓库的多个版本目录中均有对应副本如 site2/website-next/docs/adaptors-spark.md 与各 versioned_docs 目录新版本文档在用法一致的基础上补充了 Token 认证的示例。pulsar-spark适配器本身维护在 Apache Pulsar 的独立适配器仓库中本仓库主要承载文档与 Pulsar 客户端核心实现。使用前请确认三点前提一是 Pulsar 集群已启动且 Broker 的 6650 端口可达二是工程中pulsar-spark的版本与 Pulsar 服务端版本保持兼容建议取相同版本号三是若生产集群启用了认证务必配置与 Broker 端一致的认证插件参数避免因认证失败导致 Receiver 无法建立连接。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐Apache Pulsar 与 Apache Spark Streaming 集成指南使用 pulsar-spark 接收器消费消息Apache Pulsar 与 Apache Spark Streaming 集成指南使用 pulsar spark 接收器消费消息 Apache Pulsa消息队列后端流处理Apache Pulsar 与 Spark 集成指南使用 Spark Streaming Receiver 消费 Pulsar 消息Apache Pulsar 与 Spark 集成指南使用 Spark Streaming Receiver 消费 Pulsar 消息 本指南讲解 Apache消息队列后端流处理Apache Pulsar 与 Spark Streaming 集成指南pulsar-spark 接收器接入与实战配置Apache Pulsar 与 Spark Streaming 集成指南pulsar spark 接收器接入与实战配置 本篇文章基于 Apache Pulsa消息队列后端流处理上一篇Flower核心功能全揭秘任务追踪、Worker监控与性能分析完整指南下一篇Kimi-K3-0.40B架构解密混合注意力机制KDA/MLA与MoE技术详解创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
企业数字化 ERP 产品动态
相关推荐
wdcp更改网站域名实操:新手避坑指南与服务器选型哪家好 wdcp更改网站域名实操:新手避坑指南与服务器选型哪家好 域名和服务器,这两个词放在一起,90%的新手站长都会头大。 很多刚接触 Web 环境的老板或者运维小白,手里拿着一个刚买的服务器,看着宝塔面板或者 WDCP… · 2026/9/27 7:54:23
CodeQL 1.26 Python 分析改进深度解读:共享数据流库迁移与污点追踪能力增强 静态分析SAST应用安全漏洞扫描代码质量 【免费下载链接】codeql CodeQL: the libraries and queries that power security researchers around the world, as well as code scanning in GitHub Advanced Security 项目地址: https://gitcode.com/gh_mirrors/co/code… · 2026/9/27 7:54:05
南昌网站设计企业SEO避坑指南:从被黑到排名的实战 南昌网站设计企业SEO避坑指南:从被黑到排名的实战 上周刚帮一个南昌做机械设备的朋友救火。他网站后台突然弹出一堆红色警告,打开浏览器一看,页面底部挂满了赌博和色情链接,百度快照直接显示“该页面可能含有不良信息”。客户急得满头汗,第一反应是:… · 2026/9/27 8:37:54
一文搞懂如何让做树洞网站不挂马且SEO排名稳 一文搞懂如何让做树洞网站不挂马且SEO排名稳 昨天凌晨三点,我接到一个老客户的电话,声音都在抖:“哥,我的树洞网站被黑了,首页全是赌博广告,后台密码改了也没用,现在流量全掉到零了。”这种 网站被黑挂马不知道怎么办… · 2026/9/27 8:37:36
国内织梦和wordpress对比评测:选错CMS,网站做好没人访问 国内织梦和wordpress对比评测:选错CMS,网站做好没人访问 网站做好了没人访问,这不仅是流量焦虑,更是底层架构选型的失败。很多站长盯着页面美观度,却忽略了内容管理系统(CMS)对搜索引擎抓取效率的影响。今天这篇关于… · 2026/9/27 8:37:17
网站标题在哪里改对SEO和性能优化至关重要 网站标题在哪里改对SEO和性能优化至关重要 网站做好了没人访问,往往不是内容不行,而是基础配置没搞对。很多站长盯着页面看半天,找不到 网站标题在哪里 ,导致搜索引擎抓取的关键词全错。更隐蔽的是,错误的标题配置会拖累 性能优化… · 2026/9/27 8:37:11
宁夏网站制作避坑指南:搞定域名服务器,看清真实建站报价 宁夏网站制作避坑指南:搞定域名服务器,看清真实建站报价 你是不是也卡在这一步?手里拿着域名和服务器账号,看着后台那堆英文代码和配置选项,脑子一片空白。很多银川的老板做宁夏网站制作时,最大的焦虑不是设计好不好看,而是根本搞不懂域名解析、SSL… · 2026/9/27 8:37:11
MATLAB雷达信号脉冲压缩仿真:LFM线性调频、匹配滤波与距离分辨率实现 简介:这套Matlab仿真工具完整呈现雷达信号脉冲压缩过程,从线性调频(LFM)信号生成、目标回波仿真到匹配滤波压缩处理均有可运行代码支撑,面向电子信息工程、计算机、数学等专业学生,适用于课程设计、期末大作… · 2026/9/27 0:00:01
汕头网站建设制作厂家避坑指南:5大注意事项救急 汕头网站建设制作厂家避坑指南:5大注意事项救急 改个需求建站公司拖一周,这种憋屈事我见得太多了。 很多汕头老板找本地建站团队,签合同前看着方案挺美,一上线就变脸。 今天不聊虚的,直接拆解找 汕头网站建设制作厂家 时的5个核心 注意事项… · 2026/9/27 0:00:01
多模态虚假新闻检测实战:BERT+ResNet双塔与对比学习 简介:基于PyTorch的多模态虚假新闻检测项目完整代码包,面向自然语言处理与计算机视觉交叉方向的开发者、科研人员及毕业设计选题者,解决社交媒体中文本与图像联合识别虚假新闻的问题。系统以BERT预训练模型提取文本语义特征,以Res… · 2026/9/27 0:00:01
MATLAB雷达信号脉冲压缩仿真:LFM线性调频、匹配滤波与距离分辨率实现 简介:这套Matlab仿真工具完整呈现雷达信号脉冲压缩过程,从线性调频(LFM)信号生成、目标回波仿真到匹配滤波压缩处理均有可运行代码支撑,面向电子信息工程、计算机、数学等专业学生,适用于课程设计、期末大作… · 2026/9/27 0:00:01
汕头网站建设制作厂家避坑指南:5大注意事项救急 汕头网站建设制作厂家避坑指南:5大注意事项救急 改个需求建站公司拖一周,这种憋屈事我见得太多了。 很多汕头老板找本地建站团队,签合同前看着方案挺美,一上线就变脸。 今天不聊虚的,直接拆解找 汕头网站建设制作厂家 时的5个核心 注意事项… · 2026/9/27 0:00:01
多模态虚假新闻检测实战:BERT+ResNet双塔与对比学习 简介:基于PyTorch的多模态虚假新闻检测项目完整代码包,面向自然语言处理与计算机视觉交叉方向的开发者、科研人员及毕业设计选题者,解决社交媒体中文本与图像联合识别虚假新闻的问题。系统以BERT预训练模型提取文本语义特征,以Res… · 2026/9/27 0:00:01