事务消息 + 本地消息表 + 生产踩坑

事务消息 + 本地消息表 📖 前置阅读:本文是分布式事务系列的第四篇——假设你已经理解了 CAP/BASE 理论、Seata AT 的 undo_log 机制、TCC 的 Try/Confirm/Cancel 三阶段和 Saga 的补偿链。如果这些概念还陌生——先读 分布式事务本质——CAP、BASE 与四大方案、Seata AT 模式——undo_log 与二阶段原理 和 TCC + Saga——补偿型分布式事务。 一、⚡ 同步方案的瓶颈——为什么还需要异步方案 先回顾前面三篇文章我们做了什么: Seata AT:下单 → 扣库存 → 扣余额——三个操作在一个 @GlobalTransactional 中——同步执行 TCC:Try 预留 → Confirm 确认 → Cancel 回滚——三个阶段——同步执行 Saga:正向执行 → 失败逆补偿——协调者串联——同步执行 它们有一个共同特征:调用方要等所有分支都执行完——才返回结果。 order-service 调用 product-service 扣库存: → 发起 RPC 调用 → 等待 product-service 处理 → 等待 product-service 返回结果 → 拿到结果——继续下一步 如果 product-service 很慢——比如库存要查 3 个 Redis + 2 个 DB: → order-service 的线程就等着 → 线程池撑爆 → 整个链路超时 同步方案的根本矛盾:事务参与方的响应时间——直接影响调用方的吞吐量。 ...

十二月 30, 2022 · 17 分钟 · 3414 字 · yaomingye

TCC + Saga——补偿型分布式事务

TCC + Saga 📖 前置阅读:本文假设读者已理解 Seata AT 模式的原理和局限。如果还不熟悉,建议先阅读 Seata AT 模式——undo_log 与二阶段原理。 一、⚡ AT 能回滚库存——但能回滚一条"已发出的短信"吗? AT 模式的局限——上一篇说了: AT 的自动回滚依赖 undo_log——生成反向 SQL INSERT → DELETE(undo_log 记录自增 ID——反向就是 DELETE) UPDATE → UPDATE(undo_log 记录前置镜像——反向就是把值改回去) 但以下操作——数据库回滚不了: ① 发了优惠券——HTTP POST 到营销系统的 API——数据库回滚不了 HTTP 调用 ② 发了短信——调了阿里云短信 API——阿里云不会因为你的"反向 SQL"就收回短信 ③ 调了第三方支付——Payment API 已经扣了钱——不能"生成反向 HTTP"退钱 ④ 给 Redis 写了一个计数器——Redis 没有 undo_log——AT 管不了 TCC 和 Saga 就是为这而生的——手动补偿——操作本身和撤回操作都由你写代码实现。 二、🔄 TCC——Try / Confirm / Cancel——你自己管理回滚 2.1 TCC 的本质——每个操作配一个"撤销操作" TCC 把每个业务操作拆成三个方法: Try(尝试) —— 预留资源——但不真正执行 Confirm(确认) —— 真正执行——Try 预留的资源生效 Cancel(取消) —— 释放 Try 预留的资源——回滚 和 AT 的区别: AT:你写一套代码——Seata 自动生成"撤销操作"(反向 SQL) TCC:你写三套代码——Try(正向)、Confirm(确认)、Cancel(撤销) → 写了三套代码——能处理任何类型的操作——不再局限于数据库 2.2 示例——“创建订单 + 发优惠券 + 扣积分”——用 TCC // ===== 场景:下单时——创建订单 + 发优惠券 + 扣积分 ===== // 订单是 DB 操作——但发优惠券是 HTTP API——扣积分也是 HTTP API // AT 回滚不了 HTTP API——用 TCC // ===== 订单服务——TCC 接口 ===== public interface OrderTccAction { /** * Try:预创建订单——状态为 PENDING——库存还没扣——订单还不能支付 * @param businessContext 在 TM 端传入的参数——和 @BusinessActionContextParameter 对应 */ @TwoPhaseBusinessAction( name = "order-create", // TCC 资源名 commitMethod = "confirmCreateOrder", // Confirm 方法 rollbackMethod = "cancelCreateOrder" // Cancel 方法 ) boolean tryCreateOrder( @BusinessActionContextParameter(paramName = "userId") Long userId, @BusinessActionContextParameter(paramName = "items") List<OrderItemDto> items, @BusinessActionContextParameter(paramName = "totalAmount") BigDecimal totalAmount ); /** * Confirm:把订单从 PENDING 变为 CREATED——正式生效 */ boolean confirmCreateOrder(BusinessActionContext context); /** * Cancel:把 PENDING 的订单变为 CANCELLED——释放预占 */ boolean cancelCreateOrder(BusinessActionContext context); } // ===== 订单服务——TCC 实现 ===== @Service public class OrderTccActionImpl implements OrderTccAction { @Autowired private OrderMapper orderMapper; @Override @Transactional public boolean tryCreateOrder(Long userId, List<OrderItemDto> items, BigDecimal totalAmount) { // ① 预创建订单——状态为 PENDING——不是正式订单 Order order = new Order(); order.setOrderNo(generateOrderNo()); order.setUserId(userId); order.setTotalAmount(totalAmount); order.setStatus(OrderStatus.PENDING); // ← PENDING——不是正式订单——不可支付 order.setCreatedAt(LocalDateTime.now()); orderMapper.insert(order); // ② 把 orderId 存入 BusinessActionContext——Confirm/Cancel 会用到 // Seata 自动把方法返回值之外的参数存入 Context // 这里通过 RootContext 手动放 RootContext.bind("orderId_" + RootContext.getXID(), order.getId()); return true; // Try 成功——等待 TC 通知 Confirm 或 Cancel } @Override @Transactional public boolean confirmCreateOrder(BusinessActionContext context) { // ① 从 Context 中取出 orderId Long orderId = (Long) context.getActionContext() .get("orderId_" + context.getXid()); // ② 把订单状态从 PENDING → CREATED——正式生效 Order order = orderMapper.selectById(orderId); if (order == null || order.getStatus() != OrderStatus.PENDING) { // 幂等——如果已经 Confirm 过了——不再处理 return true; } order.setStatus(OrderStatus.CREATED); orderMapper.updateById(order); return true; } @Override @Transactional public boolean cancelCreateOrder(BusinessActionContext context) { Long orderId = (Long) context.getActionContext() .get("orderId_" + context.getXid()); Order order = orderMapper.selectById(orderId); if (order == null) { // 空回滚——Try 还没执行——Cancel 先到了——不做处理 return true; } if (order.getStatus() == OrderStatus.CANCELLED) { // 幂等——已经取消过了 return true; } order.setStatus(OrderStatus.CANCELLED); orderMapper.updateById(order); return true; } } // ===== 优惠券服务——TCC 接口(HTTP API——AT 回滚不了)===== public interface CouponTccAction { @TwoPhaseBusinessAction( name = "coupon-grant", commitMethod = "confirmGrantCoupon", rollbackMethod = "cancelGrantCoupon" ) boolean tryGrantCoupon( @BusinessActionContextParameter(paramName = "userId") Long userId, @BusinessActionContextParameter(paramName = "couponType") String couponType ); boolean confirmGrantCoupon(BusinessActionContext context); boolean cancelGrantCoupon(BusinessActionContext context); } @Service public class CouponTccActionImpl implements CouponTccAction { @Autowired private CouponService couponService; // 这个 Service 调外部营销 API @Override public boolean tryGrantCoupon(Long userId, String couponType) { // Try:预占优惠券——调营销 API——标记为用户——但未激活 Coupon coupon = couponService.reserveCoupon(userId, couponType); // 外部 API 返回了 couponId RootContext.bind("couponId_" + RootContext.getXID(), coupon.getId()); return true; } @Override public boolean confirmGrantCoupon(BusinessActionContext context) { // Confirm:激活优惠券——用户可用 Long couponId = (Long) context.getActionContext() .get("couponId_" + context.getXid()); couponService.activateCoupon(couponId); // HTTP PUT /coupons/{id}/activate return true; } @Override public boolean cancelGrantCoupon(BusinessActionContext context) { // Cancel:回收优惠券——把预留的优惠券放回库存 Long couponId = (Long) context.getActionContext() .get("couponId_" + context.getXid()); if (couponId == null) { return true; // 空回滚——Try 还没执行完 } couponService.recycleCoupon(couponId); // HTTP DELETE /coupons/{id} return true; } } // ===== TM——全局事务发起方——调各个 TCC 接口 ===== @Service public class OrderApplicationService { @Autowired private OrderTccAction orderTccAction; @Autowired private CouponTccAction couponTccAction; @Autowired private PointTccAction pointTccAction; @GlobalTransactional public Order createOrderWithCoupon(CreateOrderRequest request) { // ① Try:预创建订单 boolean orderTry = orderTccAction.tryCreateOrder( request.getUserId(), request.getItems(), request.getTotalAmount()); if (!orderTry) throw new BusinessException("预创建订单失败"); // ② Try:预发优惠券——不是数据库操作——是 HTTP 调外部 API boolean couponTry = couponTccAction.tryGrantCoupon( request.getUserId(), "FIRST_ORDER"); if (!couponTry) throw new BusinessException("预发优惠券失败"); // ③ Try:预扣积分——也是 HTTP 调外部 API boolean pointTry = pointTccAction.tryDeductPoints( request.getUserId(), 100); if (!pointTry) throw new BusinessException("预扣积分失败"); // ④ 所有 Try 成功——TM 通知 TC 进 Confirm // TC 依次调每个 RM 的 confirmXxx() // → orderTccAction.confirmCreateOrder() ——订单 PENDING→CREATED // → couponTccAction.confirmGrantCoupon() ——优惠券激活 // → pointTccAction.confirmDeductPoints() ——积分确认扣除 return ...; // 返回订单信息 } // 如果任何一个 Try 抛异常——TC 依次调每个 RM 的 cancelXxx() // → orderTccAction.cancelCreateOrder() ——订单 PENDING→CANCELLED // → couponTccAction.cancelGrantCoupon() ——优惠券回收 // → pointTccAction.cancelDeductPoints() ——积分退回 } 2.3 TCC 的两个致命陷阱——空回滚与悬挂 陷阱一:空回滚——Try 没执行——Cancel 先到了 时间线: ① TM 调 Order TCC 的 Try——网络超时——TM 不知道 Try 成功了没有 ② TM 决定回滚——发起 Cancel ③ Cancel 到达 order-service——但此时 Try 还没收到(网络延迟)——或者 Try 正在执行 ④ Cancel 执行时——订单不存在(Try 还没创建)——Cancel 失败 这叫"空回滚"——Cancel 先于 Try 到达 解决——控制记录表: 在 Cancel 中——如果查不到订单——不能报错——记录一条"Cancel 已执行"的空记录 当 Try 终于到达时——先查"Cancel 是否已执行"——如果是——Try 不再执行 陷阱二:悬挂——Try 超时后——Cancel 执行了——Try 又到了 时间线: ① TM 调 Try——Try 执行中——卡住了(GC 停顿——网络延迟) ② TM 等 10 秒超时——发起 Cancel ③ Cancel 到达——顺利执行——订单状态改为 CANCELLED ④ 第 30 秒——Try 终于执行完了——订单 INSERT 进去了——状态是 PENDING ⑤ 结果:Cancel 已经执行了——但 Try 把数据又写进去了——这个 Try"悬挂"了 这叫"悬挂"——Try 在 Cancel 之后到达——Cancel 的撤销效果被 Try 覆盖了 解决——同样用控制记录表: Cancel 执行时——记录一条"xid=xxx 已 Cancel" Try 执行前——先查"xid=xxx 是否已 Cancel"——如果是——拒绝执行 -- ===== TCC 防悬挂 + 空回滚控制表——每个参与 TCC 的服务都建一张 ===== CREATE TABLE tcc_operation_record ( id BIGINT AUTO_INCREMENT PRIMARY KEY, xid VARCHAR(128) NOT NULL COMMENT '全局事务 ID', branch_id BIGINT NOT NULL COMMENT '分支事务 ID', action_name VARCHAR(64) NOT NULL COMMENT 'TCC 资源名——order-create/coupon-grant', status TINYINT NOT NULL COMMENT '1-Try 2-Confirm 3-Cancel', created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, UNIQUE KEY uk_xid_branch_action (xid, branch_id, action_name) ) ENGINE=InnoDB DEFAULT CHARSET=utf8; // ===== 改进后的 TCC 实现——带防悬挂 + 空回滚 ===== @Service public class OrderTccActionImpl implements OrderTccAction { @Autowired private TccOperationRecordMapper recordMapper; @Override @Transactional public boolean tryCreateOrder(Long userId, List<OrderItemDto> items, BigDecimal totalAmount) { String xid = RootContext.getXID(); Long branchId = RootContext.getBranchId(); // ① 防悬挂——检查 Cancel 是否已执行 TccOperationRecord cancelRecord = recordMapper.selectOne( xid, branchId, "order-create", 3); // status=3 = Cancel if (cancelRecord != null) { // Cancel 先到了——Try 不能再执行——这就是"悬挂"——拒绝 return false; } // ② 记录 Try TccOperationRecord tryRecord = new TccOperationRecord(); tryRecord.setXid(xid); tryRecord.setBranchId(branchId); tryRecord.setActionName("order-create"); tryRecord.setStatus(1); // Try recordMapper.insert(tryRecord); // ③ 执行业务逻辑 Order order = new Order(); // ... 创建订单——状态 PENDING orderMapper.insert(order); RootContext.bind("orderId_" + xid, order.getId()); return true; } @Override @Transactional public boolean cancelCreateOrder(BusinessActionContext context) { String xid = context.getXid(); Long branchId = context.getBranchId(); // ① 幂等——检查 Cancel 是否已执行 TccOperationRecord existingRecord = recordMapper.selectOne( xid, branchId, "order-create", 3); if (existingRecord != null) { return true; // Cancel 已经执行过了——幂等——直接返回 } // ② 记录 Cancel——在查订单之前——防止空回滚 TccOperationRecord cancelRecord = new TccOperationRecord(); cancelRecord.setXid(xid); cancelRecord.setBranchId(branchId); cancelRecord.setActionName("order-create"); cancelRecord.setStatus(3); // Cancel recordMapper.insert(cancelRecord); // ③ 空回滚处理——查不到订单——不能报错 Long orderId = (Long) context.getActionContext().get("orderId_" + xid); if (orderId == null) { return true; // Try 没执行——空回滚——正常 } Order order = orderMapper.selectById(orderId); if (order == null) { return true; // Try 没执行完——空回滚——正常 } if (order.getStatus() == OrderStatus.CANCELLED) { return true; // 幂等 } // ④ 执行业务撤销 order.setStatus(OrderStatus.CANCELLED); orderMapper.updateById(order); return true; } } ⚠️ 新手提示:空回滚和悬挂是 TCC 的两个经典坑——90% 的 TCC 实现都有这两个问题。解决方案就是一张操作记录表——在 Cancel 执行前先记一笔"Cancel 已执行"——在 Try 执行前先查"Cancel 是否已执行"。记录表的唯一键 (xid, branch_id, action_name) 天然防并发——并发的 Try 和 Cancel 只有一个能插入成功。 ...

十二月 29, 2022 · 8 分钟 · 1671 字 · yaomingye

Seata AT 模式——undo_log 与二阶段原理

Seata AT 模式 📖 前置阅读:本文假设读者已理解分布式事务的核心问题(多数据库操作一致性)和 BASE 最终一致性概念。如果还不熟悉,建议先阅读 分布式事务本质——CAP、BASE 与四大方案。 一、⚡ Seata AT 一句话——你写你的 SQL——它自动生成反向 SQL 回想 XA 2PC 的问题——锁住数据库行等协调者——性能黑洞。Seata AT 是怎么解决的? XA 2PC 的做法(性能黑洞): ① Prepare:执行 SQL——不提交——锁住行 ② 等协调者——这期间这些行都是锁着的——其他事务不能动 ③ Commit/Rollback:提交或回滚——释放锁 Seata AT 的做法(攒反向 SQL——事后再补): ① 一阶段:执行 SQL——立即提交——释放锁——同时记录 undo_log(反向 SQL) ② 如果全局事务成功:删掉 undo_log——完事 ③ 如果全局事务失败:根据 undo_log 执行反向 SQL——把数据改回去 核心区别:XA 是锁住行等结果——Seata 是先把活干了——记下 undo_log——失败了逆向执行。 二、🏗️ Seata 架构——TC / TM / RM 三角 flowchart LR TM["TM(Transaction Manager)\n全局事务管理者\n-- 标注 @GlobalTransactional"] RM1["RM(Resource Manager)\norder-service\n-- 操作 order 数据库"] RM2["RM(Resource Manager)\nproduct-service\n-- 操作 product 数据库"] RM3["RM(Resource Manager)\naccount-service\n-- 操作 account 数据库"] TC["TC(Transaction Coordinator)\nSeata Server\n-- 协调全局事务——管理全局锁"] TM -->|"① 开启全局事务"| TC TM -->|"② 调用 order-service"| RM1 RM1 -->|"③ 一阶段:执行业务 SQL + 记录 undo_log + 向 TC 注册分支事务"| TC TM -->|"④ 调用 product-service"| RM2 RM2 -->|"⑤ 一阶段:执行业务 SQL + 记录 undo_log + 注册分支事务"| TC TM -->|"⑥ 调用 account-service"| RM3 RM3 -->|"⑦ 一阶段:执行业务 SQL + 记录 undo_log + 注册分支事务"| TC TM -->|"⑧ 全局事务成功 → 通知 TC 提交"| TC TC -->|"⑨ 二阶段:通知所有 RM 删除 undo_log"| RM1 TC -->|"⑨ 通知所有 RM 删除 undo_log"| RM2 TC -->|"⑨ 通知所有 RM 删除 undo_log"| RM3 classDef style_TM fill:#431407,stroke:#ea580c,stroke-width:2px,color:#fed7aa; classDef style_TC fill:#450a0a,stroke:#dc2626,stroke-width:2px,color:#fecaca; class TM style_TM; class TC style_TC;``` | 角色 | 全称 | 作用 | 在哪里 | |------|------|------|------| | TC | Transaction Coordinator | 协调全局事务——管理全局锁——决定提交还是回滚 | Seata Server——独立部署 | | TM | Transaction Manager | 定义全局事务边界——标 `@GlobalTransactional` 的方法 | 发起方服务(order-service) | | RM | Resource Manager | 管理分支事务——执行 undo_log 记录——向 TC 注册 | 每个参与方服务(product/account) | ## 三、🔍 undo_log 的核心原理——Seata AT 的灵魂 ### 3.1 undo_log 表结构 ```sql -- 每个参与分布式事务的数据库都需要一张 undo_log 表 -- Seata 提供了建表 SQL——直接执行即可 CREATE TABLE undo_log ( id BIGINT(20) NOT NULL AUTO_INCREMENT, branch_id BIGINT(20) NOT NULL COMMENT '分支事务 ID', xid VARCHAR(100) NOT NULL COMMENT '全局事务 ID', context VARCHAR(128) NOT NULL COMMENT '上下文', rollback_info LONGBLOB NOT NULL COMMENT '回滚信息——记录前置镜像和后置镜像', log_status INT(11) NOT NULL COMMENT '状态:0-正常 1-全局事务已完成', log_created DATETIME NOT NULL, log_modified DATETIME NOT NULL, PRIMARY KEY (id), UNIQUE KEY ux_undo_log (xid, branch_id) ) ENGINE=InnoDB DEFAULT CHARSET=utf8; 3.2 undo_log 的工作原理——前置镜像 + 后置镜像 以"扣库存"为例——product 服务执行:UPDATE product SET stock = stock - 5 WHERE id = 1 一阶段——执行 SQL + 记录 undo_log: ① Seata 拦截 SQL——先查一下当前数据: SELECT stock FROM product WHERE id = 1 → stock = 10 ② 执行你的业务 SQL: UPDATE product SET stock = stock - 5 WHERE id = 1 → stock = 5 (后置镜像) ③ 立即提交——不锁行——释放数据库锁 ④ 记录 undo_log: 前置镜像:stock = 10 (SQL 执行前的值) 后置镜像:stock = 5 (SQL 执行后的值) 反向 SQL:UPDATE product SET stock = 10 WHERE id = 1 ⑤ 向 TC 注册:我的分支事务完成了——xid=xxx——undo_log 已记录 二阶段——提交: 全局事务成功 → TC 通知所有 RM 提交 → 删掉 undo_log 记录 → 完事 二阶段——回滚: 全局事务失败 → TC 通知所有 RM 回滚 → 读 undo_log 中的反向 SQL → 执行: UPDATE product SET stock = 10 WHERE id = 1 然后把数据改回去了 → 删掉 undo_log 记录 关键——为什么 AT 比 XA 快: ...

十二月 28, 2022 · 8 分钟 · 1567 字 · yaomingye

分布式事务本质——CAP、BASE 与四大方案

分布式事务本质 一、⚡ @Transactional 在生产中失效——不是代码写错了——是底层就不是一回事 先看一个场景——最经典的"下单扣库存": // 单体应用——一个 @Transactional 搞定 @Service public class OrderService { @Transactional public void createOrder(CreateOrderRequest request) { // ① 创建订单 orderMapper.insert(order); // ② 扣库存 product.setStock(product.getStock() - quantity); productMapper.updateById(product); // ③ 扣余额 account.setBalance(account.getBalance().subtract(totalAmount)); accountMapper.updateById(account); // 这三个操作在同一个数据库中——同一个事务——要么全成功——要么全回滚 } } 拆成微服务后——同样的流程——@Transactional 失效: order-service ──→ 创建订单(自己的数据库) product-service ──→ 扣库存(product 数据库) account-service ──→ 扣余额(account 数据库) 每个服务有独立的数据库——三个 @Transactional 是三个独立的事务 → 订单创建成功——库存扣减成功——但扣余额失败 → 订单已创建——库存已扣——余额没变——钱还在——但东西已经扣了 → 数据不一致——用户赚了——公司亏了 分布式事务的本质问题:多个数据库(或服务)的操作——怎么保证"要么全成功、要么全回滚"? ...

十二月 27, 2022 · 4 分钟 · 706 字 · yaomingye
Cat Radio