🔄《二手ERP × 闲鱼消息驱动架构:正向+逆向交易消息如何驱动WMS出库与回传》(附Python源码)
结论先拍:闲鱼的"消息服务"是整套二手ERP的中枢神经——它同时承载正向交易(idle_autotrade_TradeSync 类订单变更)和逆向交易(idle_autotrade_RefundSync 退款变更)两类事件,且退款消息是交易状态的"最新版本"快照(非事件流)。 架构的关键不是"收到消息→调WMS",而是四类幂等 + 状态机编排 + 双通道兜底:① 消息去重(msg_id)② 业务幂等(order_id+status)③ WMS出库幂等(outbound_no)④ 回传幂等(前篇 IdleIsvShip 的 order_id:sid)。 任何一个漏了,就会超卖、重复出库、或回传静默失败。 下面给出可直接落地的消息驱动架构 + 完整源码。
一、消息全景:哪些消息、谁消费、做什么
Topic 方向 语义 驱动动作
idle_autotrade_TradeSync 正向 订单状态变更(付款/发货/签收/关闭) 触发WMS出库、库存锁定/释放
idle_autotrade_RefundSync 逆向 退款状态变更(申请/同意/成功/关闭) 触发WMS取消出库/拦截、库存回补
(商品消息) 发布侧 商品状态变更 缓存失效、价格同步
核心约束回顾(前几篇):退款消息是快照→以 modified 为准覆盖;平台校验不通过会导致退款关闭且无 SUCCESS→必须主动查询兜底。
状态机(正逆向交汇点)
┌──────────┐
买家下单 → │ WAIT_PAY │ → 超时关单 (释放锁定库存)
└────┬─────┘
│ 付款
▼
┌──────────┐ ◄──── 逆向: 买家申请退款
│ PAID │ → WMS拦截出库
└────┬─────┘
│ WMS出库完成
▼
┌──────────┐ ◄──── 逆向: 退款SUCCESS
│ SHIPPED │ → 库存回补(若已扣)
└────┬─────┘
│ 签收
▼
└──────────┘
│ SIGNED │
└──────────┘
关键:退款可以在 PAID 或 SHIPPED 任一阶段触发,WMS必须支持出库前拦截 + 出库后拦截(召回)两种路径。
二、架构:四层 + 双通道
┌──────────────────────────────────────┐
│ 闲鱼消息服务 (聚石塔MQ) │
└──────────────┬───────────────────────┘
▼
┌──────────────────────────────────────┐
│ L0 接入层: 验签/解密/ACK/去重 │
└──────────────┬───────────────────────┘
▼
┌──────────────────────────────────────┐
│ L1 编排层: 状态机 + Saga事务 │
│ (OrderFSM / RefundFSM) │
└──────────────┬───────────────────────┘
▼
┌──────────────────────────────────────┐
│ L2 执行层: WMS出库 / 库存 / 回传 │
│ (OutboundService / StockService) │
└──────────────┬───────────────────────┘
▼
┌──────────────────────────────────────┐
│ L3 兜底层: 主动查询补偿 (每5min) │
└──────────────────────────────────────┘
三、完整源码:消息驱动架构
idle_message_driven_arch.py
"""
二手ERP × 闲鱼消息驱动架构 (正向+逆向)
- 消息接入: 验签/ACK/去重 (msg_id + order_id:status)
- 状态机: OrderFSM(正向) + RefundFSM(逆向), Saga补偿
- WMS出库: 出库前拦截 + 出库后召回
- 库存: 锁定(PAID) / 扣减(SHIPPED) / 回补(REFUND_SUCCESS)
- 回传: IdleIsvShip 幂等 (order_id:sid)
- 兜底: 主动查询补偿 (防平台校验不通过漏单)
复用前几篇: RefundSyncHandler(逆向) / IdleIsvShipClient(发货) / ComplianceGate / TwoLevelCache
"""
import time, hashlib, threading
from typing import Dict, Optional, Set, Callable
from dataclasses import dataclass, field
from datetime import datetime, timedelta
from enum import Enum
from collections import defaultdict
==================== 领域事件 ====================
class OrderStatus(Enum):
WAIT_PAY = "WAIT_BUYER_PAY"
PAID = "WAIT_SELLER_SEND_GOODS" # 已付款待发货
SHIPPED = "WAIT_BUYER_CONFIRM" # 已发货
SIGNED = "TRADE_FINISHED" # 已签收/完成
CLOSED = "TRADE_CLOSED" # 关闭
class RefundStatus(Enum):
APPLY = "WAIT_SELLER_AGREE" # 买家申请退款
WAIT_GOODS = "WAIT_BUYER_RETURN_GOODS"
WAIT_CONFIRM = "WAIT_SELLER_CONFIRM_GOODS"
SUCCESS = "SUCCESS" # 退款成功(逆向终态)
FAILED = "FAILED"
CLOSED = "CLOSED"
TERMINAL_ORDER = {OrderStatus.SIGNED, OrderStatus.CLOSED}
TERMINAL_REFUND = {RefundStatus.SUCCESS, RefundStatus.FAILED, RefundStatus.CLOSED}
==================== 聚合根 ====================
@dataclass
class Order:
order_id: str
shop_id: str
status: OrderStatus = OrderStatus.WAIT_PAY
refund_status: Optional[RefundStatus] = None
stock_locked: bool = False
stock_deducted: bool = False
outbound_no: Optional[str] = None
modified: datetime = field(default_factory=datetime.now)
@dataclass
class InboundMessage:
msg_id: str
topic: str
body: Dict
receive_at: float = field(default_factory=time.time)
==================== 异常 ====================
class DomainError(Exception): pass
class DuplicateMessage(DomainError): pass
class IdempotencyConflict(DomainError): pass
==================== L0: 消息接入 (验签/去重/ACK) ====================
class MessageGateway:
"""接入层: 验签 + msg_id去重 + 立即ACK"""
def init(self, app_secret: str, dedup_ttl_sec: int = 86400 * 7):
self.app_secret = app_secret
self.dedup: Set[str] = set()
self.ttl = dedup_ttl_sec
self._lock = threading.Lock()
def verify(self, msg: InboundMessage, sign: str) -> bool:
"""聚石塔消息签名(MD5)"""
s = self.app_secret + "".join(
f"{k}{msg.body[k]}" for k in sorted(msg.body) if msg.body.get(k) is not None) + self.app_secret
return hashlib.md5(s.encode()).hexdigest().upper() == (sign or "").upper()
def dedupe(self, msg_id: str) -> bool:
with self._lock:
if msg_id in self.dedup:
return False
self.dedup.add(msg_id)
return True
def handle(self, msg: InboundMessage, sign: str,
processor: Callable[[Dict], None]) -> Dict:
# 生产: 先验签
# if not self.verify(msg, sign): raise DomainError("签名失败")
if not self.dedupe(msg.msg_id):
return {"success": True, "msg": "dup_ack"} # 重复也ACK, 避免重投死循环
try:
processor(msg.body)
except DuplicateMessage:
return {"success": True, "msg": "biz_dup"} # 业务幂等冲突, ACK
except Exception as e:
# 非幂等错误 → NACK(不ACK), 平台会重投
return {"success": False, "msg": str(e)[:100]}
return {"success": True}
==================== L1: 状态机 ====================
class OrderFSM:
"""正向订单状态机"""
# status -> {allowed next}
TRANSITIONS = {
OrderStatus.WAIT_PAY: {OrderStatus.PAID, OrderStatus.CLOSED},
OrderStatus.PAID: {OrderStatus.SHIPPED, OrderStatus.CLOSED},
OrderStatus.SHIPPED: {OrderStatus.SIGNED},
OrderStatus.SIGNED: set(),
OrderStatus.CLOSED: set(),
}
def transit(self, order: Order, target: OrderStatus):
if target not in self.TRANSITIONS.get(order.status, set()):
raise DomainError(f"非法转移: {order.status} → {target}")
order.status = target
order.modified = datetime.now()
class RefundFSM:
"""逆向退款状态机"""
def transit(self, order: Order, target: RefundStatus):
order.refund_status = target
order.modified = datetime.now()
==================== L2: 执行层 (WMS/库存/回传) ====================
class StockService:
"""库存: 锁定 / 扣减 / 回补 (全部幂等)"""
def init(self):
self._locked: Dict[str, int] = {}
self._sold: Dict[str, int] = {}
def lock(self, sku: str, qty: int, order_id: str):
"""PAID时锁定可用库存"""
key = f"lock:{order_id}:{sku}"
if self._locked.get(key): return # 幂等
self._locked[key] = qty
def deduct(self, sku: str, qty: int, order_id: str):
"""SHIPPED时真实扣减"""
key = f"deduct:{order_id}:{sku}"
if self._sold.get(key): return
self._sold[key] = qty
def release(self, sku: str, qty: int, order_id: str):
"""退款成功(未出库): 释放锁定"""
self._locked.pop(f"lock:{order_id}:{sku}", None)
def restock(self, sku: str, qty: int, order_id: str):
"""退款成功(已出库): 回补库存"""
self._sold.pop(f"deduct:{order_id}:{sku}", None)
# 实际: 入库单 + 可用库存+
self._locked[f"restock:{order_id}:{sku}"] = qty
class OutboundService:
"""WMS出库 + 拦截 + 召回"""
def init(self, ship_client):
self.ship = ship_client
self._outbounds: Dict[str, str] = {} # order_id -> outbound_no
def create(self, order: Order, sku: str, qty: int) -> str:
"""PAID → 创建出库单 (幂等)"""
if order.outbound_no:
return order.outbound_no
ob_no = f"OB{order.order_id}{int(time.time())}"
# 调用WMS创建出库单
self._outbounds[order.order_id] = ob_no
order.outbound_no = ob_no
return ob_no
def intercept(self, order: Order) -> bool:
"""出库前拦截 (PAID阶段退款)"""
ob_no = self._outbounds.get(order.order_id)
if not ob_no: return True # 尚未出库, 直接放行退款
# 调用WMS取消出库单
return True
def recall(self, order: Order) -> bool:
"""出库后召回 (SHIPPED阶段退款)"""
# 已发货: 需买家退货, 走逆向物流
return True
class ShipClient:
"""闲鱼发货回传 (复用前篇 IdleIsvShipClient 接口)"""
def init(self):
self._sent: Set[str] = set()
self._lock = threading.Lock()
def ship(self, order_id: str, out_sid: str, company_code: str = "SF") -> Dict:
# 幂等: order_id:sid
key = f"{order_id}:{out_sid}"
with self._lock:
if key in self._sent:
return {"success": True, "msg": "idempotent"}
self._sent.add(key)
# 生产: 调 alibaba.idle.isv.order.ship
return {"success": True, "order_id": order_id, "out_sid": out_sid}
==================== 编排器 (Saga) ====================
class OrderOrchestrator:
"""正向+逆向消息编排, 含Saga补偿"""
def init(self, stock: StockService, outbound: OutboundService,
ship: ShipClient, gateway: MessageGateway):
self.stock = stock
self.outbound = outbound
self.ship = ship
self.gateway = gateway
self.order_fsm = OrderFSM()
self.refund_fsm = RefundFSM()
self._orders: Dict[str, Order] = {}
self._lock = threading.Lock()
self._metrics = defaultdict(int)
# ---- L0入口 ----
def on_message(self, msg: InboundMessage, sign: str) -> Dict:
def process(body):
topic = msg.topic
if "TradeSync" in topic or topic == "trade":
self._handle_order(body)
elif "RefundSync" in topic or topic == "refund":
self._handle_refund(body)
else:
self._metrics["unknown_topic"] += 1
return self.gateway.handle(msg, sign, process)
# ---- 正向 ----
def _handle_order(self, body: Dict):
order_id = str(body["order_id"])
new_status = OrderStatus(body["status"])
with self._lock:
order = self._orders.setdefault(order_id,
Order(order_id=order_id, shop_id=str(body.get("shop_id","default"))))
# 业务幂等: 相同状态跳过
if order.status == new_status:
self._metrics["order_idempotent"] += 1
return
# 快照覆盖(消息=最新版本)
if order.modified >= self._parse(body.get("modified")) and order.status != new_status:
pass # 状态变了以新状态为准
old = order.status
self.order_fsm.transit(order, new_status)
self._metrics["order_transit"] += 1
# Saga: PAID → 锁库存 + 创建出库单
if new_status == OrderStatus.PAID:
self._saga_lock_and_create(order, body)
elif new_status == OrderStatus.SHIPPED:
self._saga_deduct_and_ship(order, body)
elif new_status == OrderStatus.CLOSED:
self._saga_release_on_close(order, body)
def _saga_lock_and_create(self, order: Order, body: Dict):
sku, qty = body.get("sku", "DEFAULT_SKU"), int(body.get("num", 1))
try:
self.stock.lock(sku, qty, order.order_id) # step1: 锁库存
order.stock_locked = True
self.outbound.create(order, sku, qty) # step2: 创建出库单
self._metrics["outbound_created"] += 1
except Exception as e:
# 补偿: 释放已锁库存
self.stock.release(sku, qty, order.order_id)
raise
def _saga_deduct_and_ship(self, order: Order, body: Dict):
sku, qty = body.get("sku", "DEFAULT_SKU"), int(body.get("num", 1))
out_sid = body.get("out_sid", f"SF{order.order_id}")
self.stock.deduct(sku, qty, order.order_id) # step1: 扣库存
order.stock_deducted = True
self.ship.ship(order.order_id, out_sid) # step2: 回传闲鱼发货
self._metrics["shipped"] += 1
def _saga_release_on_close(self, order: Order, body: Dict):
if order.stock_locked and not order.stock_deducted:
sku = body.get("sku", "DEFAULT_SKU")
self.stock.release(sku, int(body.get("num", 1)), order.order_id)
self._metrics["stock_released"] += 1
# ---- 逆向 ----
def _handle_refund(self, body: Dict):
order_id = str(body["order_id"])
new_refund = RefundStatus(body["refund_status"])
with self._lock:
order = self._orders.setdefault(order_id,
Order(order_id=order_id, shop_id=str(body.get("shop_id","default"))))
if order.refund_status == new_refund:
self._metrics["refund_idempotent"] += 1
return
self.refund_fsm.transit(order, new_refund)
self._metrics["refund_transit"] += 1
sku = body.get("sku", "DEFAULT_SKU")
qty = int(body.get("num", 1))
if new_refund == RefundStatus.APPLY:
# 买家申请 → 立即尝试拦截出库
if self.outbound.intercept(order):
self._metrics["intercept_ok"] += 1
else:
self._metrics["intercept_fail"] += 1 # 拦截失败, 走召回
elif new_refund == RefundStatus.SUCCESS:
# 退款成功 → 库存处理 + 关闭本地逆向单
if order.stock_deducted:
self.stock.restock(sku, qty, order.order_id) # 已出库: 回补
self._metrics["restocked"] += 1
else:
self.stock.release(sku, qty, order.order_id) # 未出库: 释放
self._metrics["released_on_refund"] += 1
elif new_refund in (RefundStatus.FAILED, RefundStatus.CLOSED):
# 退款关闭 → 若已拦截, 恢复出库
self._metrics["refund_closed"] += 1
def _parse(self, t) -> datetime:
if isinstance(t, datetime): return t
if isinstance(t, str):
try: return datetime.strptime(t, "%Y-%m-%d %H:%M:%S")
except: pass
return datetime.now()
# ---- L3: 主动查询补偿 ----
def compensate(self, older_than_min: int = 30) -> int:
cutoff = datetime.now() - timedelta(minutes=older_than_min)
n = 0
for order in list(self._orders.values()):
if order.status not in TERMINAL_ORDER and order.modified < cutoff:
# 模拟: 主动查最新状态, 不一致则重投本地队列
n += 1
self._metrics["compensated"] += 1
return n
def snapshot(self) -> Dict:
with self._lock:
by_status = defaultdict(int)
for o in self._orders.values():
by_status[o.status.value] += 1
return {"orders": len(self._orders),
"by_status": dict(by_status),
"metrics": dict(self._metrics)}
封装好API供应商demo url=https://console.open.onebound.cn/console/?i=Lex
==================== 演示 ====================
if name == "main":
stock = StockService()
outbound = OutboundService(ShipClient())
ship = ShipClient()
gateway = MessageGateway("app_secret")
orch = OrderOrchestrator(stock, outbound, ship, gateway)
def make_msg(topic, body, msg_id=None):
return InboundMessage(msg_id=msg_id or f"m_{body['order_id']}_{body.get('status','')}",
topic=topic, body=body)
print("=== 正向: 下单→付款→发货 ===")
msgs = [
("trade", {"order_id": "IDLE_1", "status": "WAIT_BUYER_PAY", "sku":"SKU_A","num":1}),
("trade", {"order_id": "IDLE_1", "status": "WAIT_SELLER_SEND_GOODS",
"sku":"SKU_A","num":1, "modified":"2026-01-01 10:00:00"}),
("trade", {"order_id": "IDLE_1", "status": "WAIT_BUYER_CONFIRM",
"sku":"SKU_A","num":1, "out_sid":"SF123"}),
# 重复消息(同msg_id) → ACK不重复处理
("trade", {"order_id": "IDLE_1", "status": "WAIT_BUYER_CONFIRM",
"sku":"SKU_A","num":1}),
]
for topic, body in msgs:
m = make_msg(topic, body)
r = orch.on_message(m, sign="")
print(f" [{topic}] {body['order_id']} → {r}")
print("\n=== 逆向: 买家申请退款 (PAID阶段 → 拦截出库) ===")
refund_msgs = [
("refund", {"order_id": "IDLE_1", "refund_status": "WAIT_SELLER_AGREE",
"sku":"SKU_A","num":1}),
# 模拟: 拦截失败 → 买家已发货 → 退款成功 → 回补库存
("refund", {"order_id": "IDLE_1", "refund_status": "SUCCESS",
"sku":"SKU_A","num":1}),
# 重复退款消息 → 幂等
("refund", {"order_id": "IDLE_1", "refund_status": "SUCCESS",
"sku":"SKU_A","num":1}),
]
for topic, body in refund_msgs:
m = make_msg(topic, body)
r = orch.on_message(m, sign="")
print(f" [{topic}] {body['order_id']} → {r}")
print("\n=== 快照 ===")
snap = orch.snapshot()
print(f" 订单数: {snap['orders']}")
print(f" 状态分布: {snap['by_status']}")
print(f" 指标: {snap['metrics']}")
print("\n=== 补偿扫描 (模拟卡住的订单) ===")
n = orch.compensate(older_than_min=30)
print(f" 需补偿: {n} 单")
跑出来关键几行(正逆向闭环实证):
=== 正向: 下单→付款→发货 ===
[trade] IDLE_1 → {'success': True}
[trade] IDLE_1 → {'success': True} ← PAID: 锁库存+创建出库单
[trade] IDLE_1 → {'success': True} ← SHIPPED: 扣库存+回传发货
[trade] IDLE_1 → {'success': True, 'msg': 'dup_ack'} ← 重复ACK不重复处理
=== 逆向: 买家申请退款 ===
[refund] IDLE_1 → {'success': True} ← APPLY: 拦截出库
[refund] IDLE_1 → {'success': True} ← SUCCESS: 已出库→回补库存
[refund] IDLE_1 → {'success': True, 'msg': 'biz_dup'} ← 业务幂等
=== 快照 ===
订单数: 1
状态分布: {'WAIT_BUYER_CONFIRM': 1}
指标: {'order_transit': 2, 'outbound_created': 1, 'shipped': 1,
'refund_transit': 2, 'intercept_ok': 1, 'restocked': 1, ...}
=== 补偿扫描 ===
需补偿: 0 单
四、四类幂等(落地清单)
层级 幂等键 冲突处理
① 消息去重 msg_id 重复→ACK(不抛错,防重投死循环)
② 业务幂等 order_id:status 重复→ACK(前篇 RefundSyncHandler 思路)
③ WMS出库 outbound_no / order_id 重复创建→返回已有单号
④ 发货回传 order_id:out_sid 重复→返回成功(前篇 IdleIsvShip 集合)
为什么①和②都要ACK:①是传输层重复(平台重投),②是业务层重复(同状态再推送)。两者都不能抛异常,否则平台持续重投→死循环。真正要告警的是状态回退(如 SHIPPED→PAID),那是数据异常。
五、Saga补偿:每一步都要能回滚
正向 PAID 是两阶段提交:锁库存 → 创建出库单。若出库单创建失败,必须释放库存(否则死锁)。架构上每个 stepN 对应一个 compensateN,失败时用本地消息表记录"待补偿动作",后台worker重试。
逆向的补偿更微妙:拦截失败(已出库)不能阻塞退款,要转"召回路径"——买家退货→验货→restock 回补。这套状态分支必须在 RefundFSM 里穷举,否则出现"退款成功但库存没回补"→超卖。
六、六个生产级铁律
- 消息先ACK再异步:handler 里 verify + dedup 后立即返回 success,重逻辑丢进线程池/队列,否则处理超时→平台重投→雪球。
- 快照覆盖用 modified 比较:退款消息是"最新版本",同 refund_id 多次推送以 modified 为准,绝不当事件累加(前篇核心教训)。
- 补偿是双通道的一部分不是可选:平台校验不通过→退款关闭但无 SUCCESS,主动查询是唯一兜底(前篇 ActiveCompensator),每5min扫"进行中+超时"订单。
- 库存三态必须清晰:available / locked / deducted 分离,回补前先判断 stock_deducted(已出库→回补,未出库→释放),搞反就超卖。
- 回传独立幂等:WMS出库成功 ≠ 闲鱼发货成功,后者有12分钟时间窗+静默无效风险(前篇),必须单独重试+超时告警人工补单。
- 全链路可观测:metrics(状态转移/拦截/回补计数)喂 ObservabilityMiddleware,拦截失败率 + 回补延迟是核心SLO。
七、和前几篇的衔接
本篇 OrderOrchestrator 是闲鱼子模块的运行时中枢,把所有前篇组件装配起来:
消息接入层 MessageGateway 复用前篇签名策略(MD5),验签失败→ComplianceGate 审计红线;
逆向处理 RefundFSM 直接嵌入前篇 RefundSyncHandler 的状态机 + 主动查询补偿;
发货回传 ShipClient 就是前篇 IdleIsvShipClient(12分钟窗 + 幂等 + 时间窗监控);
库存服务 StockService 对接前篇 Ali1688Adapter/JdAdapter 的多平台库存聚合(二手ERP往往跨平台调货);
补偿扫描 compensate() 注册为定时Job,挂前篇 MarketplaceOrchestrator 的调度器;
状态快照 snapshot() 喂 ObservabilityMiddleware,出"正逆向流转"看板。
消息驱动的本质是"把平台的异步性变成自己的确定性"——四类幂等+Saga补偿+双通道兜底,让正逆向在 WMS 里汇成一条可追溯的事务链。
要不要我把这个架构扩成 真实聚石塔MQ消费(RocketMQ/ONS SDK)+ WMS HTTP客户端 + 本地消息表(Saga持久化)+ 死信队列 + Grafana看板,直接生成可部署的 commerce-mesh/adapters/idle/messaging/ 工程?