展开菜单
首页 精品内容 本月促销 装机必备 Windows PDF教程 macOS软件 IOS软件 Android AI 专题
全部分类

当前位置:

首页 > 编程开发 > SpringBoot与ApachePulsar集成构建高性能消息系统实践应用案例

SpringBoot与ApachePulsar集成构建高性能消息系统实践应用案例

ApachePulsar采用分层架构实现高吞吐低延迟与持久化存储,通过SpringBoot集成可快速搭建消息系统。集成需配置Pulsar客户端,实现消息发送与消费,并支持分区、批处理、事务及死信队列等高级特性,适用于订单处理与实时数据分析场景,合理配置参数可优化系统性能。

引言

说到分布式系统,消息中间件可以说是整个架构的“交通枢纽”——它让系统组件之间不再紧耦合,同时还能扛住高并发、保证数据不丢。在众多消息中间件里,Apache Pulsar 算是近几年很受关注的一个新面孔。它既继承了 Kafka 的高吞吐能力,又在存储和延迟上做了不少优化,逐渐成了不少企业级项目的首选。今天这篇文章,我们就来聊聊怎么在 Spring Boot 应用里把 Pulsar 集成进来,搭一套高性能的消息系统。

一、Apache Pulsar 简介

1.1 核心特性

  • 高吞吐低延迟:Pulsar 采用分层架构,把存储和计算分离——这个设计很关键,它支持百万级消息吞吐量,延迟能控制在毫秒级。
  • 持久化存储:底层基于 Apache BookKeeper,消息存储可靠,不会丢数据。
  • 多租户支持:内置多租户隔离机制,适合大型企业级应用,不同团队可以共用同一套集群。
  • 灵活的消息模型:发布/订阅模式和队列模式都支持,可以根据业务场景灵活切换。
  • 跨地域复制:支持消息跨数据中心复制,对系统可用性和容灾能力提升很明显。

1.2 架构组成

  • Broker:负责消息的收发、路由和负载均衡,是消息系统的“调度中心”。

二、Spring Boot 集成 Apache Pulsar

2.1 添加依赖

集成第一步,先把依赖加进来。在 pom.xml 里引入 Pulsar 客户端和 Spring Boot Web 的依赖:


    org.apache.pulsar
    pulsar-client
    3.0.0


    org.springframework.boot
    spring-boot-starter-web

2.2 配置 Pulsar 连接

接下来,配置 Pulsar 的连接信息。在 application.yml 里写上服务地址:

spring:
  pulsar:
    client:
      service-url: pulsar://localhost:6650
    admin:
      service-url: http://localhost:8080

2.3 发送消息

直接上代码,创建一个消息发送服务。这里用 @PostConstruct 和 @PreDestroy 来管理客户端和生产者生命周期,省心的做法:

import org.apache.pulsar.client.api.Producer;
import org.apache.pulsar.client.api.PulsarClient;
import org.apache.pulsar.client.api.Schema;
import org.springframework.stereotype.Service;
import ja vax.annotation.PostConstruct;
import ja vax.annotation.PreDestroy;
import ja va.util.concurrent.CompletableFuture;
@Service
public class PulsarProducerService {
    private PulsarClient client;
    private Producer producer;
    @PostConstruct
    public void init() throws Exception {
        client = PulsarClient.builder()
                .serviceUrl("pulsar://localhost:6650")
                .build();
        producer = client.newProducer(Schema.STRING)
                .topic("persistent://public/default/my-topic")
                .create();
    }
    public void sendMessage(String message) throws Exception {
        producer.send(message);
    }
    public CompletableFuture sendAsyncMessage(String message) {
        return producer.sendAsync(message);
    }
    @PreDestroy
    public void close() throws Exception {
        if (producer != null) {
            producer.close();
        }
        if (client != null) {
            client.close();
        }
    }
}

2.4 消费消息

消费端同样简单,用 messageListener 处理消息,消费完记得确认:

import org.apache.pulsar.client.api.Consumer;
import org.apache.pulsar.client.api.PulsarClient;
import org.apache.pulsar.client.api.Schema;
import org.apache.pulsar.client.api.SubscriptionType;
import org.springframework.stereotype.Service;
import ja vax.annotation.PostConstruct;
import ja vax.annotation.PreDestroy;
import ja va.util.concurrent.TimeUnit;
@Service
public class PulsarConsumerService {
    private PulsarClient client;
    private Consumer consumer;
    @PostConstruct
    public void init() throws Exception {
        client = PulsarClient.builder()
                .serviceUrl("pulsar://localhost:6650")
                .build();
        consumer = client.newConsumer(Schema.STRING)
                .topic("persistent://public/default/my-topic")
                .subscriptionName("my-subscription")
                .subscriptionType(SubscriptionType.Exclusive)
                .messageListener((consumer, msg) -> {
                    try {
                        System.out.println("Received message: " + new String(msg.getData()));
                        consumer.acknowledge(msg);
                    } catch (Exception e) {
                        consumer.negativeAcknowledge(msg);
                    }
                })
                .subscribe();
    }
    @PreDestroy
    public void close() throws Exception {
        if (consumer != null) {
            consumer.close();
        }
        if (client != null) {
            client.close();
        }
    }
}

三、高级特性

3.1 消息分区

该说不说,消息分区是提高并行度的好手段。通过指定 key,Pulsar 会把消息路由到对应分区:

producer = client.newProducer(Schema.STRING)
        .topic("persistent://public/default/my-partitioned-topic")
        .create();
// 发送消息到指定分区
producer.newMessage()
        .value("Hello Pulsar")
        .key("key1") // 基于key分区
        .send();

3.2 消息批处理

想要提高吞吐量,批处理是个利器。把多条消息攒在一起发,网络开销少了很多:

producer = client.newProducer(Schema.STRING)
        .topic("persistent://public/default/my-topic")
        .batchingEnabled(true)
        .batchingMaxMessages(1000)
        .batchingMaxPublishDelay(10, TimeUnit.MILLISECONDS)
        .create();

3.3 事务支持

Pulsar 支持事务,这在需要保证消息原子性的时候特别有用。比如一次发送多条消息,要么全部成功,要么全部回滚:

// 开启事务
Transaction txn = client.newTransaction()
        .withTransactionTimeout(1, TimeUnit.MINUTES)
        .build()
        .get();
// 在事务中发送消息
producer.newMessage(txn)
        .value("Hello Transaction")
        .send();
// 提交事务
txn.commit().get();

3.4 死信队列

消息消费失败怎么办?死信队列就是个兜底方案。设置最大重试次数,超过次数就扔到死信主题里,方便后续排查:

consumer = client.newConsumer(Schema.STRING)
        .topic("persistent://public/default/my-topic")
        .subscriptionName("my-subscription")
        .deadLetterPolicy(DeadLetterPolicy.builder()
                .maxRedeliverCount(10)
                .deadLetterTopic("persistent://public/default/my-dlq")
                .build())
        .subscribe();

四、实践应用

4.1 订单处理系统

在订单处理场景里,Pulsar 可以很好地串联起各个服务:

  1. 订单创建时,把订单消息发到 Pulsar
  2. 订单处理服务消费消息,进行后续处理
  3. 处理结果再发到另一个主题,供下游服务使用

4.2 实时数据分析

实时数据分析是另一个典型场景。前端采集的用户行为数据通过 Pulsar 流入,流处理服务实时消费分析,结果写入数据库或缓存:

  1. 前端采集用户行为数据,发送到 Pulsar
  2. 流处理服务消费数据,进行实时分析
  3. 分析结果存储到数据库或缓存

五、性能优化

5.1 生产者优化

  • 启用批处理:减少网络请求次数,提升吞吐量
  • 使用异步发送:不阻塞主线程,发送效率更高
  • 合理设置消息大小:消息太大影响性能,太小又浪费带宽,需要找到平衡点

5.2 消费者优化

  • 批量接收消息:减少网络往返时间
  • 合理设置消费者数量:根据系统负载调整,避免资源浪费或消费积压
  • 使用并发消费:多线程处理消息,提高处理速度

5.3 集群配置优化

  • 增加 Broker 数量:提高系统的处理能力
  • 合理配置 BookKeeper:确保存储性能,避免成为瓶颈
  • 使用负载均衡:均匀分布消息处理压力,防止单点过载

六、常见问题与解决方案

问题原因解决方案
消息发送失败网络连接问题检查网络连接,配置重试机制
消息消费延迟消费者处理速度慢增加消费者数量,优化处理逻辑
系统吞吐量低配置不合理优化批处理设置,调整集群配置
消息丢失未正确处理确认确保消费后正确确认消息

七、总结

坦率说,Apache Pulsar 在消息中间件这个领域里,算是一个后起之秀。它把高吞吐、低延迟、持久化存储这些特性集于一身,特别适合用来构建高性能的分布式系统。通过 Spring Boot 和 Pulsar 的集成,我们可以快速搭建一套可靠的消息系统,满足各种业务场景的需求。

在实际项目中,关键是根据业务场景和系统需求,合理配置 Pulsar 的各项参数,把性能优化到位。同时,可观测性也不能忽视——及时发现和解决问题,才能保证系统稳定运行。

希望这篇文章能帮你更快地上手 Spring Boot 与 Pulsar 的集成。具体怎么用,还得看你的业务场景,灵活运用 Pulsar 的各种特性,才能构建出真正可靠、高效的消息系统。

本站声明:本文内容由网友自发贡献,版权归原作者所有,本站不承担相应法律责任。如您发现有涉嫌抄袭侵权的内容,请联系bd@zhengruan.com
作者最新文章
编程开发
相关文章 更多
精品专题 更多
装机必备

正软商城装机必备专区,精选办公、浏览器、安全防护、影音播放、压缩解压、设计创作和系统工具等电脑常用正版软件,帮助用户快速完成新电脑软件配置。

Windows

正软商城Windows软件专区,汇集适用于Windows电脑的办公、设计、安全防护、影音播放、开发工具和系统优化软件,提供软件介绍、系统要求、正版授权及购买下载服务。

PDF教程

正软商城PDF教程频道提供PDF编辑、转换、合并、拆分、压缩及格式处理方法,同时介绍常用PDF软件和工具的使用技巧。

macOS软件

正软商城macOS软件专区,精选适用于Mac电脑的办公、设计、影音、效率、开发和系统工具,提供软件功能介绍、macOS兼容版本、正版授权及购买下载服务。

IOS软件

正软商城iOS软件专区,精选适用于iPhone和iPad的办公、学习、影音、设计、效率及AI应用,提供功能介绍、适用设备、系统要求和正版获取方式等信息。

AI

正软商城AI软件专区,汇集AI写作、AI绘画、AI视频、AI办公、AI编程、AI翻译、智能客服和数据分析等人工智能工具,提供功能介绍、适用平台、收费方式及正版购买信息。

Mac软件 更多
Shapr3D macOS版
Shapr3D macOS版

Shapr3D是一款面向工业设计、机械工程、建筑概念和三维打印工作流的CAD软件。Mac版采用Parasolid建模内核,支持草图约束、实体建模、工程图、可视化渲染及常见CAD格式交换,并可通过账户在多台设备之间同步项目。

REAPER macOS版
REAPER macOS版

REAPER是Cockos开发的数字音频工作站,提供多轨音频与MIDI录制、剪辑、处理、混音和母带制作工具。Mac版兼容Intel与Apple芯片,支持AU、VST、VST3、CLAP等插件格式,并提供高度可定制的工作流程。

Ableton Live macOS版
Ableton Live macOS版

Ableton Live 是面向音乐制作人与现场表演者的数字音频工作站,提供编曲视图、独具特色的现场视图、音频录制、MIDI创作、实时变速、乐器及效果器。Mac版原生支持Apple芯片,并可连接音频接口、MIDI控制器和第三方插件。

Adobe After Effects macOS版
Adobe After Effects macOS版

Adobe After Effects 是面向动态图形、影视特效与视频合成的专业创作软件。它提供关键帧动画、遮罩抠像、运动跟踪、三维空间、文字动画和表达式等工具,并可与 Premiere Pro、Photoshop、Illustrator

Adobe Audition macOS版
Adobe Audition macOS版

Adobe Audition 是面向播客、影视后期、音乐制作及内容创作者的专业音频工作站,提供波形精修、多轨录制与混音、频谱编辑、降噪修复、响度匹配和效果处理工具,并可与 Adobe Premiere 协同完成视频声音制作。

Adobe Bridge macOS版
Adobe Bridge macOS版

Adobe Bridge 是面向摄影师、设计师及内容团队的免费数字资产管理工具,可集中预览照片、视频和 Adobe 项目文件,并利用评级、标签、关键词、元数据、筛选与收藏集完成分类检索。它还支持批量重命名、格式导出、相机导入及 Camera

Adobe Illustrator macOS版
Adobe Illustrator macOS版

Adobe Illustrator 是 Adobe 推出的专业矢量图形设计软件,适合在 Mac 上制作标志、图标、插画、包装、信息图及印刷版面。它以路径和锚点构建可无损缩放的图形,并提供钢笔、形状、文字、图像描摹、渐变、画板及生成式功能。

Adobe InDesign macOS版
Adobe InDesign macOS版

Adobe InDesign 是面向印刷品与数字出版物的专业版面设计软件,可在 Mac 上制作书籍、杂志、宣传册、海报、报告及交互式文档。它提供精细的文字样式、网格、母版、长文档管理和印前输出能力。

Adobe Lightroom Classic macOS版
Adobe Lightroom Classic macOS版

Adobe Lightroom Classic 是面向桌面摄影工作流的照片管理与后期处理软件,可将原始照片保存在本地硬盘,通过目录、关键词、评分和收藏夹高效整理图库,并提供 RAW 冲印、蒙版、镜头校正、批量同步及多格式导出等功能。

Adobe Media Encoder macOS版
Adobe Media Encoder macOS版

Adobe Media Encoder 是 Adobe 推出的专业音视频编码与转码工具,可通过队列、预设和监视文件夹批量完成格式转换、代理文件创建及多平台成片输出,并与 Premiere Pro、After Effects 等软件紧密协作。

Adobe Photoshop macOS版
Adobe Photoshop macOS版

Adobe Photoshop 是面向摄影师、设计师和内容创作者的专业图像处理软件,提供图层、蒙版、选区、修复、调色、文字排版、智能对象及生成式编辑等能力。Mac 版支持 Intel 与 Apple Silicon 处理器。

Adobe Premiere Pro macOS版
Adobe Premiere Pro macOS版

Adobe Premiere Pro 是面向专业视频制作的非线性剪辑软件,提供多轨时间线、调色、音频处理、字幕、代理工作流和多格式输出能力,并可与 After Effects、Audition、Photoshop 及 Frame.io 协同

WINDOWS 更多
3dmax(3ds max)
3dmax(3ds max)

Autodesk 3ds Max 是一款专业的三维建模、动画与渲染软件,广泛应用于建筑可视化、游戏开发、影视动画、广告设计和产品展示等领域。

photoshop
photoshop

Photoshop 2026 是 Adobe 推出的专业图像处理与视觉设计软件,支持 Windows、macOS 和 iPad 等平台,广泛应用于摄影修图、电商设计、平面海报、数字绘画及视觉合成等创作场景。

Blender
Blender

Blender 是一款免费开源、跨平台的专业 3D 创作软件,集建模、动画、渲染、视频编辑与视觉合成等功能于一体,广泛应用于影视动画、游戏设计和建筑可视化等领域。软件支持 Cycles 物理渲染器与 Eevee 实时渲染引擎,并提供多边形建模、骨骼绑定、物理模拟等专业工具。Blender 兼容 Windows、macOS 和 Linux 系统,安装包轻巧、运行流畅,依托活跃的全球开发者社区持续更新,是从初学者到专业创作者都值得选择的正版 3D 创作工具。

Windows 10
Windows 10

Windows 10 是一款微软推出的经典操作系统,拥有硬件兼容性与多任务处理能力。它更偏向把系统状态查看和常用调节动作放在一起,适合需要持续观察和微调设备状态的场景。

极度公式
极度公式

极度公式是一款跨平台专业LaTeX公式识别编辑软件,支持OCR公式识别和多平台编辑。和使用说明,避免使用,享受完整功能与稳定支持。做扫描整理、文字提取和表格转换时,它能把识别后的处理步骤接得更顺,资料录入这类场景会省下不少时间。

密码键盘
密码键盘

密码键盘是一款兼具安全性与便捷性的高效密码管理器。日常使用里的持续防护和信息管理会更突出,适合把安全控制放进长期使用流程中的场景。

思源笔记
思源笔记

思源笔记是一款本地笔记软件,提供所见即所得的编辑方式,为长文写作带来顺滑的体验。记录、整理和输出之间的过渡会更自然,适合长期写作、做笔记或持续沉淀个人内容。

傲梅轻松备份
傲梅轻松备份

傲梅轻松备份是一款专业易用的数据备份软件,为重要数据提供安全保障。日常使用里的持续防护和信息管理会更突出,适合把安全控制放进长期使用流程中的场景。

Office 365 简体中文
Office 365 简体中文

一款文字处理软件,一种订阅式的跨平台办公软件,基于云平台提供多种服务,通过将 Excel 和 Outlook 等应用与 OneDrive 和 Microsoft Teams 等强大的云服务相结合,Office 365 可让任何人使用任何设备随时随地创建和共享内容。

Mac
Wise Folder Hider Pro
Wise Folder Hider Pro

Wise Folder Hider Pro 是一款专业级文件和文件夹隐藏加密软件,为私密数据添加多重保护。高频操作更强调就近处理,浏览、整理和跨目录移动文件时,来回切换和重复点击都会少很多。

WALTR PRO
WALTR PRO

WALTR是一款电脑至iOS文件传输转换工具,操作简单,快速实现文件识别与传送。做扫描整理、文字提取和表格转换时,它能把识别后的处理步骤接得更顺,资料录入这类场景会省下不少时间。

CodeExpander
CodeExpander

CodeExpander 是一款快捷短语输入增强工具,通过键入缩写自动展开为自定义文段,提升工作效率。任务管理和过程控制会更完整,持续下载、批量同步或需要稳定传输流程的场景会更适合它。