资讯详情 分布式事务与最终一致性:双 11 扣减库存与赠送积分的本地消息表落地
📅 2026/10/11 6:21:11
大促架构设计里有一条用真金白银换来的血泪教训凡是在每秒数万 TPS 的主交易链路上搞跨服务强一致分布式事务如 2PC/XA、全局事务锁的双 11 零点基本没有不翻车的。三年前我们曾经迷信过某开源分布式事务框架。大促零点洪峰一到用户下单扣减库存框架为了维护全局一致性在跨服务的多个数据源之间协调 Prepare 与 Commit。结果跨机房网络延迟一抖动大量的数据库行锁被跨网络的长事务死死按住连接池几十秒内被占满整个交易链路连环雪崩最后只能靠紧急切断熔断开关收场。痛定思痛我们把架构原则重塑为核心链路与非核心链路坚决物理剥离能走异步最终一致性的绝不搞同步强一致。以大促最经典的“扣减实物库存”与“赠送促销积分”为例扣库存是生死线必须强事务原子性而赠送用户 10 倍积分属于履约激励允许有几秒到几分钟的延迟。放弃沉重的 2PC选择最务实、最经得起风浪的本地消息表Local Message Table / Transactional Outbox Pattern才是小厂扛住高并发的定海神针。核心解题思路用单机 ACID 换取分布式可靠性为什么直接在扣完库存后调用消息队列MQ或积分 RPC 是错的如果先扣库存、再发 MQ一旦数据库提交成功而 MQ 网络闪断积分消息彻底丢失买家必然投诉维权如果先发 MQ、再扣库存万一库存不足事务回滚MQ 已经发出去了白白送给用户积分造成资产流失如果把 RPC 嵌套在数据库事务里外部网络抖动会直接延长数据库事务持锁时间瞬间拖垮整个单库并发。本地消息表的核心逻辑是将“跨网络的不可靠通信”降维为“同一个单机数据库内部的原子操作”。┌────────────────────────────────────────────────────────┐ │ 订单/库存数据库 (单机事务) │ │ │ │ BEGIN TRANSACTION; │ │ 1. UPDATE inventory SET stock stock - 1 ... │ │ 2. INSERT INTO local_outbox_msg (statusPENDING)..│ │ COMMIT; │ └───────────────────────┬────────────────────────────────┘ │ (单机 ACID 100% 绑定) ▼ ┌────────────────────────────────────────────────────────┐ │ 异步投递 Worker (轮询或 CDC 读取 PENDING 消息) │ │ 1. 投递到消息队列 (Kafka / RocketMQ) │ │ 2. 收到投递成功 ACK更新消息表状态为 SUCCESS │ └───────────────────────┬────────────────────────────────┘ │ ▼ ┌────────────────────────────────────────────────────────┐ │ 积分服务 (Consumer 端) │ │ 1. 消费消息并提取业务订单号 │ │ 2. 插入积分明细流水 (利用唯一索引实现严格幂等防重) │ │ 3. 增加用户可用积分余额 │ └────────────────────────────────────────────────────────┘Go 1.27.1 生产落地实战我们使用 Go 1.27.1 构建这套本地消息事务投递器。通过方法级泛型抽象各种业务领域的事件载荷搭配 Go 1.26 的new(expr)将内存分配负担降至最低。1. 数据库单机事务闭环package outbox import ( context database/sql encoding/json fmt time ) type OutboxMessage struct { ID int64 BizID string json:biz_id Topic string json:topic Payload string json:payload Status string json:status // PENDING, SUCCESS, FAILED RetryCount int json:retry_count CreatedAt time.Time json:created_at } type OrderInventoryService struct { db *sql.DB } func NewOrderInventoryService(db *sql.DB) *OrderInventoryService { return OrderInventoryService{db: db} } // DeductStockAndStagePoints 利用单机事务强绑定扣库存与写本地消息表 func (s *OrderInventoryService) DeductStockAndStagePoints[T any](ctx context.Context, orderID string, skuID int64, count int, pointsEvent T) error { payloadBytes, err : json.Marshal(pointsEvent) if err ! nil { return fmt.Errorf(marshal points event failed: %w, err) } tx, err : s.db.BeginTx(ctx, sql.TxOptions{Isolation: sql.LevelReadCommitted}) if err ! nil { return err } defer tx.Rollback() // 1. 核心扣库存行级条件防超卖 invSQL : UPDATE goods_inventory SET stock stock - ?, locked locked ? WHERE sku_id ? AND stock ? res, err : tx.ExecContext(ctx, invSQL, count, count, skuID, count) if err ! nil { return fmt.Errorf(exec inventory failed: %w, err) } rows, _ : res.RowsAffected() if rows 0 { return fmt.Errorf(sku %d inventory insufficient, skuID) } // 2. 将异步履约事件落入本地消息表 outboxSQL : INSERT INTO local_outbox_msg (biz_id, topic, payload, status, retry_count, created_at, updated_at) VALUES (?, ?, ?, PENDING, 0, NOW(), NOW()) _, err tx.ExecContext(ctx, outboxSQL, orderID, topic_points_grant, string(payloadBytes)) if err ! nil { return fmt.Errorf(insert outbox msg failed: %w, err) } // 3. 提交本地事务库存与消息同生共死 return tx.Commit() }2. 异步投递与消费端的幂等防线投递 Worker 定时或基于变更抓取PENDING状态记录并投递给 MQ投递成功后更新为SUCCESS。而在下游积分消费端最重要的守则只有四个字绝对幂等。大促期间网络重试极其频繁MQ 绝对不能假设只会投递一次At-Least-Once 机制。消费端必须使用数据库级别的唯一约束来阻挡重复消息CREATE TABLE user_points_flow ( id bigint unsigned NOT NULL AUTO_INCREMENT, user_id bigint NOT NULL, order_id varchar(64) NOT NULL, biz_type varchar(32) NOT NULL, points_delta int NOT NULL, created_at datetime NOT NULL, PRIMARY KEY (id), UNIQUE KEY uk_order_biz (order_id, biz_type) -- 防重核心 ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;消费端 Go 代码处理模式package consumer import ( context database/sql encoding/json errors fmt github.com/go-sql-driver/mysql ) type PointsGrantEvent struct { OrderID string json:order_id UserID int64 json:user_id GrantPoints int json:grant_points } type PointsConsumer struct { db *sql.DB } func (c *PointsConsumer) HandlePointsMessage(ctx context.Context, msgData []byte) error { var event PointsGrantEvent if err : json.Unmarshal(msgData, event); err ! nil { return err // 格式坏死直接进死信 } tx, err : c.db.BeginTx(ctx, nil) if err ! nil { return err } defer tx.Rollback() // 1. 尝试插入唯一流水行 flowSQL : INSERT INTO user_points_flow (user_id, order_id, biz_type, points_delta, created_at) VALUES (?, ?, ORDER_REWARD, ?, NOW()) _, err tx.ExecContext(ctx, flowSQL, event.UserID, event.OrderID, event.GrantPoints) if err ! nil { var mysqlErr *mysql.MySQLError if errors.As(err, mysqlErr) mysqlErr.Number 1062 { // 触发唯一键冲突说明该订单已经赠送过积分直接视为消费成功并 ACK tx.Rollback() return nil } return fmt.Errorf(insert points flow error: %w, err) } // 2. 增加用户账户总余额 balanceSQL : UPDATE user_account SET points points ? WHERE user_id ? if _, err : tx.ExecContext(ctx, balanceSQL, event.GrantPoints, event.UserID); err ! nil { return fmt.Errorf(update user balance error: %w, err) } return tx.Commit() }生产避坑两把斧消息表的膨胀治理在双 11 期间每天几百万的订单会让local_outbox_msg快速膨胀。如果扫描线程一直在大表上做全表索引检索数据库 I/O 很快会被拉爆。解决办法是对本地消息表按状态建立复合索引KEY idx_status_id (status, id)。投递成功的消息保留 7 天后由定时任务夜间分批物理删除DELETE FROM local_outbox_msg WHERE statusSUCCESS AND created_at NOW() - INTERVAL 7 DAY LIMIT 1000或者直接使用日级分区表 DROP 分区。投递线程的并发自旋控制多个投递 Worker 实例同时去拉取PENDING消息时必须使用带有锁限制的扫描策略例如SELECT id, payload FROM local_outbox_msg WHERE statusPENDING ORDER BY id ASC LIMIT 100 FOR UPDATE SKIP LOCKED。利用 MySQL 8.0 的SKIP LOCKED特性不同 Worker 实例互不等待、并行瓜分未投递任务既避免了死锁又让投递吞吐轻松跑满几万 TPS。放弃对“完美的瞬间一致性”的执念用本地事务锁定确定性用异步消息与幂等消化不确定性。大道至简这套组合拳才是支撑业务穿越风浪最皮实的底座。