分布式事务:本地消息表实践
背景
在分布式业务系统中,网络抖动、下游超时、服务重启很容易导致业务状态不一致。这就是分布式事务要解决的问题,解决方案有多种,如题这里要介绍的是本地消息表的实现或者其变种。
我们信贷项目里授信、用信提交是个同步的重接口,内含多个下游服务调用(如创建合同这种耗时操作),迭代过程中暴露出多种问题。
- 网络抖动、下游超时报错,要依赖用户或者上游去主动重试,系统健壮性、容错能力差
- 接口耗时逐渐增加,用户体验差、合作方给的接口SLA要求不满足等
- 卡单后,人工运维成本高
亟待提升系统稳定性问题,要实现同步接口改造成收单模式,并保障分布式事务的最终一致性。相比于重运维的MQ事务消息、高侵入的Seata分布式事务,本地消息表是轻量易落地的实现方案。
这里我们设计并实现了抢跑+定时兜底,无MQ、分布式定时任务等中间件依赖的架构方案。方案兼顾了低延迟、高可用、低锁竞争,其中涉及多种知识点,如分布式事务、并发、线程池、重试、幂等、数据库锁、监控等。
一、方案介绍
1.1 本息消息表的实现思路
建一张本地重试任务表,将收单动作,和下游调用的任务作为一个事务,持久化到本地数据库。
通过新建任务异步抢跑保障响应低延迟,通过定时轮询兜底保障可靠性,下游接口至少一次调用,依靠整个链路幂等保证最终一致性。
1.2 方案设计
双链路互补,各司其职
-
主链路(抢跑执行):业务落库成功后,事务提交完毕立即异步执行,毫秒级低延迟,承担全部的正常业务流量。失败回退后统一走定时重试。
-
兜底链路(定时轮询):10s一次定时扫描,处理异常丢失任务(如线程池拒绝、重启IP变更等),保障百分百不丢任务
-
MIS运维:到达重试次数后,任务依然未成功处理(系统原因需人工运维)。
二、数据库表结构
CREATE TABLE `t_retry_task` (
`id` bigint NOT NULL AUTO_INCREMENT COMMENT '主键ID',
`event_type` varchar(64) NOT NULL COMMENT '事件类型:授信、收单等,绑定独立业务处理器和线程池',
`biz_key` varchar(64) NOT NULL COMMENT '业务唯一标识,订单ID/流水号,用于问题排查',
`biz_context` text NOT NULL COMMENT '业务JSON上下文,任务执行参数',
`status` tinyint NOT NULL COMMENT '0待执行,1成功,2失败终态,3执行中',
`retry_count` int NOT NULL DEFAULT 0 COMMENT '已重试次数',
`max_retry_count` int NOT NULL COMMENT '最大允许重试次数',
`retry_strategy` varchar(32) NOT NULL COMMENT '重试策略:FIXED固定间隔,EXPONENTIAL指数退避',
`first_retry_interval` int NOT NULL COMMENT '首次重试间隔(秒)',
`next_retry_time` datetime(3) NOT NULL COMMENT '下次可调度时间',
`ip` varchar(40) DEFAULT NULL COMMENT '当前抢占实例IP,仅用于本机捞取任务',
`create_time` datetime(3) NOT NULL DEFAULT CURRENT_TIMESTAMP(3) COMMENT '创建时间',
`update_time` datetime(3) NOT NULL DEFAULT CURRENT_TIMESTAMP(3) ON UPDATE CURRENT_TIMESTAMP(3) COMMENT '更新时间',
PRIMARY KEY (`id`),
-- 业务唯一约束:同一个event‑type下,biz_key不能重复,防止重复创建相同任务
UNIQUE KEY `uk_event_biz` (`biz_key`,`event_type`),
-- 定时调度过滤索引
KEY `idx_next_retry` (`next_retry_time`),
KEY `idx_update_time` (`update_time`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='本地任务重试表';
2.1 状态枚举定义
| 状态值 | 状态名称 | 说明 |
|---|---|---|
| 0 | 待执行 | 等待抢跑或定时调度执行 |
| 1 | 执行成功 | 终态,不再参与调度 |
| 2 | 失败终态 | 达最大重试次数,等待人工介入 |
| 3 | 执行中 | 已被实例抢占,正在执行业务 |
2.2 重试策略规则
-
FIXED 固定间隔:每次重试间隔 = 首次间隔时间
-
EXPONENTIAL 指数退避:第N次重试间隔 = 首次间隔 * 2^(n-1)
-
防护机制:指数退避设置最大间隔上限(24小时),防止时间爆炸
三、业务层统一抽象
每种业务事件(授信、用信)对应独立处理器、独立线程池,完全解耦,支持动态扩展业务类型。
3.1 任务处理器统一接口
public interface RetryTaskHandler {
// 匹配数据库event_type
String getEventType();
// 每个业务独立线程池,实现线程池隔离
Executor getExecutor();
// 业务核心处理逻辑
boolean handle(String bizContext);
// TODO 可考虑增加更新主单表的接口方法
// 主单表状态更新 + 重试表状态 一个事务
// 业务默认重试元数据
RetryMeta getRetryMeta();
// 可选:执行前预判,业务已完成则直接跳过重试
default boolean preCheckSkip(String bizKey) {
return false;
}
// 重试配置元数据
class RetryMeta {
private Integer maxRetryCount;
private String retryStrategy;
private Integer firstRetryInterval;
// getter/setter
}
}
3.2 业务实现示例
每个业务单独实现接口,自定义重试策略和线程池,互不影响。
@Component
public class PayOrderRetryHandler implements RetryTaskHandler {
@Override
public String getEventType() {
return "PAY_ORDER";
}
@Override
public Executor getExecutor() {
// 收单业务独立线程池
return payTaskExecutor;
}
@Override
public boolean handle(String bizContext) {
// 执行业务收单回调逻辑
return true;
}
@Override
public RetryMeta getRetryMeta() {
RetryMeta meta = new RetryMeta();
// 配置中心配置
meta.setMaxRetryCount(8);
meta.setRetryStrategy("EXPONENTIAL");
meta.setFirstRetryInterval(10);
return meta;
}
}
四、双链路架构实现
4.1 主链路:事务提交后异步抢跑(低延迟)
本地事务完全提交后,再提交异步任务。
// 业务主事务执行完毕,插入重试任务
payorder.insert(xx);
Long taskId = retryTaskMapper.insert(task);
// 上边事务提交后,提交线程池异步执行
asyncGrabRun(taskId);
异步抢跑逻辑(单行抢占,零锁冲突):
public void async(Long taskId) {
// 1. 单行抢占:粒度最小,无锁竞争
int rows = retryTaskMapper.grabTaskById(taskId, localIp);
if (rows == 0) {
// 已被其他实例抢占,直接退出,避免重复执行
return;
}
// 2. 执行业务预判,已完成则直接标记成功。查主订单状态
RetryTask task = retryTaskMapper.selectById(taskId);
if (handler.preCheckSkip(task.getBizKey())) {
retryTaskMapper.markSuccess(taskId);
return;
}
boolean success = false;
try {
success = handler.handle(task.getBizContext());
} catch (Throwable e) {
log.error("x", e);
success = false;
}
// 3. 根据执行结果流转状态
if (success) {
retryTaskMapper.markSuccess(taskId);
} else {
// 计算下次重试时间
LocalDateTime nextTime = calcNextRetryTime(
LocalDateTime.now(),
task.getRetryCount() + 1,
task.getRetryStrategy(),
task.getFirstRetryInterval()
);
if (task.getRetryCount() + 1 >= task.getMaxRetryCount()) {
retryTaskMapper.markFailFinal(taskId);
} else {
retryTaskMapper.markFailRetry(taskId, nextTime);
}
}
}
单行抢占SQL(WHERE id=#{taskId} AND status=0):
UPDATE t_retry_task
SET status=3, ip=#{ip}, update_time=NOW()
WHERE id=#{taskId} AND status=0;
4.2 兜底链路:定时任务
关键点:无XXL-Job分布式任务组件,依靠Spring Scheduled多实例执行,利用MySQL天然锁机制 + IP字段实现分布式争抢。定时仅兜底处理异常任务,不承载直接业务流量。
设计IP字段,能保障多实例运行时,不会互相争抢任务重复执行。
@Component
public class RetryTaskScheduler {
// 10s固定延迟执行,上一轮执行完毕再等待10s,杜绝任务堆积
@Scheduled(fixedDelay = 10000, initialDelay = 30000)
public void pollTask() {
try {
// 1. 僵死任务回收:处理机器crash、重启后IP变更后未执行情况
recycleDeadTask();
// 2. 单次抢占
int affectRows = retryTaskMapper.grabTaskBatch(localIp);
if (affectRows <= 0) {
return;
}
// 3. 只捞取本机任务,分发线程池执行
List<RetryTask> taskList = retryTaskMapper.selectLocalGrabbedTask(localIp);
taskList.forEach(this::dispatchTask);
} catch (Exception e) {
log.error("轮询任务异常", e);
//监控告警埋点
}
}
}
1)僵死任务回收SQL(3分钟超时可配置)
任务已经在处理中,但时间已经过了一个阈值时间。可能正在处理中,但确实超时了;也可能是任务僵死了(比如ip更新完了,执行前机器重启,导致查不到这些任务)。
UPDATE t_retry_task
SET status = 0, ip = NULL, update_time = NOW()
WHERE status = 3
AND retry_count < max_retry_count
AND update_time < DATE_SUB(NOW(), INTERVAL 3 MINUTE)
ORDER BY next_retry_time ASC, id ASC
LIMIT 50;
2)定时批量抢占SQL(小批量、有序、防死锁)
UPDATE t_retry_task
SET ip = #{ip}, status = 3, update_time = NOW()
WHERE next_retry_time <= NOW()
AND retry_count < max_retry_count
AND status = 0
ORDER BY next_retry_time ASC, id ASC
LIMIT 50;
3)查出来本机待执行任务 SQL
SELECT * FROM t_retry_task
WHERE ip = #{ip}
AND status = 3
ORDER BY next_retry_time ASC, id ASC
LIMIT 50;
五、状态流转SQL
where 条件包含 id,且有乐观锁 status
5.1 执行成功
UPDATE t_retry_task
SET status = 1, update_time = NOW()
WHERE id = #{xx} AND status = 3;
5.2 执行失败、未达最大次数
UPDATE t_retry_task
SET status = 0,
ip = NULL,
next_retry_time = #{calc_next_time},
retry_count = retry_count + 1,
update_time = NOW()
WHERE id = #{xx} AND status = 3;
5.3 执行失败、达到最大次数
UPDATE t_retry_task
SET status = 2,
retry_count = retry_count + 1,
update_time = NOW()
WHERE id = #{xx} AND status = 3;
5.4 线程池拒绝回退任务
比如突然一波流量,线程池接收不过来,拒绝策略命中后,要把处理中再改回待执行。待定时任务去执行。
如果不改回待执行,那就需要等待3分钟后再捞回来了。
UPDATE t_retry_task
SET status = 0, ip = NULL, update_time = NOW()
WHERE id = #{taskId} AND status = 3;
5.5 MIS人工重试(仅操作失败终态)
重置 status、next_retry_time、retry_count,待定时任务再次轮询到此任务。
UPDATE t_retry_task
SET status = 0,
ip = NULL,
next_retry_time = NOW(),
retry_count = 0,
update_time = NOW()
WHERE id = #{taskId} AND (status = 2 or (status = 3 and max_retry_count <= retry_count));
六、关键设计
6.1 抢跑约束
-
仅新建任务抢跑:失败回退、人工恢复的任务,一律不走抢跑,只走定时调度,保证重试策略生效
-
抢占互斥:单行抢占+状态校验,避免多实例重复执行
6.2 定时任务优化
-
固定使用
fixedDelay定时,避免调度线程堆积 -
更新SQL携带
ORDER BY next_retry_time ASC, id ASC,按顺序来,避免多实例执行遇到MySQL死锁 -
IP字段设计,多实例并发,但避免任务争抢、重复执行
6.3 线程池
-
按event_type独立线程池,业务隔离,互不影响
-
队列长度适当,拒绝策略使用AbortPolicy,拒绝后主动告警、不阻塞、留给定时兜底
-
禁止使用CallerRunsPolicy,防止阻塞全局定时线程
6.4 数据归档
-
定期归档到Hive
-
本地表保留7天数据,批量清理。避免表过大,降低索引效率
七、监控告警
-
僵死回收任务数(持续上涨告警)
-
线程池拒绝次数(流量过载告警)
-
失败终态任务堆积数(业务异常告警)
八、方案总结
-
低延迟:新建任务毫秒级抢跑执行,不只依赖定时延迟
-
低锁冲突:单行抢占为主、定时小批量更新。where更新
-
高可靠:抢跑+定时双兜底,100%不丢任务
-
高扩展:接口化业务抽象,新增业务无需改调度核心代码
-
易运维:支持后台人工重试、数据归档追溯