MQ 问题不能只问“消息发没发”。完整问题是:业务事务何时提交、消息何时发布、路由到了哪个消费者、消费者怎样定位业务单、失败返回 ACK 还是 NACK、重复投递是否安全、最终状态是否需要补偿。

本文覆盖 DGJ2.0 与支付中心、订单中心、调拨/配送中心、OA、商品中心、SAAS、App、渠道平台、对账单任务之间的异步链路。目标是让开发、测试和运维可以从 destination、routing key、业务单号、消息体任一线索,追到生产者、消费者、表、状态和补偿入口。

1. 业务目标

  • 建立 destination -> routing key -> producer -> consumer -> callback -> 业务表 的完整地图。
  • 写清主要消息 payload 的字段、关联键和单位,而不是只列事件名称。
  • 区分 Direct、Topic、延迟队列、外部回调和普通 CLI 任务。
  • 解释 ACK、NACK、异常被吞掉、主动忽略、迁移后空消费的差异。
  • 明确消息 ID、业务幂等键、状态前置条件、数据库唯一键各自解决什么问题。
  • 给出“未消费、重复消费、部分成功、乱序、迟到、毒消息”的排查和恢复步骤。
  • 补偿前先证明可重跑,补偿后核对主表、明细、库存、资金和下游状态。

2. 边界与基本概念

2.1 本文包含

  • application/KzData/Enums/MqEventEnums.php 中的主要 destination 和 routing key。
  • application/Services/Mq/MqSer.php 中的主要生产者。
  • application/controllers/tasks/*Notify.php 中的主要消费者。
  • 秒杀、对账单、库存预警等补偿/异步任务。
  • 消息回调事务、幂等、重试、延迟、人工重放和验证。

2.2 本文不直接给出

  • RabbitMQ 账号、密码、主机等敏感配置值。
  • 生产队列积压量、消费者实例数和真实死信策略;这些必须在目标环境确认。
  • 外部系统内部的消费实现,只描述 DGJ 已知的请求/响应边界。
  • 直接可执行的生产数据修复 SQL;本文只给只读检查和安全门槛。

2.3 四个不能混用的标识

标识示例用途是否天然幂等
destinationpaycenter_notify队列/主题逻辑命名空间否
routing keypaycenter_pay_result把事件路由到回调否
MQ message ID生产端大量使用 uniqid()区分本次发送实例否,同一业务重发会产生新 ID
业务幂等键sourceOrderNo、requestId、渠道单号识别同一业务事实需要代码、状态或唯一键实现

3. 消息生命周期总图

flowchart LR
    A["业务事务"] --> B["构造 payload"]
    B --> C["MqSer / 外部生产者"]
    C --> D["Exchange"]
    D --> E["destination.routing-key"]
    E --> F["Queue"]
    F --> G["tasks Consumer"]
    G --> H["registryCallback"]
    H --> I["业务状态/明细/库存/资金"]
    I --> J{"处理结果"}
    J -->|ACK| K["消息确认"]
    J -->|NACK| L["重投/死信,依环境配置"]
    J -->|异常被捕获后 ACK| M["消息丢弃,靠告警或补偿"]

3.1 正常消费的七个证据

  1. 生产业务事实存在。
  2. 生产者方法确实执行到 publish()。
  3. destination 和 routing key 与消费者绑定一致。
  4. 消费进程正在运行。
  5. 回调日志能找到业务键。
  6. 数据库事务提交,业务状态和副作用一致。
  7. 回调返回 ACK,或重复消息按设计被安全 ACK。

只有第 7 项没有前六项,不能证明业务完成;只有业务表变化但没有 ACK,也可能在稍后重复消费。

4. destination 总表

destination 常量值方向主要业务
DEST_DGJdgjDGJ 内部/回到 DGJSAAS 消息、库存事件、新站、采购关闭、延迟任务
DEST_ODCordercenterDGJ -> 订单中心采购订单创建
DEST_ODC_NOTIFYordercenter_notify订单中心 -> DGJ审核、售后、取消、自制订单
DEST_DISPATCHCENTER_NOTIFYdispatchcenter_notify双向约定配送出库、物流、自提、退货确认
DEST_PAYCENTER_NOTIFYpaycenter_notify支付中心 -> DGJ支付、退款、授信、提现、异常
DEST_ITEMCENTER_NOTIFYitemcenter_notify商品中心 -> DGJ限购、套包、价格、敏感词
DEST_OA_NOTIFYoa_notifyOA -> DGJ审批通知
DEST_OA_RESULT_NOTIFYoa_result_notifyDGJ/OA 结果同步OA 结果转发
DEST_DGJ_NOTIFYdgj_notifyDGJ -> App/SAAS/下游销售、出库、活动、客户、机器人
DEST_DGJ_CHANNELdgj_channelDGJ 渠道内部人工开单、发货、退款审核、订单审核
DEST_DGJ_STATEMENTdgj_statementDGJ 对账任务对账单生成和收款状态
DEST_DGJ_STATEMENT_STRATEGYdgj_statement_strategy策略任务对账策略同步
DEST_DGJ_STATEMENT_SENDdgj_statement_send发送任务对账单发送
DEST_QPAYCENTER_NOTIFYqpaycenter_notify结算中心 -> DGJ全车件支付、退款、提现
DEST_IMCENTER_NOTIFYimcenter_notifyIM -> DGJApp 业务消息
DEST_UDESKdgj_udeskDGJ -> Udesk客户异步同步

5. 生产者模型

5.1 核心代码

文件作用
application/Services/Mq/MqSer.php通用业务消息生产者集合
application/Services/SyncOrder/SyncOrderSer.php销售单/出库单同步编排
application/Services/SyncOrder/SyncOrderBuilder.php构造下游报文
各领域 Service决定发布时机和 payload
application/config/mq.phpMQ 连接和系统配置;不可把值写入文档

5.2 Direct 与 Topic

类型代码工厂当前用途特点
DirectcreateDirectProducerDGJ 内部任务、库存、对账、延迟事件等精确路由
TopiccreateTopicProducer外部通知、订单中心回调、App/SAAS 同步、渠道等

业务代码使用哪种 Producer 不足以证明 RabbitMQ 真实绑定;exchange、queue 和 binding 必须在环境中核验。

5.3 消息 ID 风险

MqSer 的大多数方法把 uniqid() 传给 Producer:

createTopicProducer(destination, routingKey, uniqid())
createDirectProducer(destination, routingKey, uniqid())

这意味着:

  • 同一业务操作重发两次会得到两个不同消息 ID。
  • 消费者不能仅按 message ID 判断业务重复。
  • 需要以订单号、来源单号、明细行号、任务 ID、状态或数据库唯一键做幂等。
  • 人工重放时更不能假设旧 message ID 会阻止重复副作用。

5.4 发布与事务的三种边界

flowchart TD
    A["业务写库"] --> B{"MQ 在何时发送"}
    B -->|事务提交后| C["库成功、发送失败 -> 需要补偿"]
    B -->|事务提交前| D["消息可能先被消费 -> 读不到已提交数据"]
    B -->|事务内部同步发送| E["MQ 超时拉长事务和锁"]
    C --> F["Outbox/任务表/状态扫描更可靠"]
    D --> G["消费者重试或延迟"]
    E --> H["确认回滚后消息是否已经发出"]

本文只能从调用位置静态判断单个链路。任何重要改动都应在具体 Service 中确认事务开始、提交和 publish() 的先后关系。

6. 支付中心消息

6.1 消费入口

destination: paycenter_notify
consumer: application/controllers/tasks/PayCenterNotify.php
start: php index.php tasks/PayCenterNotify/consume

6.2 routing key 与回调

routing key回调主要业务事实
paycenter_pay_resultpaycenterPayResult采购/预订单支付结果
paycenter_balance_adjustpaycenterBalanceAdjust成本调整
paycenter_refund_account_withdraw_resultpaycenterRefundWithdrawResult待退款账户提现结果
paycenter_baitiao_return_pay_resultpaycenterBaitiaoReturnPayResult授信还款支付结果
paycenter_baitiao_expirepaycenterBaitiaoExpire授信到期提醒
paycenter_pay_exceptionpaycenterPayException支付抵扣等异常
paycenter_installment_expirepaycenterInstallmentExpire分期到期
paycenter_refund_resultpaycenterRefundResult退款结果
paycenter_baitiao_offline_returnpaycenterBaitiaoOfflineReturnPayResult线下授信还款/退款结果
paycenter_baitiao_service_overdue_changepaycenterBaitiaoServiceOverdueChange服务站逾期状态变化

6.3 支付成功 payload 骨架

{
  "sourceOrderNo": "<采购单或预订单号>",
  "sourceOrderType": "<01采购单/04预订单>",
  "payOrderNo": "<支付中心支付单号>",
  "mainPayTypeCode": "<主支付方式>",
  "payStatus": "02",
  "payStatusMsg": "支付成功",
  "payTime": "YYYY-MM-DD HH:mm:ss",
  "detailTotalAmount": "<支付中心金额口径>",
  "details": [
    {
      "payAccountTypeCode": "<账户类型>",
      "payAccountTypeDesc": "<账户说明>",
      "payAmount": "<分账户金额>"
    }
  ]
}

必须与支付中心合同确认金额单位。DGJ 内部不同表可能使用元、分或历史字符串金额,不能只凭字段名换算。

6.4 支付回调数据流

sequenceDiagram
    participant P as 支付中心
    participant Q as paycenter_notify
    participant C as PayCenterNotify
    participant O as 采购/预订单
    participant F as Payment/PaymentInfo
    participant R as TaskRecord
    P->>Q: paycenter_pay_result
    Q->>C: payload
    C->>C: payStatus 必须为 02
    C->>O: 按 sourceOrderNo/sourceOrderType 找单
    C->>O: 检查当前状态与支付方式
    C->>F: 写支付主表、明细、采购支付关系
    C->>O: 推进订单状态/额度/活动订单
    C->>C: 提交事务
    C->>R: 查询是否有后置任务记录
    C-->>Q: ACK 或 NACK

6.5 支付幂等和危险分支

情况当前行为关注点正确判断
payStatus != 02返回 NACK 并告警可能反复重试,需确认失败状态是否应该 NACK
业务单不存在异常 -> NACK可能是消息早到、错环境或单号合同错误
已经支付成功又收到相同回调状态校验可能抛错不能把重复成功回调当普通系统故障
超时关闭后迟到支付可能恢复订单或形成状态冲突必须核对活动锁量、支付事实和订单中心状态
事务提交后后置任务失败主业务已成功不能直接回滚已提交事实,应补后置任务

6.6 核心表

表常量用途关联键
SCM_PAYMENT付款/支付主单支付单号、业务单号、sid
SCM_PAYMENT_INFO支付账户明细payment id、账户类型、金额
SCM_PO_ORDER_PAYINFO采购单与支付关系采购单 ID、支付单号
SCM_PO_ORDER采购主单billNo、状态、来源类型
SCM_PRE_ORDER预订单预订单号、状态
TaskRecord 对应表支付后置任务记录sid、订单号

7. 订单中心消息

7.1 消费入口

destination: ordercenter_notify
consumer: application/controllers/tasks/OrderCenterNotify.php

7.2 注册事件

事件组routing key业务
采购正向ordercenter_order_audit采购审核
在线支付ordercenter_order_cancel订单取消
采购售后ordercenter_aftersale_audit/update/close/confirm_warehouse售后推进
采购售后ordercenter_aftersale_refund_finish退款完成
采购售后ordercenter_aftersale_return_quota_expire退货额度过期
预订单ordercenter_preorderrefund_create/cancel预订单售后
预订单ordercenter_preorder_close预订单关闭
发货方式ordercenter_self_change_direct自提/配送切换
OPS 自制单ordercenter_order_sync_dgj自制订单下发
OPS 自制单ordercenter_order_close自制订单行关闭

7.3 所有权迁移风险

以下四个回调当前在方法开头直接 ACK,后面的 DGJ 历史代码不可达,并注明已迁移到采购服务:

  • orderCenterAfterSaleAudit
  • orderCenterAfterSaleClose
  • orderCenterAfterSaleUpdate
  • orderCenterAfterSaleConfirmWareHouse
flowchart LR
    A["订单中心售后消息"] --> B["DGJ OrderCenterNotify"]
    B --> C["立即 ACK"]
    C --> D["DGJ 历史处理代码不可达"]
    A --> E["dgj-purchase-service 应承担实际消费"]

排查这四类事件时:

  1. DGJ 返回 ACK 不能证明售后已处理。
  2. 必须查采购服务的队列绑定、消费日志和数据库。
  3. 若两个服务都绑定同一 Topic,要确认是各自独立队列还是竞争消费。
  4. 不要通过删除 DGJ 的早返回来“恢复”历史逻辑,除非完成跨服务所有权设计。

7.4 订单中心主数据关系

flowchart TD
    A["订单中心订单号/售后号"] --> B["outOrderNo/outAftersaleNo"]
    B --> C["DGJ SCM_PO_ORDER.billNo"]
    C --> D["SCM_PO_ORDER_INFO"]
    D --> E["sourceLineCode / srcOrderEntryId"]
    E --> F["关闭明细/数量表/库存预警事件"]

消息体中的 sourceChannel、外部售后号、行号、审核结果和关闭明细必须与 DGJ 订单来源、交易类型和当前状态共同校验。

7.5 事务和后置 MQ

历史回调中常见模式是:

  1. 开启采购 Model 事务。
  2. 更新采购主单、明细和关闭记录。
  3. 提交事务。
  4. 发送库存事件,重算库存预警。

如果第 4 步失败,采购事实可能已提交但预警未刷新。此时应补发库存事件或重算预警,不应重做整个售后业务。

8. 调拨/配送中心消息

8.1 入口和事件

routing key回调业务结果
dispatchcenter_delivery_stock_outdispatchCenterDeliveryStockOut采购配送出库/状态推进
dispatchcenter_logistice_info_notifydispatchCenterLogisticeInfoNotify物流信息同步
dispatchcenter_self_pick_up_notifydispatchCenterSelfPickUpNotify自提通知
退货关闭事件dispatchCenterReturnClose退货确认/关闭
dispatchcenter_order_tag_biz_change订单标签回调业务标签变化

8.2 出库链路

sequenceDiagram
    participant D as 调拨/配送中心
    participant C as DispatchCenterNotify
    participant P as 采购主单/明细
    participant O as 采购出库记录
    participant I as 库存/后置同步
    D->>C: stock_out payload
    C->>P: 按单号、来源、状态找采购单
    C->>O: 写出库数量和出库记录
    C->>P: 更新待出库/配送中/完成
    C->>C: 提交事务
    C->>I: 库存、消息、活动等副作用
    C-->>D: ACK/NACK

8.3 重复与迟到消息

  • 已完成的出库事实再次到达时,应通过已出库数量和状态避免重复累加。
  • 订单已取消或关闭后收到出库消息,需要明确谁拥有最终事实,不能无条件恢复订单。
  • 物流通知可能先于业务出库落库,需要消费者允许重试或暂存。
  • 部分分支对“目标不存在/已处理”直接 ACK,这是主动终止重试,不等于处理成功。

9. OA 消息

9.1 入口

destination: oa_notify
routing key: oa_audit
consumer: tasks/OaNotify::oaAudit

OaNotify 根据模板/回调 code 分发到不同业务。秒杀活动使用 FlashActivityEnums::OA_TEMPLATE_CODE,然后进入 FlashActivitySer::callback20260526。

9.2 秒杀 OA payload 骨架

不同 OA 平台版本字段可能不同,DGJ 至少需要识别:

{
  "callback": "<模板或业务回调code>",
  "status": "<审批结果>",
  "form_data": {
    "activity_id": "<秒杀活动ID>"
  },
  "operator": "<审批人或系统>"
}

实际字段必须以生产消息和 oaAudit 分支为准,不应根据此骨架直接伪造回调。

9.3 OA 状态数据流

flowchart LR
    A["OA 审批结果"] --> B["OaNotify::oaAudit"]
    B --> C{"callback/template"}
    C -->|秒杀| D["FlashActivitySer::callback20260526"]
    D --> E["FLASH_ACTIVITY 状态"]
    D --> F["FLASH_ACTIVITY_LOG"]
    C -->|其他业务| G["对应 Service 回调"]
    B --> H{"异常"}
    H -->|抛出| I["NACK"]
    H -->|正常| J["ACK"]

重复 OA 回调必须满足:已通过再次通过不重复启用、不重复写数量;驳回后迟到通过不能绕过当前状态;活动已结束后审批结果不能把活动恢复运行。

10. 商品中心消息

10.1 事件地图

routing key回调影响
itemcenter_sensitive_word_changesensitiveWordChange敏感词数据/缓存
itemcenter_sku_limit_purchase_configskuLimitChangeSKU 限购配置
itemcenter_package_status_syncsyncPackageList套包主明细状态
itemcenter_item_price_enableitemPriceEnable结算价生效
itemcenter_item_price_disableitemPriceDisable结算价失效
itemcenter_guide_price_limit_changesalesPriceLimitNotice销售限价

10.2 数据传播

flowchart TD
    A["商品中心事件"] --> B["ItemCenterNotify"]
    B --> C["商品/套包/限购/价格表"]
    C --> D["商品缓存"]
    C --> E["ES/搜索数据"]
    C --> F["采购可买校验"]
    C --> G["销售可卖校验"]
    C --> H["活动商品校验"]

10.3 ACK 策略风险

ItemCenterNotify 很多校验失败、数据不存在和异常分支仍返回 ACK。这样做可以避免毒消息无限重试,但代价是:

  • 消息不会自动再次投递。
  • 失败必须依赖日志、告警或定时全量同步发现。
  • 重发前要确认依赖数据已经准备好。
  • 消费成功不能只看队列清空,还要查缓存/索引和最终业务效果。

套包同步使用数据库事务,但异常回滚后也可返回 ACK。该链路应配套失败记录或重建任务。

11. DGJ 内部 dgj 消息

11.1 消费入口

SaasOrderNotify 消费 DEST_DGJ 下多类事件:

routing key回调业务
create_saas_ordercallbackApp/网页站内消息
udesk_customer_createudeskCustomerCreateUdesk/IM 账户创建
dispatch_pause_canceldispatchPauseCancel延迟关闭暂停取消
change_qtychangeQty采购明细数量异步变化
inventory_eventinventory_event实时库存/在途变化后重算预警
po_close_order_confirmpo_close_confirm采购关闭确认
search_log_addsearch_log_add搜索日志
new_stationnew_station新服务站初始化
new_application_opennewApplicationOpen新站应用开通
dgj_po_order_delivery_stock_outdelivery_stock_out采购配送出库
OPS 创建/关闭事件opsCreateOrder、opsCloseOnlineOrderOPS 自制采购单

11.2 库存事件 payload

生产者 sendInventoryEvent 明确构造:

{
  "inventoryData": [
    {
      "invId": "<库存商品ID>",
      "qty": "<变动数量及其他上下文字段>"
    }
  ],
  "iid": "<库存流水或删除撤销关联ID,未使用为0>",
  "billType": "<单据类型>",
  "transType": "<交易类型>",
  "sid": "<服务站ID>",
  "po_order_id": "<采购单ID,未使用为0>"
}

11.3 库存预警重算公式

inventory_event 不直接重做库存流水,而是收集受影响 invId,读取实时库存和采购在途后重算预警:

warning_num = 实时库存 qty + 在途 waitQty - 预警下限 lowQty
suggest_num = 预警上限 highQty - 实时库存 qty - 在途 waitQty
flowchart TD
    A["inventoryData / iid / po_order_id"] --> B["汇总 invId"]
    B --> C["商品状态和预警上下限"]
    C --> D["SCM_INVENTORY_REAL_TIME"]
    C --> E["采购在途数量"]
    D --> F["warning_num / suggest_num"]
    E --> F
    F --> G["SCM_INVENTORY_WARNING batchUpsert"]

空 inventoryData、iid=0、po_order_id=0 时直接 ACK;找不到有效商品也 ACK。排查预警不更新时,应确认 payload 至少携带一种可定位方式。

11.4 OPS 批处理风险

opsCreateOrder 对 ids 循环,单个 ID 失败只记录日志后继续,最终整体 ACK;opsCloseOnlineOrder 也会把失败状态落到日志表后整体 ACK。

这属于“部分成功 + 不自动重试”模型:

  • 不能根据 ACK 判断全部 ID 成功。
  • 必须查每个 ID 的处理状态。
  • 补偿应只重跑失败 ID,不能重发整批。

12. DGJ 发往 App/SAAS 的消息

12.1 主要生产者

MqSer 方法routing keypayload 关键键
sendSaOrderCloseToAPPapp_sa_order_closeorder_no
sendSaOrderCutOffToAPPapp_sa_order_cut_offorder_no
sendApplyReturnFinishToAPPapp_apply_return_finishorder_id、order_no、return_type
sendApplyReturnCheckedToAPPapp_apply_return_checkedorder_id、order_no
sendSaOrderOutToMqapp_sa_order_out调用方传入出库数据
sendSaOrderStatusToMqapp_sa_order_status配送状态数据
sendBindMoveStorageToAPPapp_bind_move_storage微仓绑定数据
sendActivityStartToSaasactivity_startsids
sendActivityEndToSaasactivity_endsids

新销售同步还会使用 SyncOrderBuilder 构造维修厂销售单/出库单事件。排查下游未更新时,要区分旧 APP_SA_* routing key 和新 garage_repair_* 事件。

12.2 下游同步链路

flowchart LR
    A["销售单/出库单事务"] --> B["SyncOrderBuilder"]
    B --> C["统一来源单号、明细、状态"]
    C --> D["dgj_notify"]
    D --> E["App/SAAS 消费者"]
    E --> F["下游订单事实"]
    F --> G["下游后续重试/回执"]

核对字段至少包括:DGJ 订单号、出库单号、来源单号、服务站、客户、明细行、SKU、数量、状态、事件时间。只按订单号查不到时,应继续按来源单号和出库单号追踪。

13. 渠道消息

13.1 事件与处理

routing key回调典型结果
channel_order_manual_createchannelOrderManualCreate创建平台/DGJ 订单关系
channel_order_deliverychannelOrderDelivery同步渠道发货
channel_refund_actionchannelRefundAction售后入库审核动作
channel_order_approvechannelOrderApprove渠道订单审核

13.2 ACK 语义

渠道消费者会对某些“记录不存在、状态已处理、不需要动作”情况直接 ACK。这是业务性忽略,目的通常是阻止无意义重试。

flowchart TD
    A["渠道消息"] --> B{"关系单存在?"}
    B -->|否| C["记录日志并 ACK"]
    B -->|是| D{"状态允许?"}
    D -->|已完成/无需处理| E["ACK"]
    D -->|允许| F["调用渠道 Service"]
    F -->|成功| G["ACK"]
    F -->|可重试异常| H["部分方法 NACK"]

若业务方认为“记录不存在”应该重试,就不能只改 ACK 为 NACK,还要解决消息早到、关系表创建时序和无限毒消息问题。

14. 对账单异步任务

14.1 三条队列

destinationrouting key消费者业务
dgj_statement_strategystatement_strategy_syncStatementStrategyConsumer策略同步
dgj_statementstatement_bill_createStatementNotify对账单生成
dgj_statementstatement_bill_receive_statusStatementNotify收款状态同步
dgj_statement_sendstatement_sendStatementSendConsumer对账单发送

14.2 任务化数据流

flowchart LR
    A["对账策略"] --> B["任务初始化"]
    B --> C["StatementBillTask"]
    C --> D["statement_bill_create"]
    D --> E["生成主单/明细/明细项"]
    E --> F["核对/生成PDF"]
    F --> G["statement_send"]
    G --> H["发送记录/机器人/微信"]
    E --> I["statement_bill_receive_status"]
    I --> J["收款状态"]

StatementTaskSer 和发送记录表是天然补偿抓手。补偿应按任务 ID 或对账单 ID 执行,不能仅按日期整批重发。

14.3 部分成功

对账单批量发送可能出现:部分对账单不存在、未核对、机器人离线、客户未绑定微信、指定绑定无效。正确结果应包含成功、跳过和失败明细,而不是把整批简化为一个布尔值。

15. 延迟队列

15.1 已确认的 TTL + DLX 模型

延迟 routing keyTTL到期后 routing key用途
new_station_delay_300300000 msnew_station新站延迟初始化
dispatch_pause_cancel_delay_6060000 msdispatch_pause_cancel暂停取消延迟关闭
dispatch_out_delay_20名称表达 20 秒需要在环境/初始化代码确认出库重新消费
sequenceDiagram
    participant P as 生产者
    participant D as 延迟队列
    participant X as Direct Exchange
    participant C as 正式消费者
    P->>D: delay routing key
    Note over D: x-message-ttl
    D->>X: dead-letter routing key
    X->>C: 正式事件
    C-->>X: ACK/NACK

15.2 延迟队列注意事项

  • 队列声明发生在消费者初始化代码,进程从未启动时队列可能不存在。
  • 代码中的 TTL 不代表线上队列参数一定一致,应读取实际队列声明。
  • 延迟消息到期时业务状态可能已被人工操作改变,消费者必须再次检查状态。
  • 重复发送延迟消息会产生多次到期事件,业务层必须幂等。

16. ACK、NACK 和异常语义

16.1 五种结果

结果消息层含义业务层含义
正常 ACK不再投递本次成功或被安全忽略
重复消息 ACK不再投递业务事实已存在,不需重复执行
条件不满足 ACK不再投递消息被丢弃,可能需要人工/定时补偿
异常 NACK依 MQ 配置重试/死信期望稍后重新处理
未捕获异常/进程退出是否重投依客户端和 broker消费结果不确定,需要同时查队列和事务

16.2 ACK 不能证明什么

  • 不能证明每个批次子项成功。
  • 不能证明后置 MQ 成功。
  • 不能证明下游系统已落库。
  • 不能证明消费者没有走“迁移后直接 ACK”。
  • 不能证明缓存、索引、报表已经更新。

16.3 NACK 前必须考虑

flowchart TD
    A["消费者异常"] --> B{"重试能改变结果吗?"}
    B -->|依赖暂时不可用| C["NACK 合理"]
    B -->|永久脏数据| D["无限 NACK 形成毒消息"]
    B -->|业务已部分成功| E["直接重试可能重复副作用"]
    D --> F["死信/失败记录/人工修复"]
    E --> G["先识别已完成步骤,再补剩余步骤"]
    C --> H["确认退避、次数和告警"]

17. 幂等设计字典

17.1 幂等层级

层级手段适合解决
消息层message ID 去重同一投递实例重复;当前生产端随机 ID 不足以覆盖业务重发
请求层requestId、来源事件 ID外部系统同一请求重发
业务层来源单号 + 事件类型 + 明细行同一业务事实
状态层只允许从前置状态迁移乱序和重复推进
数据库层唯一键/条件更新并发下最终兜底
副作用层资金/库存流水唯一引用防重复记账和重复增减库存

17.2 推荐幂等键

业务推荐组合
支付成功sourceOrderNo + payOrderNo + payStatus
退款结果来源退款单号 + 支付中心退款单号 + 状态
采购审核订单中心订单号 + 审核版本/结果
售后行关闭售后单号 + sourceLineCode + 关闭事件
配送出库配送单号 + 出库明细行/批次
渠道发货渠道订单号 + 物流/发货批次
库存预警事件不一定需要事件唯一,重算应天然幂等
对账单生成策略 + 服务站 + 客户 + 对账周期
对账单发送对账单 ID + 发送类型 + 目标绑定 ID + 任务批次

18. 常见故障树

18.1 消息没有到消费者

flowchart TD
    A["没有消费日志"] --> B{"生产日志有 publish 吗?"}
    B -->|否| C["业务分支未发送/事务异常"]
    B -->|是| D{"destination/routing key 正确?"}
    D -->|否| E["常量或调用参数错误"]
    D -->|是| F{"队列绑定存在?"}
    F -->|否| G["部署/声明/环境配置"]
    F -->|是| H{"消费者进程运行?"}
    H -->|否| I["启动/进程守护/崩溃"]
    H -->|是| J["积压、未注册回调、日志路径"]

18.2 消息重复消费

  1. 查生产端是否执行了两次 publish()。
  2. 查消费者第一次是否 NACK 或进程在 ACK 前退出。
  3. 查队列是否有重试/死信回投。
  4. 比较业务键,不要只比较随机 message ID。
  5. 查第一次事务是否提交。
  6. 查副作用是否有唯一引用。
  7. 只补未完成步骤。

18.3 队列清空但业务没变化

优先查以下主动 ACK 分支:

  • 订单中心售后四个已迁移方法直接 ACK。
  • 商品中心异常/数据不存在后 ACK。
  • 渠道记录不存在或状态无需处理后 ACK。
  • OPS 批次子项失败但整体 ACK。
  • 库存事件没有任何定位字段时 ACK。

18.4 部分成功

flowchart LR
    A["消息包含 N 个子项"] --> B["逐项处理"]
    B --> C["成功集合"]
    B --> D["失败集合"]
    C --> E["已产生业务副作用"]
    D --> F["日志/失败状态"]
    F --> G["按失败 ID 精确补偿"]
    G --> H["不得整批盲重放"]

19. 排查 SOP

19.1 收集事实

先记录:

环境:
发现时间:
destination:
routing key:
业务单号/来源单号:
sid:
消息 ID(如有):
请求或事件时间:
预期状态:
实际状态:
是否有库存/资金副作用:

19.2 从代码确认路由

rg -n "DEST_PAYCENTER_NOTIFY|TYPE_PAYCENTER_PAY_RESULT" application/KzData/Enums/MqEventEnums.php
rg -n "createTopicConsumer|registryCallback" application/controllers/tasks/PayCenterNotify.php
rg -n "createDirectProducer|createTopicProducer" application/Services/Mq/MqSer.php

19.3 检查消费者进程

ps aux | rg 'index.php tasks/.+Notify/consume|Statement.*Consumer'

具体进程管理器名称、容器和日志路径以部署手册为准。

19.4 查日志

组合条件:

消费者类名 + 回调方法
routing key + 业务单号
sourceOrderNo + payOrderNo
sid + 采购单号/销售单号
消息到达时间前后 5-10 分钟

不要只搜“error”;很多消费者使用 info/warn 输出失败,部分还使用标准输出。

19.5 查数据库事实

以下 SQL 只读,分表后缀和字段以目标环境表结构为准:

-- 采购主单
SELECT id, sid, billNo, billStatus, orderType, transType, srcOrderId, srcOrderType, updateTime
FROM t_scm_po_order_<shard>
WHERE billNo = '<业务单号>';

-- 采购支付关系
SELECT *
FROM t_scm_po_order_payinfo_<shard>
WHERE po_order_id = <采购单ID>;

-- 支付主单和明细
SELECT * FROM t_scm_payment_<shard> WHERE srcOrderNo = '<业务单号>';
SELECT * FROM t_scm_payment_info_<shard> WHERE payment_id = <支付主单ID>;

-- 秒杀订单
SELECT id, order_no, po_order_id, status, locked_qty, used_qty, update_time
FROM t_flash_order
WHERE order_no = '<秒杀单号>' OR po_order_id = <采购单ID>;

-- 库存预警
SELECT sid, inv_id, sku_id, warning_num, suggest_num
FROM t_scm_inventory_warning
WHERE sid = <sid> AND inv_id IN (<invIds>);

19.6 判断恢复方式

当前事实恢复方式
从未发布修复发送条件后补发
已发布但无绑定修复队列绑定,再从可靠原始数据补发
消费 NACK 且无副作用修复依赖后允许重试
消费 NACK 且部分提交先识别已完成步骤,再做定向补偿
消费 ACK 但主动忽略修复业务前置数据后使用领域补偿入口
已完成但下游未同步只补下游通知,不重做主业务
重复副作用停止重放,先做账实核对和人工修复方案

20. 补偿任务地图

20.1 秒杀

命令作用
php index.php tasks/FlashSaleTask/expireWaitPayOrders 100释放超时待支付订单
php index.php tasks/FlashSaleTask/refreshGoodsStatus 500刷新商品有效状态
php index.php tasks/FlashSaleTask/refreshInvalidCarts 500失效购物车
php index.php tasks/FlashSaleTask/compensatePaidOrders 100补支付成功状态
php index.php tasks/FlashSaleTask/compensateClosedOrders 100补关闭库存释放
php index.php tasks/FlashSaleTask/runCompensate 500顺序执行常规补偿

20.2 消费者内置 resend

部分消费者包含 resend(),例如支付、订单中心、商品中心、配送和 SaasOrderNotify。这类方法通常从 /tmp 文件读取消息并重新发布,属于人工运维工具,不是自动补偿保证。

执行前必须:

  1. 阅读具体 resend() 的允许 routing key。
  2. 确认输入文件格式和每行消息体。
  3. 去除已成功业务键。
  4. 在非生产环境验证一条。
  5. 记录原始消息、执行人、时间和数量。
  6. 执行后逐条核对业务结果。
  7. 删除或归档临时文件,避免误重跑。

20.3 补偿安全门

flowchart TD
    A["准备补偿"] --> B["冻结原始证据"]
    B --> C["导出业务键和当前状态"]
    C --> D["证明操作幂等"]
    D --> E["先单条/小批量"]
    E --> F["核对主表、明细、副作用"]
    F --> G{"结果正确?"}
    G -->|否| H["立即停止,恢复/人工分析"]
    G -->|是| I["扩大批次"]
    I --> J["最终对账和执行记录"]

21. 高风险清单

风险静态证据后果
随机 message ID 不等于业务幂等生产者普遍使用 uniqid()业务重发仍可能重复执行
四类订单售后消息直接 ACKOrderCenterNotify 方法首行返回DGJ 队列清空但实际所有权在采购服务
商品中心异常也可能 ACK多个 catch/校验分支自动重试缺失,缓存/价格可能陈旧
OPS 批处理子项失败整体 ACKopsCreateOrder 循环 continue形成部分成功
消费后发布后置 MQ多个回调提交后 sendInventoryEvent主业务成功、派生数据失败
延迟队列由消费者代码声明init*DelayedQueue消费进程未启动时基础设施可能未建立
手工 resend 使用临时文件各 Notify 的 resend()文件残留或整批重放造成重复副作用
失败输出级别不统一info/warn/error/stdout 并存仅查 error 会漏问题

22. 回归清单

22.1 生产者

  • [ ] 正常业务只发布一次预期事件。
  • [ ] 事务回滚时不会留下不可撤销消息。
  • [ ] 事务提交后发布失败有补偿抓手。
  • [ ] payload 包含稳定业务键、sid、来源和事件时间。
  • [ ] 新字段对旧消费者可选兼容。
  • [ ] 不在 payload 和日志中发送密钥或隐私数据。

22.2 消费者

  • [ ] 正常消息推进主表、明细和副作用。
  • [ ] 相同业务消息重复投递不会重复记账、扣库存或建单。
  • [ ] 乱序消息不会把终态回退到中间态。
  • [ ] 业务单暂时不存在时重试策略明确。
  • [ ] 永久脏消息不会无限 NACK。
  • [ ] ACK、NACK 和主动忽略都有日志和指标。
  • [ ] 批处理能记录每个子项结果。
  • [ ] 事务提交与 ACK 先后关系经过故障注入验证。

22.3 延迟和补偿

  • [ ] 延迟队列 TTL、DLX 和正式 routing key 与代码一致。
  • [ ] 同一延迟事件重复到期仍幂等。
  • [ ] 补偿任务 limit 有上限并可小批运行。
  • [ ] 补偿失败能被调度系统识别,不只看退出码。
  • [ ] 人工重放只包含失败业务键。
  • [ ] 补偿后完成业务、库存、资金和下游四层对账。

22.4 关键专项

  • [ ] 支付成功、失败、异常、重复成功、迟到成功、退款分别验证。
  • [ ] 订单中心审核、关闭、售后所有权迁移分别验证。
  • [ ] 商品限购、套包、价格生效/失效、敏感词分别验证。
  • [ ] 配送出库、物流早到、自提、退货关闭分别验证。
  • [ ] 秒杀超时与支付回调并发验证。
  • [ ] 对账单生成、部分发送失败、重复发送分别验证。

23. 已确认与待环境确认

23.1 静态代码已确认

  • MqEventEnums 集中定义主要 destination 和 routing key。
  • MqSer 同时使用 Direct 和 Topic Producer。
  • 大多数生产者使用随机 uniqid() 作为消息 ID。
  • PayCenterNotify 注册 10 类支付中心事件。
  • 订单中心四个售后回调已迁移并在 DGJ 直接 ACK。
  • ItemCenterNotify 多个失败分支返回 ACK。
  • SaasOrderNotify::inventory_event 重算预警而非重做库存流水。
  • 新站 300 秒和暂停取消 60 秒延迟队列使用 TTL + dead-letter routing。
  • 对账单拆为策略、生成/收款状态、发送三类 destination。
  • 多个消费者提供基于临时文件的人工 resend()。

23.2 必须在环境确认

  • RabbitMQ exchange、queue、binding、durable、prefetch、消费者并发配置。
  • NACK 是重新入队、进入重试队列还是死信。
  • 自动重试次数、退避策略和死信告警。
  • 各消费者进程管理器、实例数、日志路径和告警规则。
  • dispatch_out_delay_20 的实际 TTL 和 DLX。
  • 消息生产是否启用 publisher confirm。
  • 消费 ACK 是否在 callback 返回后可靠执行。
  • 采购服务对已迁移订单中心事件的真实队列绑定。
  • 核心业务表是否有业务唯一键或事件处理记录。
  • 对账、商品中心和渠道 ACK 丢弃后的定时兜底任务。

24. 证据来源

主题代码文件
destination/routing keyapplication/KzData/Enums/MqEventEnums.php
DGJ 生产者application/Services/Mq/MqSer.php
销售同步报文application/Services/SyncOrder/SyncOrderSer.php、SyncOrderBuilder.php
支付消费application/controllers/tasks/PayCenterNotify.php
订单中心消费application/controllers/tasks/OrderCenterNotify.php
OA 消费application/controllers/tasks/OaNotify.php、OaResultNotify.php
商品中心消费application/controllers/tasks/ItemCenterNotify.php
配送消费application/controllers/tasks/DispatchCenterNotify.php
DGJ 内部消费application/controllers/tasks/SaasOrderNotify.php
渠道消费application/controllers/tasks/ChannelOrderNotify.php
对账消费者StatementNotify.php、StatementStrategyConsumer.php、StatementSendConsumer.php
秒杀补偿application/controllers/tasks/FlashSaleTask.php、application/Services/Activity/FlashSaleSer.php
库存预警SaasOrderNotify::inventory_event、库存实时/预警 Model

25. 一页式排查结论

先用 destination + routing key 找消费者,不要先猜业务表。
没有消费日志:查 publish、binding、进程和注册回调。
有消费日志:查业务键、事务提交和 ACK/NACK。
队列清空但业务没变:重点查主动 ACK、迁移空消费和部分成功。
重复消费:比较业务键,不比较随机 message ID。
主业务完成但派生数据缺失:只补后置 MQ/缓存/预警,不重做主单。
任何人工重放:先导出当前状态,只重放失败键,小批验证后扩大。

请求-日志-数据变更追踪卡

多入口请求链路

场景调用方与入口请求载荷/上下文Controller/ConsumerService/Provider汇合点最终业务事实
DGJ 主动生产订单/库存/财务/渠道 Servicedestination、routing key、业务快照领域 ServiceServices/Mq/MqSer.php、SyncOrderSer业务单号下游收到已提交的本地业务事实
中心回调消费ODC、支付、配送、商品、OA消息 ID、事件、来源单号、状态对应 tasks/*Notify.php领域 Service来源单号/外部业务号回调更新本地主表、明细或关系表
DGJ/渠道消费DEST_DGJ、DEST_DGJ_CHANNELinventory/po_close/channel action 等事件SaasOrderNotify、ChannelOrderNotify采购/库存/渠道 Service事件业务键手工或跨系统动作被执行
定时补偿crontab/调度平台task 方法、时间窗、批量限制tasks/*Task.php补偿 Service业务 ID + 状态 + 更新时间修复超时、漏回调、漏发消息

日志证据矩阵

| 链路段 | 日志来源 | 可检索锚点 | 成功信号 | 失败信号 | 与下一段关联方式 | | --- | --- | --- | --- | --- | --- | --- | | 生产 | MqSer/Provider 日志 | destination、routing key、业务单号、message ID | Broker 接受/生产成功 | 序列化失败、连接错误、发送异常 | message ID + 业务键查消费端 | | 消费入口 | *Notify 方法 | destination、event、消息 ID、来源单号 | 命中目标 handler | 未知事件、迁移空消费 | handler 名 + 业务键进入 Service | | 业务处理 | Consumer 调用的领域 Service | 业务单号、当前状态、幂等键 | commit 且状态/数量正确 | rollback、找不到关系、顺序冲突 | 单号回查主表、流水、关系表 | | ACK/NACK | Consumer 框架/Broker | message ID、重试次数 | ACK;重复消费零副作用 | NACK、反复重投、死信 | 重试记录与同一业务键对照 | | 补偿 | Task 日志 | task、批次、业务 ID、扫描/成功/失败数 | 重跑后失败数下降,已完成跳过 | 扫描无边界、重复扣加、任务中断 | 批次清单回查业务和消息事实 |

环节数据变更台账

步骤代码位置事务读取事实写入表/缓存/MQ字段或数量变化回查证据
构造消息MqSer、SyncOrderBuilder通常在业务事务尾部或 commit 后已提交业务对象消息 payload不应修改核心业务;生成稳定业务键payload 抽样、routing key、单号
Broker 投递MQ SDK/Provider本地 DB 外部边界payload、destinationBroker本地已成功,消息可能尚未消费生产确认、message ID
消费校验tasks/*Notify.php消费事务开始消息版本、事件、来源单号、当前状态无或消费记录无效/旧消息零业务写入handler 日志、当前状态
消费落库领域 Service单消息本地事务当前业务事实、幂等条件主表/明细/流水/关系表status/qty/amount: old -> new;相同业务键仅一次影响行数、业务键、流水条数
ACK/补偿Consumer/TaskACK 不属于 DB 事务commit 结果或失败记录ACK/NACK、重发 MQcommit 后 ACK;不确定状态先查事实再重放Broker 状态、重跑前后零重复副作用

子模块追踪:paycenter-mq 支付中心消息

环节入口/触发请求/业务键代码链路读取事实写入与字段变化日志证据异常与补偿
消费落账DEST_PAYCENTER_NOTIFY 支付结果事件message ID、payOrderNo、来源单、payStatusapplication/controllers/tasks/PayCenterNotify.phppayStatus=02、本地支付/来源单状态、已有明细回调事务首次写 Payment/Info/采购支付关系并推进状态;重复为 0 增量message ID + routing key + payOrderNo + sourceOrderNo非成功状态不按成功落账;迟到/重复先查外部最终态和本地流水
退款异常退款、异常、额度调整事件原支付单、退款/调整单、金额application/controllers/tasks/PayCenterNotify.php原支付成功额、已退/已调额、账户类型反向事务 refunded/adjusted: old -> old+n;ACK 在 commit 后message ID + 原支付/退款单号部分职责需环境合同确认;重放以原支付关系幂等

子模块追踪:ordercenter-mq 订单中心消息

环节入口/触发请求/业务键代码链路读取事实写入与字段变化日志证据异常与补偿
状态消费DEST_ODC_NOTIFY 审核、更新、关闭、取消来源单号、事件名、状态、message IDapplication/controllers/tasks/OrderCenterNotify.php来源映射、本地主单当前态、数量/资金事实单消息事务按映射 old -> new,必要时写关系/关闭事实message ID + event + source/local billNo未知、乱序状态零回退;关系缺失记录失败并按来源单补偿
后置通知本地 commit 后触发库存/下游同步本地单号、已提交状态application/Services/Mq/MqSer.php新状态和后置任务是否已有DB 事务外发布 MQ;本地主状态不因发布超时回滚business key + routing key + publish result主业务已成只补发消息,不重做审核/关闭

子模块追踪:dispatch-mq 调拨与配送中心消息

环节入口/触发请求/业务键代码链路读取事实写入与字段变化日志证据异常与补偿
配送消费出库、物流、签收、自提、退货关闭事件message ID、配送/业务单号、eventapplication/controllers/tasks/DispatchCenterNotify.php本地配送关系、出入库状态、已处理事件回调事务写配送/物流状态 old -> new,需要时推进来源单message ID + event + dispatch/business billNo重复/迟到不回退终态;找不到关系先回查创建请求
库存衔接配送出库或退货结果触发库存业务业务单、SKU、数量、仓位application/Services/Mq/MqSer.php本地过程单与库存流水是否存在库存领域事务 qty: old +/- n;ACK 与库存 commit 分界明确billNo + SKU + transType + message ID库存已变则只补关系/状态;库存未变才重放领域动作

子模块追踪:oa-mq OA 消息

环节入口/触发请求/业务键代码链路读取事实写入与字段变化日志证据异常与补偿
OA 回调DEST_OA_NOTIFY 审批结果templateCode、processNo、业务 ID、状态application/controllers/tasks/OaNotify.php模板码、活动/业务当前申请态、OA 关联回调事务按模板分发,审批 pending -> approved/rejectedmessage ID + templateCode + processNo + activity ID未知模板零业务写入;重复/乱序回调保持终态
结果回传本地动作向 OA 返回结果业务 ID、处理结果application/Services/Mq/MqSer.php本地 commit 结果DB 不变;commit 后发布 OA result MQbusiness ID + routing key + publish result发送失败只补 OA 结果,不重做本地审批处理

子模块追踪:itemcenter-mq 商品中心消息

环节入口/触发请求/业务键代码链路读取事实写入与字段变化日志证据异常与补偿
商品事件SKU 限制、套包状态、价格启停、敏感词事件message ID、event、SKU/invId、版本application/controllers/tasks/ItemCenterNotify.php商品本地映射、当前状态/版本、影响范围单消息事务更新商品派生字段或缓存标记 old -> newmessage ID + event + SKU/invId老版本或未知事件不得覆盖新状态;失败按 SKU 小批重放
派生刷新商品变更后刷新搜索/价格/活动视图SKU/invId、站点application/Services/Mq/MqSer.phpDB 商品事实和待刷新对象DB 主事实不变;异步缓存/索引追平SKU + refresh routing key + resultACK 只证明消费完成;页面仍旧时查缓存/索引更新时间

子模块追踪:dgj-mq DGJ 内部消息

环节入口/触发请求/业务键代码链路读取事实写入与字段变化日志证据异常与补偿
内部消费DEST_DGJ 的 inventory_event、po_close_confirm 等event、业务单、sid、SKU、message IDapplication/controllers/tasks/SaasOrderNotify.php当前业务/库存状态、来源关系和幂等事实单消息事务执行库存/关闭动作 old -> new,或已完成零写入message ID + event + billNo/SKU事件名错或状态不符不强推;按稳定业务键回查后重发
OPS 批处理手工批量事件或历史补偿批次、业务键集合application/controllers/tasks/SaasOrderNotify.php每条当前状态和已存在副作用每条独立事务;成功项不随失败项回滚batch ID + item key + success/fail count禁止无边界全量重跑;导出失败键后小批继续

子模块追踪:saas-mq DGJ 发往 App 与 SAAS 的消息

环节入口/触发请求/业务键代码链路读取事实写入与字段变化日志证据异常与补偿
构造发布销售、出库、库存等本地 commit 后业务单号、事件、状态快照application/Services/SyncOrder/SyncOrderBuilder.php -> application/Services/SyncOrder/SyncOrderSer.php已提交主明细和目标事件字段本地核心 DB 不变;发布 dgj_notify/SAAS MQbusiness billNo + event + routing key + message ID序列化/发布失败只补消息;避免重复做本地订单或库存
下游核对App/SAAS 消费结果或业务反馈来源单号、下游状态application/Services/Mq/MqSer.php本地事实和已发布记录查询只读 不写;必要时幂等补发同一业务快照两端业务单号 + publish/consume timestamp下游环境日志需验证;随机 messageId 不作为业务幂等键

子模块追踪:channel-mq 渠道消息

环节入口/触发请求/业务键代码链路读取事实写入与字段变化日志证据异常与补偿
渠道消费DEST_DGJ_CHANNEL 创建、发货、退款、审核动作event、渠道订单/售后单、message IDapplication/controllers/tasks/ChannelOrderNotify.php渠道主明细、来源映射、当前状态单消息事务 insert/update old -> new;重复来源不重复建单message ID + event + channelOrderNo/aftersaleNo未知事件/状态冲突按框架 ACK/NACK;生产策略需环境确认
渠道发送本地状态变更通知渠道本地/渠道单号、状态、数量application/Services/Mq/MqSer.php已提交业务快照和目标渠道commit 后发 MQ,DB 状态不因 publish 超时回滚local/channel billNo + routing key只补失败事件;先确认渠道是否已处理再重发

子模块追踪:statement-mq 对账单异步任务

环节入口/触发请求/业务键代码链路读取事实写入与字段变化日志证据异常与补偿
任务拆分对账生成/核对/导出队列statement ID、批次、站点、时间窗application/Services/Statement/StatementSer.php对账范围、当前任务状态、已生成明细任务事务 pending -> processing 并写分片任务;业务订单只读batch/statement ID + range + task count重复启动复用任务键;时间窗不完整时停止而非生成半张对账单
汇总完成子任务完成后汇总statement ID、成功/失败数application/Services/Mq/MqSer.php所有子任务状态和金额合计汇总事务 processing -> success/partial/failed;文件/通知异步statement ID + item counts + amount部分成功保留失败清单,只重跑失败分片并重新验算总额

子模块追踪:delay-compensate 延迟队列、重试与补偿任务

环节入口/触发请求/业务键代码链路读取事实写入与字段变化日志证据异常与补偿
延迟投递TTL + DLX 或定时扫描到期记录原业务键、到期时间、重试次数application/config/mq.php -> application/controllers/tasks/FlashSaleTask.php当前状态、更新时间、支付/采购最终事实到期前 DB 不变;补偿事务关闭/释放 old -> neworiginal business key + dueAt + retry/batch队列积压和 TTL 配置需环境确认;不得只凭消息时间执行反向动作
幂等补偿expireWaitPayOrders、compensatePaidOrders 等秒杀/采购单、扫描批次application/Services/Activity/FlashSaleSer.php支付、采购、locked/used 和已有补偿每单事务仅补缺失字段/数量;重跑影响行数 0task + batch + order IDs + affected rows先小批、先快照;中断保留失败键,禁止再次处理成功项