《订单同步的幂等设计:闲鱼逆向消息重复投递下,如何做到不出库两次》(附Python源码)
先拍结论:
消息中间件/平台回调本质是 at-least-once delivery:网络超时、ACK 丢失、消费者崩溃,都会让同一条消息再来一次。“不出库两次”不是靠“别重试”,而是靠:消息去重键 + 业务状态机 + 出库指令唯一号 + 原子提交。闲鱼逆向消息(退款/售后)还是状态快照,不是事件计数——同一条refund_id推 5 次,也要以modified最新值覆盖,而不是“又退一次”。
一、重复从哪来(闲鱼场景)
1. 平台重投:idle_autotrade_TradeSync / RefundSync 网络抖动重发 2. 消费方崩溃:WMS 已扣库存,ACK 没发出去 → 重投 3. 主动查询兜底:轮询拉到的数据和消息重复 4. 多节点部署:两个 worker 同时拿到同一条消息
结果如果不防重:
PAID → 出库单A → 扣库存1 消息重投 PAID → 出库单B → 扣库存1 ❌ 一台 iPhone 发两台
二、四层幂等(前面系列已经埋过伏笔)
层级 | 幂等键 | 防什么 |
|---|---|---|
消息去重 | msg_id | 同一条传输消息处理两次 |
业务去重 | order_id + biz_type + status_version | 同一状态变更重复生效 |
出库幂等 | outbound_no / wms_outbound_id | WMS 生成两张出库单 |
回传幂等 | order_id:shipment_id | 运单号回传平台重复 |
只做msg_id去重是不够的:平台可能用不同msg_id推同一条业务状态。必须再叠一层业务幂等键。
三、统一订单状态机(正向+逆向交汇)
from enum import Enum
class OrderState(Enum):
PENDING = "PENDING" # 待付款
PAID = "PAID" # 已付款(可出库)
PICKING = "PICKING" # 拣货中
SHIPPED = "SHIPPED" # 已出库
SIGNED = "SIGNED" # 已签收
REFUNDING = "REFUNDING" # 退款中(未完结)
RETURNED = "RETURNED" # 退货入库
CANCELLED = "CANCELLED" # 关单/退款成功
# 允许的状态转移
TRANSITIONS = {
"PENDING": {"paid": "PAID", "cancel": "CANCELLED"},
"PAID": {"pick": "PICKING", "refund_apply": "REFUNDING", "cancel": "CANCELLED"},
"PICKING": {"ship": "SHIPPED", "refund_apply": "REFUNDING", "intercept": "PAID"},
"SHIPPED": {"sign": "SIGNED", "refund_success": "RETURNED"},
"REFUNDING": {"refund_success": "RETURNED", "refund_close": "PAID"},
"RETURNED": {},
"SIGNED": {"refund_success": "RETURNED"},
"CANCELLED": {},
}关键点:
PAID → PICKING → SHIPPED只能走一次逆向消息在
PAID/PICKING/SHIPPED都能来,但动作不同:出库前:拦截拣货
出库后:召回 + 退货入库 + 库存回补
四、核心:幂等消费器(SQLite/PG 可直落)
import sqlite3, time, json
from dataclasses import dataclass
from typing import Optional
@dataclass
class IdleMessage:
msg_id: str # 传输层ID(可能变)
topic: str # idle_autotrade_TradeSync / RefundSync
order_id: str
refund_id: str | None
biz_type: str # order_paid / order_shipped / refund_apply / refund_success
status: str
modified: int # 毫秒时间戳
payload: dict
class IdempotentOrderConsumer:
def __init__(self, db_path=":memory:"):
self.conn = sqlite3.connect(db_path, isolation_level=None)
self._migrate()
def _migrate(self):
c = self.conn.cursor()
c.executescript("""
CREATE TABLE IF NOT EXISTS processed_msg (
msg_id TEXT PRIMARY KEY,
order_id TEXT,
biz_type TEXT,
received_at INTEGER
);
CREATE TABLE IF NOT EXISTS order_state (
order_id TEXT PRIMARY KEY,
state TEXT NOT NULL,
state_version INTEGER NOT NULL,
last_modified INTEGER NOT NULL,
updated_at INTEGER NOT NULL
);
CREATE TABLE IF NOT EXISTS outbound_order (
outbound_no TEXT PRIMARY KEY,
order_id TEXT UNIQUE,
wms_status TEXT,
created_at INTEGER
);
CREATE TABLE IF NOT EXISTS inbound_return (
return_id TEXT PRIMARY KEY,
order_id TEXT,
processed INTEGER DEFAULT 0
);
""")
# -------- 1. 消息去重 --------
def is_duplicate_msg(self, msg: IdleMessage) -> bool:
row = self.conn.execute(
"SELECT 1 FROM processed_msg WHERE msg_id=?", (msg.msg_id,)
).fetchone()
return row is not None
# -------- 2. 业务幂等键 --------
def biz_key(self, msg: IdleMessage) -> str:
# refund 消息用 refund_id;订单消息用 order_id
if msg.refund_id:
return f"{msg.order_id}:refund:{msg.refund_id}:{msg.biz_type}"
return f"{msg.order_id}:order:{msg.biz_type}"
# -------- 3. 状态机推进(原子) --------
def apply_state(self, order_id: str, target: str, modified: int) -> str:
"""
返回:
'applied' 真转移
'ignored' 重复/非法
'stale' 老快照
"""
row = self.conn.execute(
"SELECT state, state_version, last_modified FROM order_state WHERE order_id=?",
(order_id,)
).fetchone()
cur_state = row[0] if row else "PENDING"
cur_ver = row[1] if row else 0
# 找事件类型
event = self._event_for_target(target)
next_state = TRANSITIONS.get(cur_state, {}).get(event)
if next_state is None:
return "ignored"
if modified < row[2] if row else False:
return "stale"
self.conn.execute(
"""INSERT INTO order_state(order_id,state,state_version,last_modified,updated_at)
VALUES(?,?,?,?,?)
ON CONFLICT(order_id) DO UPDATE SET
state=excluded.state,
state_version=excluded.state_version,
last_modified=excluded.last_modified,
updated_at=excluded.updated_at
""",
(order_id, next_state, cur_ver + 1, modified, int(time.time()))
)
return "applied"
def _event_for_target(self, target: str) -> str:
return {
"PAID":"paid","PICKING":"pick","SHIPPED":"ship","SIGNED":"sign",
"REFUNDING":"refund_apply","RETURNED":"refund_success",
"CANCELLED":"cancel","PAID_INTERCEPT":"intercept"
}.get(target, target.lower())
# 封装好API供应商demo url=https://console.open.onebound.cn/console/?i=Lex
# -------- 4. 出库幂等 --------
def create_outbound_if_needed(self, order_id: str) -> dict:
existing = self.conn.execute(
"SELECT outbound_no,wms_status FROM outbound_order WHERE order_id=?",
(order_id,)
).fetchone()
if existing:
return {"created": False, "outbound_no": existing[0], "wms_status": existing[1]}
outbound_no = f"OB-{order_id}-{int(time.time()*1000)}"
self.conn.execute(
"INSERT INTO outbound_order(outbound_no,order_id,wms_status,created_at) VALUES(?,?,?,?)",
(outbound_no, order_id, "CREATED", int(time.time()))
)
# >>> 这里才调用 WMS <<<
return {"created": True, "outbound_no": outbound_no, "wms_status": "CREATED"}
# -------- 5. 逆向:退款成功召回 --------
def handle_refund_success(self, msg: IdleMessage):
# 1) 状态机
res = self.apply_state(msg.order_id, "RETURNED", msg.modified)
if res != "applied":
return {"refund": res}
ob = self.conn.execute(
"SELECT outbound_no,wms_status FROM outbound_order WHERE order_id=?",
(msg.order_id,)
).fetchone()
if ob is None:
# 还没出库:仅释放预留库存
return {"refund": "applied", "wms": "inventory_released"}
outbound_no, wms_status = ob
if wms_status in ("CREATED", "PICKING"):
# 出库前拦截
self.conn.execute(
"UPDATE outbound_order SET wms_status='INTERCEPTED' WHERE outbound_no=?",
(outbound_no,)
)
return {"refund": "applied", "wms": "intercept_picking"}
elif wms_status in ("SHIPPED",):
# 出库后:买家寄回 → 入库 → 回补库存
self.conn.execute(
"UPDATE outbound_order SET wms_status='RECALL' WHERE outbound_no=?",
(outbound_no,)
)
self.conn.execute(
"INSERT OR IGNORE INTO inbound_return(return_id,order_id,processed) VALUES(?,?,0)",
(msg.refund_id, msg.order_id)
)
return {"refund": "applied", "wms": "recall_and_restock"}
return {"refund": "applied", "wms": "no_action"}
# 封装好API供应商demo url=https://console.open.onebound.cn/console/?i=Lex
# -------- 主入口 --------
def consume(self, msg: IdleMessage) -> dict:
# 消息级去重
if self.is_duplicate_msg(msg):
return {"level": "msg_dup", "action": "ack_only"}
# 业务级去重(退款快照:同 refund_id 多次推,只按 modified 覆盖)
# 用事务把「记消息 + 改状态 + 出库」绑一起
self.conn.execute("BEGIN")
try:
# 先落 msg_id(原子插入,并发时只有一个能插成功)
try:
self.conn.execute(
"INSERT INTO processed_msg(msg_id,order_id,biz_type,received_at) VALUES(?,?,?,?)",
(msg.msg_id, msg.order_id, msg.biz_type, int(time.time()))
)
except sqlite3.IntegrityError:
self.conn.execute("COMMIT")
return {"level": "msg_dup", "action": "ack_only"}
if msg.topic == "idle_autotrade_TradeSync":
if msg.biz_type == "order_paid":
st = self.apply_state(msg.order_id, "PAID", msg.modified)
ob = None
if st == "applied":
ob = self.create_outbound_if_needed(msg.order_id)
self.conn.execute("COMMIT")
return {"level": "biz", "state": st, "outbound": ob}
if msg.biz_type == "order_shipped":
st = self.apply_state(msg.order_id, "SHIPPED", msg.modified)
if st == "applied":
self.conn.execute(
"UPDATE outbound_order SET wms_status='SHIPPED' WHERE order_id=?",
(msg.order_id,)
)
self.conn.execute("COMMIT")
return {"level": "biz", "state": st}
if msg.topic == "idle_autotrade_RefundSync":
r = self.handle_refund_success(msg)
self.conn.execute("COMMIT")
return {"level": "biz", **r}
self.conn.execute("COMMIT")
return {"level": "noop"}
except Exception:
self.conn.execute("ROLLBACK")
raise五、为什么这样“绝对不会出库两次”
重复消息进来时:
第1次: msg_id 不存在 → 插 processed_msg 成功 → PAID 状态机推进 → 生成 OB-xxx(outbound_order.order_id UNIQUE) → COMMIT 第2次(同 msg_id): is_duplicate_msg = True → 直接 ack,不进业务 第2次(不同 msg_id,同 order_id+paid): msg_id 插成功 → apply_state(PAID) :当前已是 PAID → ignored → create_outbound_if_needed :order_id 已存在 → 返回旧 outbound_no → WMS 不会被再调一次
三个保险:
processed_msg.msg_id主键 → 传输层去重order_state状态机 → 同状态不回退、不重放outbound_order.order_id UNIQUE→ 一个订单只有一张出库单
六、闲鱼逆向消息的“快照覆盖”细节
退款消息不是:
refund_apply ++ refund_success ++
而是:
{ refund_id, refund_status, modified, refund_fee, reason }正确做法:
def upsert_refund_snapshot(self, msg: IdleMessage):
self.conn.execute("""
INSERT INTO refund_snapshot(refund_id,order_id,status,fee,reason,modified,updated_at)
VALUES(:refund_id,:order_id,:status,:fee,:reason,:modified,strftime('%s','now'))
ON CONFLICT(refund_id) DO UPDATE SET
status=CASE WHEN excluded.modified >= status_table.last_modified THEN excluded.status ELSE status END
""")简化版:
row = self.conn.execute( "SELECT modified FROM refund_snapshot WHERE refund_id=?", (msg.refund_id,) ).fetchone() if row is None or msg.modified >= row[0]: # 覆盖 else: # 老消息,丢弃
规则:以平台modified最大的那一条为准,不是“收到几条算几次”。
七、生产环境替换点
内存/SQLite | 生产替代 |
|---|---|
processed_msg | PG 表 + (msg_id) 主键 / Redis SET msg_id 1 EX 86400 |
order_state | 订单主表字段 state / state_version / last_modified |
BEGIN/COMMIT | 单 DB 事务;跨系统用 Outbox + CDC |
多节点并发 | INSERT msg_id 用 DB 唯一约束或 Redis SETNX 抢锁 |
主动查询兜底 | 定时拉 modified 窗口,走同一 consume() |
分布式场景推荐 Outbox 模式:
1. 本地事务:写业务表 + 写 outbox 表(含 msg_id/biz_key) 2. 后台线程把 outbox 发到 MQ 3. 消费端再按 biz_key 幂等
这样“写库”和“发消息”要么都成、要么都败,不会出现在“已出库但消息没记”的灰态。
八、验收用例(直接跑)
c = IdempotentOrderConsumer()
m1 = IdleMessage("m1","idle_autotrade_TradeSync","O123",None,
"order_paid","PAID",1700000000000,{})
m2 = IdleMessage("m2","idle_autotrade_TradeSync","O123",None,
"order_paid","PAID",1700000000001,{}) # 不同msg_id同业务
m3 = IdleMessage("m1","idle_autotrade_TradeSync","O123",None,
"order_paid","PAID",1700000000002,{}) # 同msg_id
print(c.consume(m1))
# {'level':'biz','state':'applied','outbound':{'created':True,...}}
print(c.consume(m2))
# {'level':'biz','state':'ignored','outbound':{'created':False,...}} ← 不出新出库单
print(c.consume(m3))
# {'level':'msg_dup','action':'ack_only'}
# 逆向:退款成功
rf = IdleMessage("r1","idle_autotrade_RefundSync","O123","RF1",
"refund_success","SUCCESS",1700000100000,{})
print(c.consume(rf))
# {'level':'biz','refund':'applied','wms':'recall_and_restock'}九、一句话收口
闲鱼订单同步的幂等 =msg_id 去重(传输层)
order_id/refund_id + modified(业务层)
状态机只允许向前转移(语义层)
outbound_no 唯一(WMS 层)
本地事务 / Outbox(原子层)
少一层,就会在“网络抖动”那天出库两次、发两台、财务对不上。
要不要我接着把这篇并进
commerce-mesh/core/:做成 IdempotentConsumer + OrderStateMachine + OutboxPublisher,并和前篇的 EventBus / GrayReleaseRouter / GradeIntegrityLoop 串成“订单域统一运行时”?