一、先搞懂为啥会有这个麻烦事儿
做分布式系统的朋友应该都遇过这么个糟心场景:比如你要给10个商家发促销补贴,每个商家的补贴得从账户池扣钱、给商家余额加钱、再发通知,这一套操作得同时成或者同时败——这就是分布式事务要干的事儿。之前很多人会用XXL-JOB来拆成10个小任务,每个任务对应一个商家的操作,想着一个任务错了就全停,结果发现根本不行:比如第5个任务失败了,前4个已经执行完的任务没法回滚,钱都扣了,商家也加了余额,通知也发了,想撤回来比登天还难。
为啥XXL-JOB拆的子任务没法直接回滚?核心原因是XXL-JOB的任务调度是“异步+独立”的:每个子任务是单独跑的,跑完就跟其他任务没关系了,既没有全局的事务上下文,也没有“任务执行的状态链”——前一个任务跑成什么样,后一个任务不知道,更别说把之前的操作撤销了。这时候就得搞“补偿任务”:就是专门给失败的子任务“擦屁股”的任务,比如第5个任务失败了,补偿任务就把前4个已经扣的钱加回去、商家余额减回去、通知撤回。但补偿任务最容易踩的坑就是“重复执行”:比如补偿任务跑了一次,又因为网络卡或者调度错,又跑了一次,结果本来要加回去的钱又加了一次,商家反而赚了,乱套了。所以补偿任务必须保证“幂等”——不管跑多少次,效果跟跑一次一样。
二、补偿任务的幂等到底要做啥
先给幂等说个大白话:就像你给微信发红包,点一次发出去,再点一次就显示“你已经发过这个红包了”,不会重复扣钱。补偿任务的幂等,核心就是要记住“哪些操作已经做过了”,下次再遇到同样的操作,直接说“我知道了,不用再做了”。
要实现这个,得抓住两个关键:一是得有个“唯一的操作标识”,能把每一次需要补偿的操作跟其他操作区分开;二是得有个“状态记录库”,专门存这些标识的状态,比如是“待补偿”“已补偿”还是“补偿失败”。
2.1 先搞懂唯一操作标识怎么来
这个标识得满足三个要求:唯一、固定、能对应到具体的操作。比如你之前给10个商家发补贴,每个商家的操作可以用“补贴批次ID+商家ID”来做标识,比如批次ID是20240520001,商家ID是123,那标识就是“20240520001_123”,这个组合肯定不会跟其他操作重复。
为啥不能用时间戳?因为两个操作可能在同一毫秒发生,就重复了;为啥不能用商家ID单独?因为同一商家可能多次参与补贴,比如这次发了,下次又发,用商家ID就分不清哪次的操作了。
2.2 状态记录库的核心要求
这个库不用复杂,就是一张表,专门存“唯一标识”和对应的“操作状态”,表的结构也很简单,举个例子(用MySQL,因为大家用得最多):
CREATE TABLE compensation_task_record (
id BIGINT AUTO_INCREMENT PRIMARY KEY COMMENT '主键',
task_key VARCHAR(128) NOT NULL COMMENT '唯一操作标识,比如补贴批次+商家ID',
task_status TINYINT NOT NULL COMMENT '状态:0=待补偿,1=已补偿,2=补偿失败',
create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '记录创建时间',
update_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '记录更新时间',
UNIQUE KEY uk_task_key (task_key) COMMENT '唯一索引,防止同一个标识重复插入'
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='补偿任务状态记录表';
这里的唯一索引uk_task_key是关键,它能保证同一个task_key只能有一条记录,不会出现两个状态的情况。
三、完整的实现步骤(带详细代码)
先明确用的技术栈:Java 17 + Spring Boot 3.2 + XXL-JOB 2.4.1 + MySQL 8.0,所有代码都用这个栈,不混其他的。
3.1 第一步:主任务拆分子任务时,提前插入状态记录
主任务就是负责拆分10个商家的子任务的任务,在拆分之前,先给每个商家的操作插入一条“待补偿”的状态记录,这样后续补偿的时候有依据。
比如主任务的代码(XXL-JOB的任务方法,返回字符串,失败返回null):
import com.xxl.job.core.context.XxlJobHelper;
import com.xxl.job.core.handler.annotation.XxlJob;
import org.springframework.stereotype.Component;
import java.util.List;
@Component
public class SubsidyMainTask {
// 假设这里是查询需要发补贴的商家列表,实际项目中可以从数据库查
private final List<Long> shopIds = List.of(123L, 456L, 789L, 1011L, 1213L, 1415L, 1617L, 1819L, 2021L, 2223L);
// 补贴批次ID,实际项目中可以从请求参数或者配置取
private final String BATCH_ID = "20240520001";
@XxlJob("subsidyMainTask")
public String execute() {
try {
// 第一步:给每个商家插入待补偿的状态记录
for (Long shopId : shopIds) {
String taskKey = BATCH_ID + "_" + shopId;
// 调用服务插入状态记录,状态为0(待补偿)
compensationTaskService.insertTaskRecord(taskKey, 0);
// 第二步:拆分子任务,每个商家一个子任务,把taskKey传进去
XxlJobHelper.shardParam(shopId.toString() + "|" + taskKey);
XxlJobHelper.shard(10); // 拆成10个分片,每个分片一个子任务
}
return "主任务拆分成功";
} catch (Exception e) {
XxlJobHelper.log("主任务拆分失败:" + e.getMessage());
return null;
}
}
}
这里的insertTaskRecord方法,因为状态表有唯一索引,所以如果同一个taskKey重复插入,会报唯一键冲突,这时候我们可以忽略,因为已经有记录了,不用再插。
3.2 第二步:子任务执行失败时,标记状态为待补偿
子任务就是每个商家的具体操作,比如扣钱、加余额、发通知,只要有一步失败,就把对应的taskKey的状态改成“待补偿”,这样补偿任务就能找到它。
子任务的代码:
import com.xxl.job.core.context.XxlJobHelper;
import com.xxl.job.core.handler.annotation.XxlJob;
import org.springframework.stereotype.Component;
@Component
public class SubsidySubTask {
@XxlJob("subsidySubTask")
public String execute() {
// 从分片参数中获取商家ID和taskKey
String shardParam = XxlJobHelper.getShardParam();
String[] params = shardParam.split("\\|");
Long shopId = Long.parseLong(params[0]);
String taskKey = params[1];
try {
// 第一步:扣账户池的钱(实际项目中是调用账户服务)
accountService.deduct(shopId, 100.0);
// 第二步:给商家加余额(实际项目中是调用商家服务)
shopService.addBalance(shopId, 100.0);
// 第三步:发通知(实际项目中是调用通知服务)
notifyService.send(shopId, "补贴已到账");
// 所有操作成功,标记状态为已补偿(或者标记为已完成,根据需求)
compensationTaskService.updateTaskStatus(taskKey, 1);
return "子任务执行成功";
} catch (Exception e) {
// 任何一步失败,标记状态为待补偿,让补偿任务处理
compensationTaskService.updateTaskStatus(taskKey, 0);
XxlJobHelper.log("子任务执行失败,商家ID:" + shopId + ",错误:" + e.getMessage());
return null;
}
}
}
这里要注意,子任务的操作要尽量保证“原子性”,比如扣钱、加余额、发通知,要么全成,要么全败,不然状态标记就不准了。
3.3 第三步:补偿任务的幂等实现
补偿任务就是专门处理状态为“待补偿”的taskKey的,核心逻辑是:先查这个taskKey的状态,如果是待补偿,就执行补偿操作,然后把状态改成已补偿;如果不是待补偿,就直接跳过,啥也不做。
补偿任务的代码:
import com.xxl.job.core.context.XxlJobHelper;
import com.xxl.job.core.handler.annotation.XxlJob;
import org.springframework.stereotype.Component;
import java.util.List;
@Component
public class CompensationTask {
@XxlJob("compensationTask")
public String execute() {
try {
// 第一步:查询所有状态为待补偿的taskKey
List<String> taskKeys = compensationTaskService.queryPendingTaskKeys();
if (taskKeys.isEmpty()) {
return "没有需要补偿的任务";
}
for (String taskKey : taskKeys) {
try {
// 第二步:执行补偿操作,先查taskKey对应的状态
int currentStatus = compensationTaskService.queryTaskStatus(taskKey);
// 只有状态是待补偿,才执行补偿
if (currentStatus == 0) {
// 拆分taskKey,得到批次ID和商家ID(之前的taskKey是批次+商家)
String[] parts = taskKey.split("_");
String batchId = parts[0];
Long shopId = Long.parseLong(parts[1]);
// 第三步:执行补偿操作,跟子任务的操作反过来
// 1. 把扣的钱加回去
accountService.add(shopId, 100.0);
// 2. 把商家加的余额减回去
shopService.subBalance(shopId, 100.0);
// 3. 撤回通知(如果支持的话)
notifyService.revoke(shopId, batchId);
// 第四步:把状态改成已补偿
compensationTaskService.updateTaskStatus(taskKey, 1);
XxlJobHelper.log("补偿任务执行成功,taskKey:" + taskKey);
}
} catch (Exception e) {
// 补偿失败,标记为补偿失败,下次再处理
compensationTaskService.updateTaskStatus(taskKey, 2);
XxlJobHelper.log("补偿任务执行失败,taskKey:" + taskKey + ",错误:" + e.getMessage());
}
}
return "补偿任务执行完成";
} catch (Exception e) {
XxlJobHelper.log("补偿任务整体执行失败:" + e.getMessage());
return null;
}
}
}
这里的核心是“先查状态再执行”,如果状态不是待补偿,直接跳过,不管补偿任务跑多少次,都不会重复执行补偿操作,这就是幂等的核心。
四、这种方案的优缺点和注意事项
4.1 优点
首先,实现简单,不用改XXL-JOB的核心逻辑,只是加了一个状态表和补偿任务,成本很低;其次,灵活性高,补偿操作可以根据具体场景定制,比如有的操作不能撤回,就可以调整补偿逻辑;最后,可靠性强,状态记录在数据库里,不会因为系统重启或者网络问题丢失,补偿任务可以重复跑,直到所有待补偿的任务都处理完。
4.2 缺点
首先,需要额外维护一个状态表,增加了数据库的负担;其次,补偿任务是异步的,不能保证实时性,比如子任务失败了,补偿任务可能要等调度周期到了才会执行;最后,补偿操作本身也可能失败,比如账户服务挂了,补偿任务没法把钱加回去,这时候需要人工干预。
4.3 注意事项
第一,状态表的唯一索引必须加,不然会出现同一个taskKey多条记录,导致状态混乱;第二,补偿操作要跟子任务的操作完全相反,比如子任务扣了钱,补偿就要加钱,子任务加了余额,补偿就要减余额,不能搞反;第三,补偿任务的调度周期要合理,不能太长也不能太短,太长会导致数据不一致的时间太长,太短会增加系统负担;第四,状态的更新要跟补偿操作绑定,比如补偿操作成功了,再更新状态,不能先更新状态再执行补偿操作,不然补偿操作失败了,状态已经改成已补偿,就不会再处理了;第五,要处理补偿失败的情况,比如补偿任务跑了几次都失败,要发告警,让开发人员人工处理。
五、应用场景总结
这种方案最适合的场景是:分布式事务的场景比较复杂,比如跨多个服务、跨多个数据库,用传统的分布式事务(比如TCC、XA)实现起来难度大、成本高,而且对实时性要求不是特别高的场景。比如电商的促销补贴、积分发放、批量订单处理等,这些场景即使有几分钟的不一致,也不会影响核心业务,而且补偿任务跑几次就能恢复一致。
另外,这种方案也适合那些“最终一致性”要求的场景,就是不要求所有操作同时成或者同时败,只要最终能一致就行,比如用户的积分过期、商品的库存调整等。
六、方案总结
XXL-JOB拆分子任务后没法直接回滚的问题,核心是因为子任务是独立的,没有全局的事务上下文,所以需要用补偿任务来处理失败的子任务。而补偿任务的幂等,核心是通过“唯一操作标识”和“状态记录库”来实现的,先查状态再执行补偿操作,保证不管跑多少次,效果跟跑一次一样。
这种方案虽然有一些缺点,比如需要额外维护状态表、补偿是异步的,但实现简单、灵活性高、可靠性强,适合大部分的分布式事务场景,尤其是那些对实时性要求不高、最终一致性要求的场景。
最后要注意的是,这种方案只是解决了“子任务回滚困难”的问题,不是完美的分布式事务解决方案,在一些对实时性要求很高的场景,比如支付、转账,还是需要用传统的分布式事务方案,比如TCC、XA,或者用消息队列来实现最终一致性。
评论
围绕“XXL-JOB结合分布式事务时子任务回滚困难,设计补偿任务的幂等性要点”参与讨论