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 四个不能混用的标识
| 标识 | 示例 | 用途 | 是否天然幂等 |
|---|---|---|---|
| destination | paycenter_notify | 队列/主题逻辑命名空间 | 否 |
| routing key | paycenter_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 正常消费的七个证据
- 生产业务事实存在。
- 生产者方法确实执行到
publish()。 - destination 和 routing key 与消费者绑定一致。
- 消费进程正在运行。
- 回调日志能找到业务键。
- 数据库事务提交,业务状态和副作用一致。
- 回调返回 ACK,或重复消息按设计被安全 ACK。
只有第 7 项没有前六项,不能证明业务完成;只有业务表变化但没有 ACK,也可能在稍后重复消费。
4. destination 总表
| destination 常量 | 值 | 方向 | 主要业务 |
|---|---|---|---|
DEST_DGJ | dgj | DGJ 内部/回到 DGJ | SAAS 消息、库存事件、新站、采购关闭、延迟任务 |
DEST_ODC | ordercenter | DGJ -> 订单中心 | 采购订单创建 |
DEST_ODC_NOTIFY | ordercenter_notify | 订单中心 -> DGJ | 审核、售后、取消、自制订单 |
DEST_DISPATCHCENTER_NOTIFY | dispatchcenter_notify | 双向约定 | 配送出库、物流、自提、退货确认 |
DEST_PAYCENTER_NOTIFY | paycenter_notify | 支付中心 -> DGJ | 支付、退款、授信、提现、异常 |
DEST_ITEMCENTER_NOTIFY | itemcenter_notify | 商品中心 -> DGJ | 限购、套包、价格、敏感词 |
DEST_OA_NOTIFY | oa_notify | OA -> DGJ | 审批通知 |
DEST_OA_RESULT_NOTIFY | oa_result_notify | DGJ/OA 结果同步 | OA 结果转发 |
DEST_DGJ_NOTIFY | dgj_notify | DGJ -> App/SAAS/下游 | 销售、出库、活动、客户、机器人 |
DEST_DGJ_CHANNEL | dgj_channel | DGJ 渠道内部 | 人工开单、发货、退款审核、订单审核 |
DEST_DGJ_STATEMENT | dgj_statement | DGJ 对账任务 | 对账单生成和收款状态 |
DEST_DGJ_STATEMENT_STRATEGY | dgj_statement_strategy | 策略任务 | 对账策略同步 |
DEST_DGJ_STATEMENT_SEND | dgj_statement_send | 发送任务 | 对账单发送 |
DEST_QPAYCENTER_NOTIFY | qpaycenter_notify | 结算中心 -> DGJ | 全车件支付、退款、提现 |
DEST_IMCENTER_NOTIFY | imcenter_notify | IM -> DGJ | App 业务消息 |
DEST_UDESK | dgj_udesk | DGJ -> Udesk | 客户异步同步 |
5. 生产者模型
5.1 核心代码
| 文件 | 作用 |
|---|---|
application/Services/Mq/MqSer.php | 通用业务消息生产者集合 |
application/Services/SyncOrder/SyncOrderSer.php | 销售单/出库单同步编排 |
application/Services/SyncOrder/SyncOrderBuilder.php | 构造下游报文 |
| 各领域 Service | 决定发布时机和 payload |
application/config/mq.php | MQ 连接和系统配置;不可把值写入文档 |
5.2 Direct 与 Topic
| 类型 | 代码工厂 | 当前用途特点 |
|---|---|---|
| Direct | createDirectProducer | DGJ 内部任务、库存、对账、延迟事件等精确路由 |
| Topic | createTopicProducer | 外部通知、订单中心回调、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_result | paycenterPayResult | 采购/预订单支付结果 |
paycenter_balance_adjust | paycenterBalanceAdjust | 成本调整 |
paycenter_refund_account_withdraw_result | paycenterRefundWithdrawResult | 待退款账户提现结果 |
paycenter_baitiao_return_pay_result | paycenterBaitiaoReturnPayResult | 授信还款支付结果 |
paycenter_baitiao_expire | paycenterBaitiaoExpire | 授信到期提醒 |
paycenter_pay_exception | paycenterPayException | 支付抵扣等异常 |
paycenter_installment_expire | paycenterInstallmentExpire | 分期到期 |
paycenter_refund_result | paycenterRefundResult | 退款结果 |
paycenter_baitiao_offline_return | paycenterBaitiaoOfflineReturnPayResult | 线下授信还款/退款结果 |
paycenter_baitiao_service_overdue_change | paycenterBaitiaoServiceOverdueChange | 服务站逾期状态变化 |
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 历史代码不可达,并注明已迁移到采购服务:
orderCenterAfterSaleAuditorderCenterAfterSaleCloseorderCenterAfterSaleUpdateorderCenterAfterSaleConfirmWareHouse
flowchart LR
A["订单中心售后消息"] --> B["DGJ OrderCenterNotify"]
B --> C["立即 ACK"]
C --> D["DGJ 历史处理代码不可达"]
A --> E["dgj-purchase-service 应承担实际消费"]
排查这四类事件时:
- DGJ 返回 ACK 不能证明售后已处理。
- 必须查采购服务的队列绑定、消费日志和数据库。
- 若两个服务都绑定同一 Topic,要确认是各自独立队列还是竞争消费。
- 不要通过删除 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
历史回调中常见模式是:
- 开启采购 Model 事务。
- 更新采购主单、明细和关闭记录。
- 提交事务。
- 发送库存事件,重算库存预警。
如果第 4 步失败,采购事实可能已提交但预警未刷新。此时应补发库存事件或重算预警,不应重做整个售后业务。
8. 调拨/配送中心消息
8.1 入口和事件
| routing key | 回调 | 业务结果 |
|---|---|---|
dispatchcenter_delivery_stock_out | dispatchCenterDeliveryStockOut | 采购配送出库/状态推进 |
dispatchcenter_logistice_info_notify | dispatchCenterLogisticeInfoNotify | 物流信息同步 |
dispatchcenter_self_pick_up_notify | dispatchCenterSelfPickUpNotify | 自提通知 |
| 退货关闭事件 | 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_change | sensitiveWordChange | 敏感词数据/缓存 |
itemcenter_sku_limit_purchase_config | skuLimitChange | SKU 限购配置 |
itemcenter_package_status_sync | syncPackageList | 套包主明细状态 |
itemcenter_item_price_enable | itemPriceEnable | 结算价生效 |
itemcenter_item_price_disable | itemPriceDisable | 结算价失效 |
itemcenter_guide_price_limit_change | salesPriceLimitNotice | 销售限价 |
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_order | callback | App/网页站内消息 |
udesk_customer_create | udeskCustomerCreate | Udesk/IM 账户创建 |
dispatch_pause_cancel | dispatchPauseCancel | 延迟关闭暂停取消 |
change_qty | changeQty | 采购明细数量异步变化 |
inventory_event | inventory_event | 实时库存/在途变化后重算预警 |
po_close_order_confirm | po_close_confirm | 采购关闭确认 |
search_log_add | search_log_add | 搜索日志 |
new_station | new_station | 新服务站初始化 |
new_application_open | newApplicationOpen | 新站应用开通 |
dgj_po_order_delivery_stock_out | delivery_stock_out | 采购配送出库 |
| OPS 创建/关闭事件 | opsCreateOrder、opsCloseOnlineOrder | OPS 自制采购单 |
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 key | payload 关键键 |
|---|---|---|
sendSaOrderCloseToAPP | app_sa_order_close | order_no |
sendSaOrderCutOffToAPP | app_sa_order_cut_off | order_no |
sendApplyReturnFinishToAPP | app_apply_return_finish | order_id、order_no、return_type |
sendApplyReturnCheckedToAPP | app_apply_return_checked | order_id、order_no |
sendSaOrderOutToMq | app_sa_order_out | 调用方传入出库数据 |
sendSaOrderStatusToMq | app_sa_order_status | 配送状态数据 |
sendBindMoveStorageToAPP | app_bind_move_storage | 微仓绑定数据 |
sendActivityStartToSaas | activity_start | sids |
sendActivityEndToSaas | activity_end | sids |
新销售同步还会使用 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_create | channelOrderManualCreate | 创建平台/DGJ 订单关系 |
channel_order_delivery | channelOrderDelivery | 同步渠道发货 |
channel_refund_action | channelRefundAction | 售后入库审核动作 |
channel_order_approve | channelOrderApprove | 渠道订单审核 |
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 三条队列
| destination | routing key | 消费者 | 业务 |
|---|---|---|---|
dgj_statement_strategy | statement_strategy_sync | StatementStrategyConsumer | 策略同步 |
dgj_statement | statement_bill_create | StatementNotify | 对账单生成 |
dgj_statement | statement_bill_receive_status | StatementNotify | 收款状态同步 |
dgj_statement_send | statement_send | StatementSendConsumer | 对账单发送 |
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 key | TTL | 到期后 routing key | 用途 |
|---|---|---|---|
new_station_delay_300 | 300000 ms | new_station | 新站延迟初始化 |
dispatch_pause_cancel_delay_60 | 60000 ms | dispatch_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 消息重复消费
- 查生产端是否执行了两次
publish()。 - 查消费者第一次是否 NACK 或进程在 ACK 前退出。
- 查队列是否有重试/死信回投。
- 比较业务键,不要只比较随机 message ID。
- 查第一次事务是否提交。
- 查副作用是否有唯一引用。
- 只补未完成步骤。
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 文件读取消息并重新发布,属于人工运维工具,不是自动补偿保证。
执行前必须:
- 阅读具体
resend()的允许 routing key。 - 确认输入文件格式和每行消息体。
- 去除已成功业务键。
- 在非生产环境验证一条。
- 记录原始消息、执行人、时间和数量。
- 执行后逐条核对业务结果。
- 删除或归档临时文件,避免误重跑。
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() | 业务重发仍可能重复执行 |
| 四类订单售后消息直接 ACK | OrderCenterNotify 方法首行返回 | DGJ 队列清空但实际所有权在采购服务 |
| 商品中心异常也可能 ACK | 多个 catch/校验分支 | 自动重试缺失,缓存/价格可能陈旧 |
| OPS 批处理子项失败整体 ACK | opsCreateOrder 循环 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 key | application/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/Consumer | Service/Provider | 汇合点 | 最终业务事实 |
|---|---|---|---|---|---|---|
| DGJ 主动生产 | 订单/库存/财务/渠道 Service | destination、routing key、业务快照 | 领域 Service | Services/Mq/MqSer.php、SyncOrderSer | 业务单号 | 下游收到已提交的本地业务事实 |
| 中心回调消费 | ODC、支付、配送、商品、OA | 消息 ID、事件、来源单号、状态 | 对应 tasks/*Notify.php | 领域 Service | 来源单号/外部业务号 | 回调更新本地主表、明细或关系表 |
| DGJ/渠道消费 | DEST_DGJ、DEST_DGJ_CHANNEL | inventory/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、destination | Broker | 本地已成功,消息可能尚未消费 | 生产确认、message ID |
| 消费校验 | tasks/*Notify.php | 消费事务开始 | 消息版本、事件、来源单号、当前状态 | 无或消费记录 | 无效/旧消息零业务写入 | handler 日志、当前状态 |
| 消费落库 | 领域 Service | 单消息本地事务 | 当前业务事实、幂等条件 | 主表/明细/流水/关系表 | status/qty/amount: old -> new;相同业务键仅一次 | 影响行数、业务键、流水条数 |
| ACK/补偿 | Consumer/Task | ACK 不属于 DB 事务 | commit 结果或失败记录 | ACK/NACK、重发 MQ | commit 后 ACK;不确定状态先查事实再重放 | Broker 状态、重跑前后零重复副作用 |
子模块追踪:paycenter-mq 支付中心消息
| 环节 | 入口/触发 | 请求/业务键 | 代码链路 | 读取事实 | 写入与字段变化 | 日志证据 | 异常与补偿 |
|---|---|---|---|---|---|---|---|
| 消费落账 | DEST_PAYCENTER_NOTIFY 支付结果事件 | message ID、payOrderNo、来源单、payStatus | application/controllers/tasks/PayCenterNotify.php | payStatus=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 ID | application/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、配送/业务单号、event | application/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/rejected | message ID + templateCode + processNo + activity ID | 未知模板零业务写入;重复/乱序回调保持终态 |
| 结果回传 | 本地动作向 OA 返回结果 | 业务 ID、处理结果 | application/Services/Mq/MqSer.php | 本地 commit 结果 | DB 不变;commit 后发布 OA result MQ | business ID + routing key + publish result | 发送失败只补 OA 结果,不重做本地审批处理 |
子模块追踪:itemcenter-mq 商品中心消息
| 环节 | 入口/触发 | 请求/业务键 | 代码链路 | 读取事实 | 写入与字段变化 | 日志证据 | 异常与补偿 |
|---|---|---|---|---|---|---|---|
| 商品事件 | SKU 限制、套包状态、价格启停、敏感词事件 | message ID、event、SKU/invId、版本 | application/controllers/tasks/ItemCenterNotify.php | 商品本地映射、当前状态/版本、影响范围 | 单消息事务更新商品派生字段或缓存标记 old -> new | message ID + event + SKU/invId | 老版本或未知事件不得覆盖新状态;失败按 SKU 小批重放 |
| 派生刷新 | 商品变更后刷新搜索/价格/活动视图 | SKU/invId、站点 | application/Services/Mq/MqSer.php | DB 商品事实和待刷新对象 | DB 主事实不变;异步缓存/索引追平 | SKU + refresh routing key + result | ACK 只证明消费完成;页面仍旧时查缓存/索引更新时间 |
子模块追踪:dgj-mq DGJ 内部消息
| 环节 | 入口/触发 | 请求/业务键 | 代码链路 | 读取事实 | 写入与字段变化 | 日志证据 | 异常与补偿 |
|---|---|---|---|---|---|---|---|
| 内部消费 | DEST_DGJ 的 inventory_event、po_close_confirm 等 | event、业务单、sid、SKU、message ID | application/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 MQ | business billNo + event + routing key + message ID | 序列化/发布失败只补消息;避免重复做本地订单或库存 |
| 下游核对 | App/SAAS 消费结果或业务反馈 | 来源单号、下游状态 | application/Services/Mq/MqSer.php | 本地事实和已发布记录 | 查询只读 不写;必要时幂等补发同一业务快照 | 两端业务单号 + publish/consume timestamp | 下游环境日志需验证;随机 messageId 不作为业务幂等键 |
子模块追踪:channel-mq 渠道消息
| 环节 | 入口/触发 | 请求/业务键 | 代码链路 | 读取事实 | 写入与字段变化 | 日志证据 | 异常与补偿 |
|---|---|---|---|---|---|---|---|
| 渠道消费 | DEST_DGJ_CHANNEL 创建、发货、退款、审核动作 | event、渠道订单/售后单、message ID | application/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 -> new | original business key + dueAt + retry/batch | 队列积压和 TTL 配置需环境确认;不得只凭消息时间执行反向动作 |
| 幂等补偿 | expireWaitPayOrders、compensatePaidOrders 等 | 秒杀/采购单、扫描批次 | application/Services/Activity/FlashSaleSer.php | 支付、采购、locked/used 和已有补偿 | 每单事务仅补缺失字段/数量;重跑影响行数 0 | task + batch + order IDs + affected rows | 先小批、先快照;中断保留失败键,禁止再次处理成功项 |