← 返回 Worker 当前架构与演进边界

Outbox + RocketMQ 接入技术方案

v3.0 · 2026-07-18 · 承接 worker-architecture-evolution-plan.html §8 的演进条件,给出具体落地方案
这不是要不要上消息队列的讨论文档,是"怎么做"的技术方案。broker 最终确定用阿里云 RocketMQ(决策过程和成本取舍见 §2)。Outbox 优先、先接非资金任务、消费者幂等这几条核心原则不变;§4 的可替换发布/订阅接口设计仍然保留,不是白做——以后如果 broker 还要再换,同样只改实现类。

1. 结论与范围

范围边界:第一阶段只迁移 worker-architecture-evolution-plan.html §3 中标注"中/低"风险的任务;标"资金高风险"的 7 个任务(5 个平台订单同步 + 佣金到期释放 + Xendit 提现重试)和标"高"风险的 TikTok Token 刷新,本方案不动,继续留在现有的轮询+锁机制上。每个任务的现状代码、接入 MQ 后的流程、逐项优缺点见 §15,本节只给结论。
第一阶段迁移(5 个)

Push 消息派发、优惠券过期扫描、TikTok Partner Campaign 商品同步、App Store 版本同步、Ops 监控刷新。

不碰钱
暂不迁移(8 个)

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/年 — 已确认接受,按此预算采购
采购口径:Phase 1 这 5 个非资金任务量级不大,选最小可用规格即可,不需要为将来 Phase 2(资金任务)预留容量买大规格——真到 Phase 2 评估通过时再单独升配或加实例,现在按最小够用规格采购。

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 个消费者业务逻辑一行都不用改
这不是过度设计:这两个接口就是 Relay 本来就该有的内部方法签名,抽出接口只是多了两行 interface 声明,不多花额外工作量,但换来了以后不被今天的技术选型锁死——下一节就直接用上这个设计。

5. 测试用 Redis Streams, 生产用 RocketMQ

正好用上 §4 的抽象:测试环境不用额外买一份 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 = "...") 按这个配置值选择装配哪一套实现,跟现有"环境变量驱动配置、无配置时安全降级"的风格一致。

这不只是省钱,还多验证了一层:测试环境用 Redis Streams 把 Outbox 表结构、Relay 轮询、5 个消费者的幂等逻辑完整跑一遍,能提前发现这套机制本身的 bug;到生产换成 RocketMQ 实现类时,需要重新验证的只是"RocketMQ 这个实现是否正确调用了 SDK",业务逻辑早就在测试环境验证过了,不是从零开始上生产验证。
2026-07-19 已实现并在 id-test 全链路验证通过: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_eventPushEventConsumer 通过跟轮询路径共用的原子租约(leaseMessageIfDue)消费,避免两条路径重复派发。全部行为收在 seahub.mq.push-dispatch.enabled 开关后面(生产默认 false,等 RocketMQ 采购完成再实现 §5 表格里的 RocketMqEventPublisher/RocketMqEventConsumer 并切换;id-test 已打开)。
id-test 实测:插入一条待发送 Push 消息 + 对应 outbox 事件后,outbox_event.status 变为 PUBLISHED,Redis Stream 消费者组 GID_pushentries-read 从 0→1、lag=0push_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 表设计

2026-07-18 修正:主键不用 BIGINT AUTO_INCREMENT项目里所有业务表(commission_recordswithdraw_requestsreferral_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)
);
  • idIds.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(idIds.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_DISPATCHPush 消息派发GID_push
COUPON_EXPIRE优惠券过期扫描GID_coupon
CAMPAIGN_PRODUCT_SYNCTikTok Partner Campaign 商品同步GID_campaign
APP_VERSION_SYNCApp Store 版本同步GID_appversion
OPS_SNAPSHOTOps 监控刷新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 快照消费者都是覆盖式写最新状态,重复消费无副作用。
结论: 第一阶段这 5 个消费者全部是"重复消费无副作用"或"有状态字段天然去重",不需要额外设计一张通用的消息去重表——这是选它们做第一阶段的另一个原因,复杂度最低。RocketMQ 的消费失败重试、死信队列是平台原生能力,不用自己实现。

10. 分阶段路线图

Phase 0 · 验证机制(已完成)

接 Push 消息派发 1 个任务,已在 id-test 验证 Outbox 表、Relay、Redis Streams、消费者组全链路跑通(见 §5 绿色标注),观察 1-2 周稳定后进入 Phase 1。

1 个任务·已实现
Phase 1 · 全量非资金任务

Phase 0 稳定后,把剩下 4 个任务(优惠券、商品同步、版本同步、Ops 监控)全部接入,同时下线对应的旧轮询代码路径(但保留代码,用开关控制,见 §10)。

4 个任务
Phase 2 · 资金相关(触发后再评估)

只有 §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 的价值要等真的有下游动作需要挂载时才体现出来,现在接等于"为将来可能用到的扩展性买单"。
结论:技术上可行、幂等成本低,但当前实现已经足够简单高效,收益主要是"面向未来扩展",不是"解决现有问题"。按 Phase 1 原计划接入,但不用抱有"接了之后会变快/变可靠"的预期。

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 都跑完了)。
结论:Phase 1 里优先级应该提到最前面——不只是"顺便接一下 MQ",而是直接解决已经实测验证过的调度线程池阻塞问题,收益比其它几个任务都明确。

15.4 应用商店版本同步(注:跟 H5 包分发是两回事,见下方澄清)

澄清:这个任务由管理后台服务每小时查询 Apple iTunes Lookup 和 Google Play Android Publisher API,获取原生 App 已发布版本并持久化。iOS 版本当前用于控制"苹果登录"功能是否对老版本客户端隐藏(AppRuntimeConfigService.applyAppleLoginGate)。H5 小程序包的分发(h5_package_versions 等表)是另一套完全独立、纯被动拉取的机制——客户端每次启动主动调 GET /api/v1/app/h5-package/latest 查询、下载更新,服务端没有任何轮询,不属于这次 MQ 迁移的候选范围。

改造后现状:轮询任务已迁移到管理后台服务。管理后台每小时同步 iOS App Store 和 Android Google Play production track,将最后一次成功结果写入 runtime_feature_configSYSTEM 配置;业务服务只读配置,不再保存实例内存版本或直接访问应用商店。

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 一次性算完"拆成"一部分轮询 + 一部分事件累加"之后,监控代码本身的复杂度会上升(两套数据来源要合并成一份快照),现在这种"一个方法算完所有指标"的写法反而更简单、更容易看懂。
结论:这是 5 个 Phase 1 任务里跟"事件驱动"这个概念最不匹配的一个——它的核心价值就是"周期性检查系统里是不是有东西卡住了",这件事本身离不开轮询。建议保持现状不改造,或者只挑"最近一小时失败数"这类真正是滚动计数的指标做小范围事件化,不用整体迁移。

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 个非资金任务。
结论:技术上不是不能做,但"资金准确性"这个要求下,乱序处理、重复投递、消费者并行度控制这些问题都要从头设计和验证,风险和收益不成正比,继续维持"只有触发条件真正出现才评估"的原判断,不纳入这次范围。