Outbox + RocketMQ 接入技术方案
1. 结论与范围
worker-architecture-evolution-plan.html §3 中标注"中/低"风险的任务;标"资金高风险"的 7 个任务(5 个平台订单同步 + 佣金到期释放 + Xendit 提现重试)和标"高"风险的 TikTok Token 刷新,本方案不动,继续留在现有的轮询+锁机制上。每个任务的现状代码、接入 MQ 后的流程、逐项优缺点见 §15,本节只给结论。Push 消息派发、优惠券过期扫描、TikTok Partner Campaign 商品同步、App Store 版本同步、Ops 监控刷新。
不碰钱5 个平台订单同步 + 佣金到期释放 + Xendit 提现重试 + TikTok Token 刷新。
留在现状1 张 Outbox 表 + 1 个 Relay 定时任务(复用 worker 现有调度) + 1 个 RocketMQ 4.x 实例(标准版/专业版, 雅加达无 Serverless 可用区)。
需采购2. 决策记录:为什么直接上 RocketMQ
过程中评估过用现有生产 Redis 的 Streams 能力做零成本替代(详见本文档修订历史),但最终决定直接采购 RocketMQ,接受下面这笔成本,理由是:用专业消息队列从第一天就有正确的可靠投递、事务消息、消息轨迹等原生能力,不用把这些能力自己在 Redis 上补一遍,也不用背 Redis 身兼缓存/锁/事件流三个角色的容量和资源隔离风险。
| 核查项 | 结果 |
|---|---|
| ap-southeast-5(雅加达) 是否有 RocketMQ 5.0 Serverless 版可用区 | 产品名录标注支持(supportRocketmqV5: true),但 ListAvailableZones 实测为空列表,实际机房资源还没有铺开这一代产品 |
| 能买到的版本 | RocketMQ 4.x 标准版/专业版,固定预留容量,不管用多少都要付这份钱 |
| 实际报价(最小规格,年付) | ≈ ¥17,000/年 — 已确认接受,按此预算采购 |
3. 总体架构
flowchart LR
subgraph Producer["业务代码(不变)"]
BIZ["业务方法内: 写业务表 + 写 outbox_event, 同一个事务"]
end
subgraph Relay["Relay(worker 内新增一个 @Scheduled)"]
POLL["轮询 outbox_event WHERE status=PENDING"]
PUB["EventPublisher.publish(...)"]
MARK["成功: 标记 PUBLISHED / 失败: retry_count+1"]
end
subgraph MQ["RocketMQ(4.x 标准版/专业版)"]
TOPIC["Topic: beex-id-events, 按 Tag 区分事件类型"]
end
subgraph Consumer["消费者(新增,各自独立)"]
PUSH["Push 消费者"]
COUPON["优惠券消费者"]
CAMPAIGN["商品同步消费者"]
VERSION["版本同步消费者"]
OPS["Ops 快照消费者"]
end
BIZ --> POLL
POLL --> PUB
PUB --> TOPIC
PUB --> MARK
TOPIC --> PUSH
TOPIC --> COUPON
TOPIC --> CAMPAIGN
TOPIC --> VERSION
TOPIC --> OPS
业务代码只多做一件事——在原来的数据库事务里顺手插一行 outbox_event,不直接调用 MQ SDK。真正跟 RocketMQ 打交道的只有 Relay 这一个组件,出问题时排查面很窄。
4. 可替换的发布/订阅接口(以后 broker 还要再换,只换实现类)
Relay 是唯一"跟 broker 打交道"的组件,把它对 broker 的调用收敛成两个接口,Outbox 表和 5 个消费者的业务代码完全不知道底层具体是哪个 MQ 产品:
public interface EventPublisher {
void publish(String topic, String tag, String eventKey, String payloadJson);
}
public interface EventConsumer {
void subscribe(String topic, String tag, String consumerGroup, MessageHandler handler);
}
- 生产:
RocketMqEventPublisher/RocketMqEventConsumer,内部调用 RocketMQ SDK。 - 以后如果要换别的 broker:新写一个实现类,换掉 Spring 装配的 Bean,Outbox 表结构、Relay 轮询骨架、5 个消费者业务逻辑一行都不用改。
interface 声明,不多花额外工作量,但换来了以后不被今天的技术选型锁死——下一节就直接用上这个设计。5. 测试用 Redis Streams, 生产用 RocketMQ
beex-id-shared-redis(Redis 7.0, 支持 Streams, 已确认版本兼容)实现同一套 EventPublisher/EventConsumer 接口。两边跑的是同一套 Outbox 表结构、同一套 Relay 轮询骨架、同一套 5 个消费者业务代码——唯一的区别是 Spring 装配了哪个实现类,不是两套方案。| 环境 | broker 实现 | 成本 |
|---|---|---|
| 测试 | RedisStreamsEventPublisher / RedisStreamsEventConsumer,写到 beex-id-shared-redis 的 Stream | 零新增——复用已有测试 Redis |
| 生产 | RocketMqEventPublisher / RocketMqEventConsumer | ≈¥17,000/年,已确认接受 |
seahub.mq.provider=redis-streams # 测试环境 application-id-test.properties seahub.mq.provider=rocketmq # 生产环境 application.properties
Spring 用 @ConditionalOnProperty(name = "seahub.mq.provider", havingValue = "...") 按这个配置值选择装配哪一套实现,跟现有"环境变量驱动配置、无配置时安全降级"的风格一致。
OutboxEventModel/OutboxEventRepository/EventPublisher/EventConsumer/MessageHandler 接口、OutboxSchemaInitializer(outbox_event 表)、OutboxRelayService(Relay)、RedisStreamsEventPublisher/RedisStreamsEventConsumer 已落地并部署到 id-test。第一个接入的消费者是 Push 消息派发(不是本节原计划的 Ops 监控刷新——Push 有真正离散的"事件发生→通知下游"语义,Ops/优惠券这类纯周期扫描没有天然的触发点,改用 Push 作为 Phase 0 验证对象):PushService.enqueueUserNotification 在写 push_messages 的同一事务里落 outbox_event;PushEventConsumer 通过跟轮询路径共用的原子租约(leaseMessageIfDue)消费,避免两条路径重复派发。全部行为收在 seahub.mq.push-dispatch.enabled 开关后面(生产默认 false,等 RocketMQ 采购完成再实现 §5 表格里的 RocketMqEventPublisher/RocketMqEventConsumer 并切换;id-test 已打开)。id-test 实测:插入一条待发送 Push 消息 + 对应 outbox 事件后,
outbox_event.status 变为 PUBLISHED,Redis Stream 消费者组 GID_push 的 entries-read 从 0→1、lag=0,push_messages 状态从 QUEUED 变为终态——证明 Outbox 写入→Relay 发布→Redis Streams→消费者组→业务派发 全链路真实跑通(该消息最终因 id-test 未配置真实 APNs 密钥而发送失败,这跟 MQ 链路无关,是预期之外的环境限制)。顺带发现并修复一个更基础的问题:Spring Boot 调度线程池默认只有 1 个线程,TikTok 商品同步这类长任务(逐商品调用第三方 API,单次耗时可达十几分钟)会把同一线程上排队的所有其它
@Scheduled 任务(含 Push 派发、Outbox Relay)全部堵住。已把 spring.task.scheduling.pool.size 从默认 1 调到 10 并上线 id-test,这是排查本次验证时发现的、影响范围更广的既有问题,不只针对 Outbox。6. Outbox 表设计
BIGINT AUTO_INCREMENT。项目里所有业务表(commission_records、withdraw_requests、referral_relation 等)主键都是 VARCHAR(64),由 Ids.newId(prefix) 在应用层生成,不用数据库自增。id-design.html §8 明确写了:分片路由要基于业务 ID 的哈希或应用层生成的主键,不能依赖单点数据库自增计数器——AUTO_INCREMENT 只在单个数据库实例内部自增,未来一旦分库分表,多个分片各自独立自增会互相冲突,而且没法在写入前就知道最终 ID。改成跟其它表一致的 Ids.newId("obe") 方案。CREATE TABLE outbox_event (
id VARCHAR(64) PRIMARY KEY, -- Ids.newId("obe"),跟项目其它表主键生成方式一致
event_id VARCHAR(64) NOT NULL, -- 业务方生成的幂等键,如 push:{taskId}、coupon:{userId}:{couponId}:expire
event_type VARCHAR(64) NOT NULL, -- 如 PUSH_DISPATCH / COUPON_EXPIRE
aggregate_type VARCHAR(32) NOT NULL, -- 业务实体类型,如 push_task / user_coupon
aggregate_id VARCHAR(64) NOT NULL,
payload_json TEXT NOT NULL, -- 消费者需要的最小数据,不是整行业务数据
status VARCHAR(16) NOT NULL DEFAULT 'PENDING', -- PENDING / PUBLISHED / FAILED
retry_count INT NOT NULL DEFAULT 0,
created_at DATETIME NOT NULL,
published_at DATETIME NULL,
UNIQUE KEY uk_event_id (event_id),
KEY idx_status_created (status, created_at)
);
id用Ids.newId("obe")生成:时间戳(base36) + 12 位随机字符,任何一台机器、任何一个未来的分片都能独立生成,不需要跟别的分片协商,天然适合分库分表。event_id全局唯一,由业务方按"事件本身"生成(不是主键id),这是后面消费者去重的第一道防线,跟id是两个不同用途的字段:id是这一行数据本身的身份,event_id是这个事件在业务语义上的幂等键。- 写入
outbox_event必须和业务表写入在同一个@Transactional方法里,要么都成功要么都不写,这是整套方案成立的前提,不能后补。 payload_json只放消费者需要的字段,不要把整行业务实体序列化进去——避免表结构变了还要跟着改历史消息格式。
Ids.newId 是时间戳+12位随机字符(约 62 位随机熵),不是数学上绝对不会碰撞,但概率级别接近 UUID,项目里所有表都是这个方案、一直这么用。真出现极小概率碰撞时,PRIMARY KEY 本身会在插入时报重复键错误,Relay/业务方按这个错误重试一次(重新生成一个新 id)即可——这也是项目里其它表的标准处理方式,不是这张表要单独解决的新问题。7. Relay 投递
不新增一个独立部署的服务,直接在 worker 里加一个新的 @Scheduled 方法(第 14 个),复用现有的调度和监控框架:
| 步骤 | 做法 |
|---|---|
| 取待发布事件 | SELECT ... WHERE status='PENDING' ORDER BY created_at, id LIMIT 100 FOR UPDATE SKIP LOCKED——按 created_at 排序,不是按主键 id(id 是 Ids.newId 生成的字符串,不保证严格递增排序,不能拿来当时间顺序用);配合现有的 DistributedTaskLockExecutor 保证多 worker 实例不重复投递。 |
| 发布 | 调用 EventPublisher.publish(...)(内部是 RocketMQ SDK 发送),event_id 作为消息 Key(不是 MessageId),方便消费侧和控制台按业务键排查。 |
| 成功 | 标记 status=PUBLISHED,记录 published_at。 |
| 失败 | retry_count+1,下一轮继续重试;超过阈值(如 10 次)转 FAILED 并告警飞书,不能无限重试掩盖问题。 |
| 投递间隔 | 建议 5-10 秒一轮,跟现有 Push 派发的量级相当,不需要专门优化到秒级以下。 |
8. Topic / Tag 设计
第一阶段量级小,用一个 Topic + 多个 Tag,不按业务域拆多个 Topic——量小的时候多 Topic 只会增加控制台管理和权限配置的负担。
| Tag | 对应任务 | 消费者组 |
|---|---|---|
PUSH_DISPATCH | Push 消息派发 | GID_push |
COUPON_EXPIRE | 优惠券过期扫描 | GID_coupon |
CAMPAIGN_PRODUCT_SYNC | TikTok Partner Campaign 商品同步 | GID_campaign |
APP_VERSION_SYNC | App Store 版本同步 | GID_appversion |
OPS_SNAPSHOT | Ops 监控刷新 | GID_ops |
Topic 命名:beex-id-events(按国家环境区分,马来上线后是 beex-my-events,不同国家不共用 Topic,和现有"部署级国家隔离"原则一致)。
测试环境用 §5 的 Redis Streams 实现时,对应关系是 Topic→Stream Key(beex:id-test:events)、Tag→消息字段、消费者组名不变,上面这张表的 Tag 和消费者组两列在两个环境通用。
9. 消费者与幂等
幂等不是新发明,复用现有数据层已经有的唯一约束和状态机,消费者只是"从轮询触发"改成"从消息触发":
| 消费者 | 幂等依据 |
|---|---|
| Push 消费者 | 推送任务表已有的状态字段(已发送/待发送),消费者处理前先查状态,已处理过直接跳过。 |
| 优惠券消费者 | 用户券表的状态字段,重复消息只是重复判断一次"是否已过期",天然幂等,不需要额外去重表。 |
| 商品同步消费者 | 商品缓存表按商品 ID upsert,重复消息等于重复刷新一次缓存,无副作用。 |
| 版本同步/Ops 快照消费者 | 都是覆盖式写最新状态,重复消费无副作用。 |
10. 分阶段路线图
接 Push 消息派发 1 个任务,已在 id-test 验证 Outbox 表、Relay、Redis Streams、消费者组全链路跑通(见 §5 绿色标注),观察 1-2 周稳定后进入 Phase 1。
1 个任务·已实现Phase 0 稳定后,把剩下 4 个任务(优惠券、商品同步、版本同步、Ops 监控)全部接入,同时下线对应的旧轮询代码路径(但保留代码,用开关控制,见 §10)。
4 个任务只有 §12 的升级触发条件真正出现,才评估把订单同步、佣金结算、提现重试接入,且必须额外设计资金对账兜底机制。
暂不启动11. 开关与回滚
每个迁移的任务保留原有 @Scheduled 轮询代码,新增一个配置开关,不是删代码替换:
seahub.mq.push-dispatch.enabled=false # 默认关闭,灰度打开 seahub.mq.coupon-expire.enabled=false seahub.mq.campaign-sync.enabled=false seahub.mq.app-version-sync.enabled=false seahub.mq.ops-snapshot.enabled=false
开关关闭时,原有轮询方法正常执行,完全不受影响;开关打开时,原有轮询方法内部改为"只在消费者连续失败/超时时兜底触发一次",不是完全禁用轮询——相当于事件驱动是主路径,轮询降级为兜底安全网。
12. 监控告警
- Outbox 积压:
PENDING状态且创建时间超过 5 分钟的行数,超过阈值告警飞书(复用现有OpsMonitorService框架和告警通道)。 FAILED状态事件数:出现即告警,不能静默积压。- RocketMQ 消费延迟:每个消费者组的堆积消息数和最老消息年龄,阿里云控制台自带这个指标,先用控制台看,不用重复造轮子。
- 沿用现有告警格式要求:必须带环境、国家、任务名、trace/requestId,不能只写"任务失败"。
13. 采购与后续复核
- 采购规格:按 Phase 1 的量级选最小可用规格(标准版即可,不需要专业版的更高吞吐),年付 ≈ ¥17,000,走控制台采购(这条产品线的下单也可能遇到跟 RDS 只读库类似的售卖组件限制,如果 API/控制台下单失败,按同样思路排查或提工单)。
- 雅加达补齐 RocketMQ 5.0 Serverless 可用区后:值得重新核价,Serverless 按量计费在低量级场景通常比标准版固定预留容量便宜,到时候可以评估是否迁移到 Serverless(同一产品线内迁移,不涉及 §4 的 broker 抽象切换)。
- Phase 2 真的要启动时:佣金结算、提现这类资金路径接入事件驱动,对"消息绝对不丢、能完整审计回放"的要求更高,评估是否需要升级到专业版或调整容量。
14. 验收清单
- Outbox 写入和业务表写入在同一事务,手动模拟业务写入失败,确认不会产生孤立的 outbox_event。
- 双 worker 实例同时跑 Relay,同一条 outbox_event 不会被投递两次(靠
SKIP LOCKED+ 分布式锁双重验证)。 - 消费者重复收到同一条消息(手动重放),业务结果不变(幂等验证)。
- 关闭开关,系统行为完全回退到当前的纯轮询模式,功能不受影响。
- Outbox 积压、发布失败、消费延迟三类告警都验证过能正常触发飞书通知。
- Phase 0 单任务观察满 1-2 周、无异常后,才进入 Phase 1 全量迁移。
15. 各业务详细方案与优缺点
本节逐个业务给出:现状(轮询)代码位置和流程、接入 Outbox+MQ 后的流程、这个业务接 MQ 具体的优点和缺点/风险。所有代码位置和配置项名字都是 2026-07-19 逐个核对代码得出的,不是转述计划。
15.1 Push 消息派发(Phase 0,已实现并在 id-test 验证,见 §5)
现状:PushService.enqueueUserNotification(seahub-core/.../push/PushService.java)写 push_messages;PushDispatcherService.dispatchDueMessages每 5 秒轮询一次未处理消息并派发。
接入 MQ 后:写 push_messages 的同一事务里插一条 outbox_event(eventType=PUSH_DISPATCH);OutboxRelayService 5 秒轮询发布到 beex-id-events;PushEventConsumer 订阅后,通过跟轮询路径共用的原子租约(leaseMessageIfDue)取到消息就立刻派发,不用等下一次轮询。轮询路径保留作兜底,不删除。
| 优点 | 缺点/风险 |
|---|---|
| 从"最多等 5 秒"变成"消息一到就发",延迟从秒级降到毫秒级,对用户体感有实际改善。 | 多了一套基础设施(Outbox 表、Relay、消费者)要维护,故障面比纯轮询宽——但这套基础设施是所有 5 个 Phase 1 任务共用的,分摊下来单个任务的边际成本很低。 |
| 有真正离散的"事件发生→通知下游"语义,是这批任务里最适合做事件驱动的一个,选它做 Phase 0 验证对象是对的。 | 需要处理两条路径(轮询+消息)可能同时抢到同一条消息的竞态——已经用共享的原子租约解决,但这是接入 MQ 前不存在的复杂度。 |
15.2 优惠券过期扫描
现状:CouponExpirationJob(seahub-core/.../coupon/CouponExpirationJob.java)每 5 分钟(seahub.coupon.expire-scan-delay-ms,默认 300000)跑一条 SQL:update user_coupons set status='EXPIRED' where status in ('CLAIMABLE','ISSUED','ACTIVE') and expires_at <= now(),配合 DistributedTaskLockExecutor 防止多实例重复跑。没有任何下游副作用(不写 coupon_ledger,不发通知)。
接入 MQ 后:这里必须澄清一个容易理解错的地方——不是"过期变成事件触发",过期本身没有天然的触发时刻,只有"时间到了"这一个条件,MQ 解决不了"怎么知道时间到了"这个问题。真正的改法是:扫描任务保留(还是要定期跑,只是不再直接改 status),改成给每个到期的优惠券写一条 outbox_event(COUPON_EXPIRE),交给消费者异步做状态变更。
| 优点 | 缺点/风险 |
|---|---|
| 状态变更和(未来如果要加的)下游动作——比如过期通知、过期后的营销触达——解耦,以后要加这类下游动作不用改扫描任务本身。 | 核心收益很有限:现在这个任务是单条 SQL、5 分钟跑一次、毫秒级完成、零下游副作用,接 MQ 除了多一层间接性,不会让它更快或更可靠。 |
幂等零成本:status 字段本身就是天然幂等键,消费者重复收到同一条消息只是重复判断一次"是否已过期",不需要额外去重表。 | 如果只是为了"过期",继续用现在这条 UPDATE 语句最简单可靠;接 MQ 的价值要等真的有下游动作需要挂载时才体现出来,现在接等于"为将来可能用到的扩展性买单"。 |
15.3 TikTok Partner Campaign 商品同步
现状:TiktokPartnerCampaignProductSyncService(seahub-core/.../tiktok/)每 30 分钟跑一次,单线程串行:遍历最多 5 页 Campaign,每个 Campaign 下最多 100 页(5000 个)商品,每个商品都单独发一次 HTTP 请求去校验 Creator 分享链接,再逐条 upsert 到 tiktok_partner_campaign_products。这就是之前验证 Outbox 时实测发现的、把调度线程池堵住十几分钟的那个任务(已通过把线程池从 1 调到 10 缓解,但任务本身逐商品调用外部 API 的低效率没变)。
接入 MQ 后:拆成两层——保留一个轻量的"列 Campaign 列表"周期任务(便宜,只需分页 5 次),对每个 Campaign 写一条 outbox_event(CAMPAIGN_PRODUCT_SYNC,eventKey=国家:CampaignId);真正耗时的"翻商品页 + 逐个校验分享链接 + upsert"下放到消费者,多个 Campaign 可以并行处理,不再挤在一个调度线程上。
| 优点 | 缺点/风险 |
|---|---|
| 这是 Phase 1 里收益最实在的一个:直接解决"一个长任务堵死调度线程池"的根因,把单线程串行改成多 Campaign 并行,总耗时能大幅缩短。 | TikTok 这边本身没有 Webhook,拿不到"Campaign/商品变更"的推送通知,还是得靠周期性"列列表"去发现变化,MQ 只是把"发现变化之后要做的重活"并行化,不能消除轮询本身。 |
"哪些商品从 Campaign 里消失了要关闭"的全量对比逻辑(markMissingAsClosed)不用变,消费者按 Campaign 粒度做,天然幂等(按 country_code+campaign_id+product_id upsert)。 | 要把现在"一次跑完整个同步就知道结果"的模式,改成"多个 Campaign 异步并行完成",中间状态的可观测性要重新设计(比如怎么知道这一轮同步全部 Campaign 都跑完了)。 |
15.4 应用商店版本同步(注:跟 H5 包分发是两回事,见下方澄清)
AppRuntimeConfigService.applyAppleLoginGate)。H5 小程序包的分发(h5_package_versions 等表)是另一套完全独立、纯被动拉取的机制——客户端每次启动主动调 GET /api/v1/app/h5-package/latest 查询、下载更新,服务端没有任何轮询,不属于这次 MQ 迁移的候选范围。改造后现状:轮询任务已迁移到管理后台服务。管理后台每小时同步 iOS App Store 和 Android Google Play production track,将最后一次成功结果写入 runtime_feature_config 的 SYSTEM 配置;业务服务只读配置,不再保存实例内存版本或直接访问应用商店。
MQ 判断:当前消费方仍然很少,数据库配置已经解决多实例一致性和重启丢失问题,暂不需要为版本变化增加 Outbox/MQ 事件。
| 优点 | 缺点/风险 |
|---|---|
| 如果以后有多个异步消费方需要"应用商店版本变化"信号,可基于已持久化配置补充 Outbox 事件。 | 现在主要消费方是 AppRuntimeConfigService 的登录方式开关,直接读取数据库配置更简单。 |
| 同步失败保留最后一次成功结果,Apple 和 Google 失败互不影响。 | Google Play 同步依赖具有目标应用权限的 Service Account,凭据缺失时只跳过 Android,不会影响 iOS。 |
15.5 Ops 监控刷新
现状:OpsMonitorService每 60 秒跑 12 条聚合 SQL(push_messages/wa_inbound_messages/wa_outbound_messages/payment_webhook_events/commission_records/affiliate_orders/risk_events等表),算出积压量、超时时长、失败数等指标,写 Micrometer 监控指标,超阈值发飞书告警,还带一个"是否建议上 MQ"的自评估开关(讽刺的是,这个任务本身恰好属于"不太适合改成事件驱动"的那一半)。
接入 MQ 后(需要区分两类指标,不能一概而论):一半指标(积压量、最老等待时长、超时状态)本质是"在检测某件事没有按时发生",天生就得靠周期性重新比对当前时间才能算出来,MQ 帮不上忙,只能继续轮询;另一半指标(过去一小时失败数、过去一小时新增订单数)是对已经发生的离散事件做滚动计数,理论上可以让上游(比如 Push 失败、Webhook 失败发生的那一刻)顺手发个事件,由监控这边订阅累加,不用每次都重新 COUNT 一遍。
| 优点 | 缺点/风险 |
|---|---|
| "过去一小时XX次数"这类滚动计数指标,改成事件驱动累加后,不用每分钟对大表做一次 COUNT,数据库压力更小。 | 接近一半的指标(积压、超时、过期)结构上就是"检测异常的缺席",不存在天然事件可订阅,这部分无论如何都得保留周期扫描,不是"全部改成事件驱动"能覆盖的。 |
| 已有的告警节流(30 分钟内同一告警不重复发)、Feishu 通道都能直接复用,不用重新做。 | 把"12 条 SQL 一次性算完"拆成"一部分轮询 + 一部分事件累加"之后,监控代码本身的复杂度会上升(两套数据来源要合并成一份快照),现在这种"一个方法算完所有指标"的写法反而更简单、更容易看懂。 |
15.6~15.8 资金相关(Phase 2,暂不启动)
订单同步(Shopee/TikTok CAP/TikTok Creator/Lazada/Traveloka 共 5 个平台,各一个 @Scheduled 服务,15 分钟轮询)、佣金到期释放(CommissionSettlementService,5 分钟轮询)、Xendit 提现重试(WithdrawService.retryFundsHeldScheduled,10 分钟轮询)这三类任务,现状都已经有一套很成熟的资金安全机制在跑:
wallet_ledger.uk_ledger_ref(ref_type, ref_id, type)唯一约束 +user_wallets行锁(for update),是"同一笔钱不会入两次账"的根本保证。- 订单/佣金记录用业务键(平台+平台订单号)算出确定性 ID(SHA-256),配合
ON DUPLICATE KEY UPDATE,天然幂等。 - 状态流转全部用 CAS(
WHERE status=期望值)做,不是无条件覆盖,防止重复消息重复推进状态。 - 提现重试专门维护了一份"哪些第三方错误码确定没有真的扣款、可以安全重试"的白名单(
FUNDS_HOLD_FAILURE_CODES),其它错误一律当作失败处理,不敢重试——这条判断逻辑是这套机制里最关键、也最脆弱的一环,必须原封不动保留。
| 接入 MQ 的潜在优点 | 缺点/风险(这是暂不启动的原因) |
|---|---|
| 能把"检测到订单/佣金/提现需要处理"和"实际处理"解耦,理论上能提升处理时效、减轻单个长事务/单个调度线程的压力。 | MQ 是"至少一次投递",天然可能重复、可能乱序;上面这套 CAS+唯一约束+行锁机制虽然能兜住重复消费,但"乱序"是新增风险——比如订单状态 PENDING→COMPLETED→CANCELLED 这种转换,如果消息乱序到达,现在的轮询模式(每次重新拉取当前状态)天然没有这个问题,改成事件驱动之后必须显式处理乱序,否则可能算错钱。 |
| 把"扫描"和"处理"解耦后,理论上能减少数据库全表扫描频率。 | 现在的分布式锁(DistributedTaskLockExecutor)保证"同一时刻只有一个实例在处理"这个强顺序保证,换成多消费者并行消费之后需要重新设计成"按分区键(比如订单ID/用户ID哈希)顺序消费",复杂度和出错代价都显著高于 Phase 1 那 5 个非资金任务。 |