一、核心核心价值:异步、解耦、削峰
1.1 业务痛点
传统同步业务架构中,主流程需要串行执行多个非核心附属业务,存在三大核心问题:
- 耦合严重:新增业务需要修改主流程代码,迭代风险高、扩展性差
- 响应缓慢:主流程阻塞等待附属业务执行,接口rt变长,用户体验差
- 流量雪崩:瞬时高并发流量直接打垮后端服务,无缓冲机制
- 数据不一致:主事务成功、附属业务失败,导致业务数据状态错乱
1.2 rocketmq核心优势
rocketmq 是阿里开源的分布式消息中间件,天然适配微服务架构,完美解决上述痛点,核心能力对应业务价值:
- 解耦:生产者只负责发消息,消费者独立消费处理,上下游服务无代码依赖、无调用依赖
- 异步:主流程发送消息后立即返回,无需等待下游执行,大幅提升接口响应速度
- 削峰填谷:瞬时高流量缓存至消息队列,消费者匀速消费,保护后端服务稳定性
- 高可靠:支持消息重试、死信队列、事务消息,保障分布式场景下最终数据一致性
二、基础异步解耦实现方案(springboot 整合 rocketmq)
2.1 环境依赖(maven)
采用springboot官方适配的rocketmq starter,适配主流2.x、3.x版本:
xml
<dependency>
<groupid>org.apache.rocketmq</groupid>
<artifactid>rocketmq-spring-boot-starter</artifactid>
<version>2.2.3</version>
</dependency>2.2 配置文件(application.yml)
yaml
rocketmq:
# nameserver地址,集群多个地址用分号分隔
name-server: 127.0.0.1:9876
# 生产者配置
producer:
group: business-producer-group
# 消息发送超时时间
send-timeout: 3000
# 异步发送失败重试次数
retry-times-when-send-failed: 2
# 消费者配置
consumer:
group: business-consumer-group2.3 异步消息发送(生产者实现)
核心:主业务事务执行完成后,异步发送消息,不阻塞主流程,实现业务解耦。提供普通异步发送、带回调异步发送两种方式。
import org.apache.rocketmq.spring.core.rocketmqtemplate;
import org.springframework.beans.factory.annotation.autowired;
import org.springframework.stereotype.service;
import org.apache.rocketmq.client.producer.sendcallback;
import org.apache.rocketmq.client.producer.sendresult;
@service
public class businessproducerservice {
@autowired
private rocketmqtemplate rocketmqtemplate;
// 消息主题
private static final string business_topic = "business_async_topic";
/**
* 异步发送消息(无阻塞,主流程直接返回)
*/
public void sendasyncmessage(object message) {
// 异步发送,带回调机制,处理发送成功/失败场景
rocketmqtemplate.asyncsend(business_topic, message, new sendcallback() {
// 发送成功回调
@override
public void onsuccess(sendresult sendresult) {
// 可记录消息发送日志、消息id,用于后续追踪
system.out.println("消息发送成功,msgid:" + sendresult.getmsgid());
}
// 发送失败回调
@override
public void onexception(throwable e) {
// 发送失败兜底:记录异常日志、人工重试、定时任务补偿
system.err.println("消息发送失败:" + e.getmessage());
}
});
}
}2.4 消息消费(消费者实现)
消费者独立监听主题,异步处理附属业务,与主业务完全解耦,支持重试机制。
import org.apache.rocketmq.spring.annotation.rocketmqmessagelistener;
import org.apache.rocketmq.spring.core.rocketmqlistener;
import org.springframework.stereotype.service;
@service
@rocketmqmessagelistener(
topic = "business_async_topic",
consumergroup = "business-consumer-group"
)
public class businessconsumerservice implements rocketmqlistener<string> {
/**
* 消费方法:消费成功自动签收,异常自动重试
* @param message 消息体
*/
@override
public void onmessage(string message) {
// 执行解耦后的附属业务:如日志记录、消息通知、数据统计、积分发放等
system.out.println("异步消费业务消息:" + message);
// 业务异常会被rocketmq捕获,触发重试机制
}
}2.5 基础方案问题总结
上述基础异步解耦方案,能解决解耦、异步、削峰问题,但存在致命缺陷:无法保证最终一致性,典型异常场景:
- 主业务数据库事务提交成功,消息发送失败,下游消费者未执行,导致多做事务、少做附属业务
- 消息发送成功,消费者消费失败,无可靠补偿机制,数据状态不一致
- 主事务回滚,但消息已发送,下游执行了无效业务,导致数据错乱
三、分布式最终一致性核心解决方案
微服务分布式场景下,最终一致性核心目标:主db事务 和 消息发送 要么同时成功,要么同时失败,下游消费最终一定执行成功。
rocketmq 专属解决方案:rocketmq 事务消息(半消息机制),是解决消息与数据库事务一致性的最优方案,优于本地消息表、可靠消息最终一致性方案。
3.1 事务消息核心原理
rocketmq 事务消息分为三个阶段,完美绑定数据库事务与消息发送:
- 发送半消息(预处理):生产者发送半消息到mq,半消息对消费者不可见,不会被消费
- 执行本地事务:半消息发送成功后,生产者执行自身数据库业务事务
- 事务状态回查/提交:
- 本地事务成功:向mq发送提交指令,消息对消费者可见,允许消费
- 本地事务失败:向mq发送回滚指令,mq删除半消息,无消息产生
- 状态未知(超时、宕机):mq定时回查生产者本地事务状态,根据结果提交/回滚
3.2 事务消息完整落地实现
3.2.1 生产者事务监听实现(核心)
import org.apache.rocketmq.spring.annotation.rocketmqtransactionlistener;
import org.apache.rocketmq.spring.core.rocketmqlocaltransactionlistener;
import org.apache.rocketmq.spring.core.rocketmqlocaltransactionstate;
import org.apache.rocketmq.client.producer.transactionsendresult;
import org.springframework.beans.factory.annotation.autowired;
import org.springframework.stereotype.service;
import org.springframework.transaction.annotation.transactional;
@service
@rocketmqtransactionlistener(txproducergroup = "business-tx-producer-group")
public class businesstxtransactionlistener implements rocketmqlocaltransactionlistener {
@autowired
private businessservice businessservice;
private static final string tx_topic = "business_tx_topic";
/**
* 阶段1:发送半消息后,执行本地数据库事务
*/
@override
@transactional(rollbackfor = exception.class)
public rocketmqlocaltransactionstate executelocaltransaction(org.apache.rocketmq.common.message.message msg, object arg) {
try {
// 解析消息体,执行核心本地业务事务(db操作)
string message = new string(msg.getbody());
businessservice.executecorebusiness(message);
// 本地事务执行成功,返回提交状态
return rocketmqlocaltransactionstate.commit;
} catch (exception e) {
// 本地事务失败,返回回滚状态,mq删除消息
return rocketmqlocaltransactionstate.rollback;
}
}
/**
* 阶段2:mq定时回查本地事务状态(解决宕机、超时未知场景)
*/
@override
public rocketmqlocaltransactionstate checklocaltransaction(org.apache.rocketmq.common.message.message msg) {
// 根据消息唯一标识,查询数据库事务执行状态
string msgid = msg.getmsgid();
boolean issuccess = businessservice.checkbusinessstatus(msgid);
if (issuccess) {
return rocketmqlocaltransactionstate.commit;
}
return rocketmqlocaltransactionstate.rollback;
}
/**
* 对外暴露:发送事务消息入口
*/
public transactionsendresult sendtxmessage(string message) {
return rocketmqtemplate.sendmessageintransaction(tx_topic, message, null);
}
}3.2.2 核心业务服务(带事务状态记录)
import org.springframework.stereotype.service;
import org.springframework.transaction.annotation.transactional;
@service
public class businessservice {
// 模拟数据库事务执行
@transactional(rollbackfor = exception.class)
public void executecorebusiness(string message) {
// 1. 执行核心业务db操作
// 2. 记录事务状态(用于mq回查):存储msgid、事务状态、创建时间
// db.savetxlog(msgid, success);
}
// 回查接口:校验本地事务是否执行成功
public boolean checkbusinessstatus(string msgid) {
// 根据msgid查询事务日志表,判断事务状态
// txlog log = db.getbymsgid(msgid);
// return log != null && log.getstatus().equals(success);
return true;
}
}3.2.3 事务消息消费者(通用消费逻辑)
import org.apache.rocketmq.spring.annotation.rocketmqmessagelistener;
import org.apache.rocketmq.spring.core.rocketmqlistener;
import org.springframework.stereotype.service;
@service
@rocketmqmessagelistener(topic = "business_tx_topic", consumergroup = "business-tx-consumer-group")
public class businesstxconsumer implements rocketmqlistener<string> {
@override
public void onmessage(string message) {
// 最终一致性保障:只有主事务成功,消息才会被消费
// 执行下游解耦业务,消费异常自动重试
doasyncbusiness(message);
}
private void doasyncbusiness(string message) {
// 附属业务逻辑:通知、统计、积分、日志等
}
}3.3 最终一致性完整闭环方案(全异常兜底)
仅靠事务消息无法覆盖所有场景,需搭配重试机制、死信队列、定时补偿、幂等性形成闭环,实现100%最终一致。
3.3.1 消费者幂等性保障(必做)
rocketmq会重试消费消息,必须防止重复消费导致业务异常,幂等方案:
- 基于消息唯一id(msgid)做全局去重
- 基于业务唯一主键(订单号、用户id)做幂等校验
- 消费成功后记录消费日志,重复消息直接跳过
3.3.2 消息重试机制配置
消费者业务异常时,rocketmq自动重试,默认重试16次,间隔递增,可自定义配置:
- 业务异常(非系统异常):抛出异常,触发重试
- 无需重试异常:捕获异常,正常返回,不触发重试
3.3.3 死信队列兜底
消息重试次数耗尽仍消费失败,自动进入死信队列,不会丢失消息:
- 单独监听死信主题,接收异常消息
- 记录异常日志,触发人工告警
- 修复问题后,手动重试死信消息
3.3.4 定时任务补偿
针对极端异常(mq宕机、消息丢失、回查失败),新增定时补偿任务:
- 定时查询事务日志表,筛选主事务成功、未消费的业务数据
- 自动补发消息或直接执行下游业务
- 对账校验,修复数据不一致问题
四、核心场景方案选型对比
方案 | 优点 | 缺点 | 适用场景 |
普通异步消息 | 实现简单、性能高、完全解耦 | 无事务一致性保障,易数据错乱 | 非核心业务、无需强一致场景(日志、统计) |
rocketmq事务消息 | 无侵入、高性能、原生支持最终一致、无需本地消息表 | 实现稍复杂,需处理回查、幂等 | 核心业务、分布式事务、需要数据一致场景(订单、支付) |
本地消息表 | 通用性强、适配所有mq | 侵入业务、需维护数据表、性能差 | 老旧项目兼容、非rocketmq环境 |
五、生产环境避坑规范
- 禁止同步发核心消息:核心业务一律使用事务消息,非核心业务使用异步消息,杜绝同步阻塞
- 必须做消费幂等:所有消费者强制实现幂等,避免重试导致数据重复、脏数据
- 事务回查必实现:事务消息必须实现本地事务状态查询,否则宕机场景会导致消息悬挂
- 区分异常类型:参数异常、业务校验异常直接捕获,不重试;系统异常、超时异常触发重试
- 消息体轻量化:禁止传输大文本、文件,消息体仅存业务唯一标识,业务数据从db查询
- 日志全链路追踪:记录msgid、traceid,实现消息发送、消费、异常全链路溯源
六、方案总结
1. 异步解耦能力:通过rocketmq普通异步消息,实现主从业务分离,提升接口响应速度、解除服务耦合、实现流量削峰,适配绝大多数非核心业务场景。
2. 最终一致性能力:通过rocketmq事务消息+幂等消费+死信队列+定时补偿的闭环方案,彻底解决分布式场景下,数据库事务与消息发送的一致性问题,实现业务数据最终一致。
3. 落地价值:方案无强业务侵入、性能优异、可靠性高,是微服务架构下异步业务解耦、分布式最终一致性的最优落地方案,可直接用于订单、支付、积分、通知等核心业务。
以上就是springboot+rocketmq异步解耦与最终一致性解决方案的详细内容,更多关于springboot rocketmq异步解耦与一致性的资料请关注代码网其它相关文章!
发表评论