×

《订单同步的幂等设计:闲鱼逆向消息重复投递下,如何做到不出库两次》(附Python源码)

万邦科技Lex 万邦科技Lex 发表于2026-09-16 09:18:02 浏览21 评论0

抢沙发发表评论

《订单同步的幂等设计:闲鱼逆向消息重复投递下,如何做到不出库两次》(附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 不会被再调一次
三个保险:
  1. processed_msg.msg_id 主键 → 传输层去重

  2. order_state 状态机 → 同状态不回退、不重放

  3. 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 串成“订单域统一运行时”?


群贤毕至

访客