当前位置: 首页 > news >正文

Pulsar 消息同步机制

Pulsar 就是一个消息中间件(消息队列),和 RabbitMQ、Kafka 是同一类东西。它解决的核心问题就一个:让两个服务之间不用直接互相调用,而是通过一个"信箱"传话

放在你的项目里,场景特别具体:

谁发给谁:雇员服务在导入员工数据后,需要通知薪资服务这个员工的信息变了,你同步一下。但雇员服务不想、也不能直接去调薪资服务的接口(服务解耦、避免互相依赖、避免一方挂了拖垮另一方)。于是雇员服务把变更的员工数据"投进信箱",薪资服务自己去"取信"处理。

这个"信箱"就是 Pulsar

三个关键概念,对应你代码里的位置:

  1. Topic(主题/频道)—— 信箱的名字。所有发给这个频道的消息,订阅它的人都能收到。

  2. Producer(生产者)—— 寄信人。负责往 topic 里塞消息。核心就是send()方法把对象转 JSON 发出去。雇员服务那边也有对应的 producer。

  3. Consumer(消费者)—— 取信人。负责从 topic 里拿消息并处理。就是专门收员工消息、同步自然人数据的消费者。

用你的街道数据走一遍完整流程:

雇员服务导入 Excel(/import接口)→ 写好EmployeeBasicInfo(含residenceStreet/resideStreet)→ 落自己的库 →通过 pulsar producer 把这个 Employee 对象发到employeetopic

薪资服务启动时, 就连上 Pulsar,订阅这个 topic。每当有新消息进来,基类messageListener回调被触发:

  • 把消息体反序列化成Employee对象
  • 校验通过就调handlerData(object)这个方法各消费者自己重写
  • 处理完consumer.acknowledge(msg)告诉 Pulsar"这条我处理好了,可以删了";处理失败negativeAcknowledge会让 Pulsar 稍后重投。

handlerDataIndividualIncomeTaxConsumer里就是那个把Employee/EmployeeBasicInfo的字段往IndividualDeclareInfo上搬的逻辑

这也正好解释了街道为什么断在同步层handlerData搬地址时搬了省市区和详细地址,却漏了residenceStreet/resideStreet两行——而薪资服务的 feign 副本EmployeeBasicInfo里也没声明这两个字段,所以即便想搬也取不到值。这就是上次说的缺口。

两个配置开关,开发时很实用,值得你知道:

  • offFlag:设成1完全跳过 pulsar 初始化(本地没 pulsar 环境也能把服务跑起来,不用连那台180.100.212.181。开发环境一般开着,生产关掉。
  • 基类里还配了死信队列:一条消息连续失败 3 次(maxRedeliverCount(3))就不再重投,扔到xxx-DLQ主题里,避免一条坏消息把消费者卡死。这是消息队列的标配容错机制。

一句话总结:Pulsar 是雇员服务和薪资服务之间的传声筒,雇员那边改了数据,通过 topic 喊一声,薪资服务听到后自己同步——双方互不直接调用。街道数据的断点不在 Pulsar 本身,而在薪资服务收到消息后"搬字段"那步漏了两行。

一、Pulsar 是什么

Pulsar 是一个消息中间件(消息队列),和 RabbitMQ、Kafka 属同一类。核心作用就一个:让两个服务之间不用直接互相调用,而是通过一个"信箱"传话。

本系统场景:雇员服务导入或变更员工数据后,要通知薪资服务同步自然人报送信息。雇员服务不直接调薪资服务的接口,而是把变更数据投递到 Pulsar,薪资服务自己订阅消费。好处是服务解耦——一方故障不拖垮另一方,双方也无直接代码依赖。

服务地址在

pulsar: serviceUrl: pulsar://180.100.212.181:6650

二、三个核心概念

概念含义类比
Topic消息频道信箱名字
Producer发送消息寄信人
Consumer接收处理消息取信人

员工主题是persistent://public/salary/employee

Producer基类把对象转 JSON 发出。跨服务场景真正的发送入口是雇员服务

它组装EmployeeSyncVO后调 SalaryProviderService.java

Consumer基类 用泛型<T>,子类重写validate()handlerData()。[IndividualIncomeTaxConsumer.java](D:\\EngmaProject\\salary-system\\src\\main\\java\\com\\engma\\salary\\service\\pulsar\\IndividualIncomeTaxConsumer.java) 就是同步自然人信息的消费者,泛型是Employee

三、消息消费流程

启动 → 连接 Pulsar → 订阅 topic → 注册 messageListener 回调 │ 每条消息到达时触发 │ ┌───────────┴────────────┐ │ 1. 反序列化为泛型对象 T │ │ 2. validate() 校验 │ │ 3. handlerData() 处理 │ │ 4. acknowledge() 确认 │ └────────────────────────┘
步骤行号说明
初始化连接49PulsarClient.builder().serviceUrl(serviceUrl).build()
创建消费者61client.newConsumer().topic(topic).subscribe()
接收回调64messageListener((consumer, msg) -> {...})
反序列化72JSONObject.parseObject(data, tClass)
业务分发78handlerData(object)调子类实现
成功确认80consumer.acknowledge(msg)告诉 Pulsar 可删
失败重投75consumer.negativeAcknowledge(msg)稍后重发

确认机制(ack/nack):ack 表示处理成功、消息可删;nack 表示失败、1 秒后重投

死信队列(第65-68行):连续失败 3 次的消息不再重投,扔进{topic}-DLQ,避免坏消息卡死消费者。

共享订阅(第68行SubscriptionType.Shared):同一订阅名下多消费者实例间负载均衡。

四、跨服务同步链路(以街道为例)

┌─────────────── 雇员服务 (employee) ───────────────┐ │ 1. POST /import 上传 Excel │ │ └─ EmployeePayTaxesPlusExcel 解析报税 sheet │ │ 含「户籍所在地(街道)」「居住地街道」两列 │ │ 2. EmployeePayTaxesPlusExcelListener.invoke() │ │ └─ buildBasicInfoForPayPlusTaxes() 搬字段 │ │ info.setResidenceStreet(...) 户籍街道 │ │ info.setResideStreet(...) 居住街道 │ │ 3. updateBatchById() 落库 employee_basic_info │ │ 4. NaturalReportHandle.sendDataToPulsar() │ │ └─ SalaryProviderService.send() 发到 employee topic │ └───────────────────────┬───────────────────────────┘ │ Pulsar 消息(JSON) ▼ ┌─────────────── 薪资服务 (salary-system) ──────────┐ │ 5. IndividualIncomeTaxConsumer 收到消息 │ │ └─ 反序列化为 Employee → handlerData() │ │ └─ 搬字段到 IndividualDeclareInfo │ │ (216-227行搬地址,漏了街道) │ │ 6. 落库 individual_declare_info │ │ 7. /exportList → getDatas() 按 code 取值 │ │ └─ residenceStreet / streetOfResidence │ └───────────────────────────────────────────────────┘

各环节代码位置:

特殊处理:省/市/区在导入时经isAddress()([EmployeePayTaxesPlusExcelListener.java:441](D:\\EngmaProject\\employee\\src\\main\\java\\com\\engma\\employee\\excel\\EmployeePayTaxesPlusExcelListener.java))做级联校验、取标准码值;街道不校验,直接原值透传,因为乡镇街道这级没有统一编码标准。

五、街道数据断链问题

现象:导出 Excel 里两列街道始终为空。

根因:链路在薪资服务搬字段那步断了,两个原因:

  1. feign 副本 [EmployeeBasicInfo.java](D:\\EngmaProject\\salary-system\\src\\main\\java\\com\\engma\\salary\\entity\\feign\\EmployeeBasicInfo.java) 没声明residenceStreet/resideStreet,反序列化取不到值
  2. 消费者 [IndividualIncomeTaxConsumer.java:216-227](D:\\EngmaProject\\salary-system\\src\\main\\java\\com\\engma\\salary\\service\\pulsar\\IndividualIncomeTaxConsumer.java) 搬地址时漏了街道两行

修复(改薪资服务两处):feign 副本补两个字段,消费者第227行后补两行info.setResidenceStreet(basicInfo.getResidenceStreet())info.setStreetOfResidence(basicInfo.getResideStreet())。注意字段名差异:雇员侧resideStreet对应薪资侧streetOfResidence

六、配置开关

offFlag([salary-dev.yaml:65](D:\\EngmaProject\\salary-system\\src\\main\\resources\\salary-dev.yaml)):1=关闭 pulsar 初始化,本地无 pulsar 也能启动;0=开启。开发一般设 1,生产必须设 0。

mockPaidSwitch(第67行):1=注入 mock 数据,0=关闭。

七、Topic 命名坑(重要)

两端 topic 命名不一致:

  • 薪资服务订阅persistent://public/salary/employee(无后缀)
  • 雇员服务发送persistent://public/salary/employee-{active}active默认dev

默认配置下雇员发到...-dev,薪资服务订阅的是无后缀版本,收不到。排查同步问题优先确认两端 topic 是否匹配。

八、组件一览

消费者:IndividualIncomeTaxConsumer(自然人同步)、LaborContractHandleConsumer(劳动合同)、IndividualContractRenewConsumer(合同续签)、SalaryBatchPayInfoConsumer(薪资发放)、AttendanceHandleConsumer(考勤)。

生产者:PulsarProducerService(基类)、AttendanceUsedStatusProducer、SalaryBatchCalculatedProducer、RecruitSalaryChannelProducer。

九、排查指引

消费失败:日志搜[PULSAR]message error;连续失败 3 次进死信队列{topic}-DLQ;确认两端 topic 一致。

本地开发offFlag=1跳过初始化即可启动。

同步验证:导入带街道的数据 → 看日志有无PULSAR---msg---确认收到消息→ 查表residence_street有无值 → 调导出接口验证 Excel。

http://www.jsqmd.com/news/1250596/

相关文章:

  • 矿山破碎设备液压监测采购,液压传感器十大厂家推荐,广东犸力靠谱使用寿命更长 - 品牌速递
  • APS 需求计划(Demand Planning)技术拆解
  • Go http.Client 实战:连接池复用、超时设置与 context 取消
  • HarmonyOS7 数据持久化:Preferences 和 RDB 到底选哪个?
  • 广州财税服务行业GEO优化公司选型指南丨生成式引擎优化服务商深度测评2026本地榜单解析 - 企业新闻快传
  • 59.嵌入式C语言高级宏定义实战:多行宏、字符串化与符号拼接
  • 【RT-DETR涨点改进】TGRS 2026 | 特征融合改进篇 |引入CDSF跨域协同融合模块,增强特征互补性与语义一致性,助力高光谱目标检测、遥感目标检测、多模态融合目标检测任务,高效涨点
  • NTC 热敏电阻全解:计算、曲线读数、跨品牌替换一次讲透
  • 2026年陕西学化妆选校攻略:正规权威机构甄别与避坑指南 - 产业观察报
  • 2026 年现阶段泾县评价高的膜结构张拉膜施工厂家推荐公司怎么联系,颠覆认知:膜结构张拉,比你想象的更省钱! - 行业鉴选官
  • 椰林海鲜码头企业新愿景是什么?:良正基业永昌 - 18002239949
  • 智能装备整机采购指南,2026转矩传感器品牌排行更新,广东犸力销量排名逐年攀升 - 品牌速递
  • ESC框架知识点
  • D2D商业化初期的博弈:成本、价格与增长的“三重奏”
  • 沁园春·数智潮
  • 汽车轮重检测仪品牌排名汇总揭晓,浙江润鑫头部品牌一致好评实力获市场认可 - 品牌速递
  • 老板必看!AI落地就做这两件事,ROI瞬间翻倍!
  • 联合利华CIO任命揭示企业数字化转型战略
  • 多元微积分
  • 【RT-DETR涨点改进】CVPR 2026 | 卷积创新改进篇 | 轻量化改进!引入YOLO-ULM中轻量D3C2f轻量化模块,含二次创新模块,5种创新改进点,助力目标检测任务,高效涨点
  • 浪琴天津售后服务中心2026年7月最新地址与客服热线重磅发布 - 浪琴服务中心
  • 2026西安美甲创业培训学校选型盘点及避坑指南:哪家正规靠谱且适配本地创业需求?附正规机构实力解析 - 行业观察网
  • LLM-Pruner: On the Structural Pruning of Large Language Models 解读
  • Android Camera Sensor Mode 分类与选择详解
  • [godot ]
  • linux_x86_64 虚拟机交叉编译aarch64平台qt源码
  • 一位直博生的 AI 编程困境,会语法,却写不出项目?|AI悦创VibeCoding/Python一对一辅导分享
  • HarmonyOS7 传感器开发:加速度计做个摇一摇,原来这么简单
  • 2026 年 7 月新发布:常熟有实力的无机纤维喷涂工程公司电话实力厂家推荐,揭秘:高强度无机纤维喷涂的秘密工厂电话曝光-佑辰隔音保温 - 鉴选官
  • WordPress国际化多语言建站插件:解决Hreflang标签的3个报错