×

《1688消息服务落地方案:alibaba.message.* 与 Webhook 回调的可靠性与重试策略》(附Python源码)

万邦科技Lex 万邦科技Lex 发表于2026-10-09 15:28:48 浏览33 评论0

抢沙发发表评论

《1688消息服务落地方案:alibaba.message.* 与 Webhook 回调的可靠性与重试策略》(附Python源码)

先拍结论:
1688 / 阿里系消息服务有两种形态:
① TMC 类消息通道(taobao.tmc.messages.consume / confirm,长连接拉消息,消费后确认,不确认会重发,7 天未确认可能被丢弃);
② HTTP Webhook / 回调 URL(订单/退款状态变更 POST 到你的公网地址,收到就回 200,超时/非 200 会重试)。
生产系统里:消息只是触发器,不是真相源;Webhook 必须“快返回、异步处理、幂等去重、失败进死信、5 分钟增量兜底”。
把 Webhook 当“权威数据源”+ 在 HTTP 线程里同步调 1688 下单 = 重复采购、幽灵订单、消息堆积三连。

一、消息模型先定对

1688 业务事件(订单支付 / 发货 / 退款创建 / 退款完成)
        │
        ▼
开放平台消息服务
   ├─ TMC/消息网关:consume → 业务处理 → confirm
   └─ HTTP Webhook:POST /webhook/1688 → 200 → 入队
        │
        ▼
本地消息表 / 队列
        │
        ▼
用 alibaba.trade.order.get / get.buyerView / refund.get 拉真相
        │
        ▼
更新本地状态机(正向单 / 逆向单解耦)
铁律:
  • msgId / eventId 只用于去重,不用于状态判定

  • 状态判定永远用 (order_id, status, status_version)

  • 消息丢了不可怕,怕的是“收到消息就信、就下单、就退款”


二、Webhook 接收端:快返回,别干重活

# webhook/receiver.py
import hashlib
import hmac
import json
import time

from fastapi import FastAPI, Request, Response

app = FastAPI()

WEBHOOK_SECRET = b"your_1688_webhook_secret"


def verify_top_sig(raw_body: bytes, sig: str) -> bool:
    # 1688/淘宝系常见 top-sig:HMAC-SHA256(raw_body, secret)
    expect = hmac.new(WEBHOOK_SECRET, raw_body, hashlib.sha256).hexdigest()
    return hmac.compare_digest(expect, sig or "")


@app.post("/webhook/1688")
async def handle_1688_webhook(request: Request):
    raw = await request.body()
    sig = request.headers.get("top-sig", "")

    # 1) 验签失败:401,不让平台以为成功
    if not verify_top_sig(raw, sig):
        return Response(content="invalid signature", status_code=401)

    # 2) 解析失败:返回 200 + 记死信,避免平台死重试
    try:
        payload = json.loads(raw)
    except json.JSONDecodeError:
        save_to_dead_letter(raw, reason="bad_json")
        return Response(content="accepted", status_code=200)

    # 3) 快返回:只做入队
    enqueue_inbound_message(
        msg_id=payload.get("msgId") or payload.get("messageId"),
        msg_type=payload.get("messageType") or payload.get("type"),
        order_id=str(payload.get("orderId") or payload.get("tradeId") or ""),
        body=payload,
        received_at=time.time(),
    )
    # 4) 永远先回 200;业务成败后面异步决定
    return Response(content="accepted", status_code=200)
要点:
  • 同步逻辑 ≤ 毫秒级:验签、解析、写 Inbound 表、入队列

  • 绝不在 Webhook 里调 fastCreateOrder / refund.agree

  • 平台超时(如 1000ms)或你抛异常 → 它会重试 → 你没幂等就炸


三、入站消息表:去重 + 触发器

# webhook/models.py
from dataclasses import dataclass
from enum import Enum


class InboundState(str, Enum):
    RECEIVED = "received"
    DEDUPED = "deduped"
    ENRICHING = "enriching"
    APPLIED = "applied"
    FAILED = "failed"
    DLQ = "dlq"


@dataclass
class InboundMessage:
    msg_id: str
    msg_type: str
    order_id: str
    body: dict
    state: InboundState = InboundState.RECEIVED
    retry_count: int = 0
建表唯一索引:
CREATE UNIQUE INDEX uq_inbound_msg ON inbound_message (msg_id);
CREATE INDEX idx_order_type ON inbound_message (order_id, msg_type);
消费侧伪代码:
def consume_inbound(msg: InboundMessage):
    # 1) msgId 去重
    if inbound_repo.exists(msg.msg_id):
        return  # 已处理过,直接丢

    # 2) 写去重行(唯一索引兜底并发)
    try:
        inbound_repo.insert(msg)
    except DuplicateKeyError:
        return

    # 3) 用消息里的 order_id 去拉真相
    order = Ali1688TradeClient().get_buyer_view(msg.order_id)
    if not order:
        raise RetryableError("order pull failed")   # 可重试

    # 4) 状态机用 (order_id, status, status_version)
    apply_order_snapshot(order)

四、幂等核心:别用 msgId 推进业务状态

# order/state_machine.py
from dataclasses import dataclass


@dataclass
class OrderSnapshot:
    order_id: str
    status: str
    status_version: int
    paid_amount: float
    refund_status: str | None
    raw: dict

# 封装好API供应商demo url=https://console.open.onebound.cn/console/?i=Lex
def apply_order_snapshot(snap: OrderSnapshot):
    cur = order_repo.get(snap.order_id)

    # 第一条:直接写
    if cur is None:
        order_repo.insert(snap)
        emit_domain_event(snap)
        return

    # 旧版本或同版本:忽略
    if snap.status_version < cur.status_version:
        return
    if snap.status_version == cur.status_version and snap.status == cur.status:
        return  # 幂等:同版本同状态,不重复发事件

    # 状态回退保护:REFUNDED / CLOSED 等终态不允许随意覆盖
    if cur.status in TERMINAL_STATES and snap.status not in TERMINAL_TRANSITIONS[cur.status]:
        log.warn("illegal rollback blocked", cur.status, snap.status)
        return

    order_repo.update(snap)
    emit_domain_event(snap)
1688 订单消息不保证有序,同 order_id 可能先收到 SHIPPED 再收到 PAID。
所以用 status_version / modified 时间做单调推进,不用“消息到达顺序”。

五、TMC / 消息网关模式(如果你走 consume/confirm)

阿里系 TMC 是:
  • taobao.tmc.messages.consume 拿消息

  • 处理完调 taobao.tmc.messages.confirm 确认

  • 不确认 → 平台择机重发

  • 7 天不确认 → 消息可能被清理

  • 消费失败不要无脑抛异常让消息狂重发,字段缺失/主键冲突不可重试,网络/权限问题才可重试

class TmcConsumer:
    def poll_once(self):
        msgs = tmc_client.messages_consume(group_name="erp_main")
        for m in msgs:
            try:
                self.handle(m)
                tmc_client.messages_confirm(s_message_ids=[m.id])
            except RetryableError:
                # 不 confirm,等平台重发;但别把所有异常都走这条路
                pass
            except FatalError as e:
                self.save_dlq(m, reason=str(e))
                tmc_client.messages_confirm(s_message_ids=[m.id])  # 确认掉,别死循环

    def handle(self, m):
        if not m.body.get("orderId"):
            raise FatalError("missing orderId")   # 不可重试
        if not pull_order_from_1688(m.body["orderId"]):
            raise RetryableError("1688 pull failed")  # 可重试
规则:
  • 可重试:网络超时、限流、1688 拉单失败

  • 不可重试:消息字段不全、业务主键冲突、签名错误、订单已终态却收到非法回退

  • 不可重试的别让它一直重发,进 DLQ


六、重试策略:指数退避 + 死信,不原地狂重试

# queue/consumer.py
import asyncio


BACKOFF = [1, 2, 4, 8, 16]   # 秒
MAX_RETRY = len(BACKOFF)


async def process_job(job):
    for attempt in range(MAX_RETRY):
        try:
            do_business(job)   # 调 1688 / 写单 / 发货
            mark_applied(job)
            return
        except RetryableError:
            if attempt == MAX_RETRY - 1:
                push_to_dlq(job, reason="retry_exhausted")
                return
            await asyncio.sleep(BACKOFF[attempt])
        except FatalError as e:
            push_to_dlq(job, reason=str(e))
            return
Webhook 场景里:
  • HTTP 层:立即 200

  • 队列层:指数退避

  • 限流(429 / `ISP_*`` / 超时):退避

  • 业务非法(重复采购、状态回退):直接 DLQ,不重试


七、兜底:消息会丢,所以“推拉结合”

哪怕 Webhook 覆盖率 99.7%,剩下 0.3% 会要命:
def reconciliation_loop():
    while True:
        window_start = now() - 300   # 最近 5 分钟
        for order in list_erp_orders(updated_within=300):
            remote = Ali1688TradeClient().get_buyer_view(order.order_id)
            if remote and remote.status_version > order.status_version:
                apply_order_snapshot(remote)
        sleep(300)
架构口诀:
Webhook:低延迟触发器
轮询:5 分钟增量兜底
trade.order.get:真相源
msgId:只去重不决策

八、可靠性检查表

Webhook 接收:
  • [ ] HTTPS 公网可达

  • [ ] 用原始 body 验签,不先 parse 再 serialize

  • [ ] 同步逻辑不调 1688 写接口

  • [ ] 收到就回 200,业务失败也别让平台无限冲

  • [ ] msgId 唯一索引去重

消息处理:
  • [ ] 消息 → 拉 trade.order.get / get.buyerView

  • [ ] 状态推进用 (order_id, status, status_version)

  • [ ] 正向单 / 逆向单两套状态机

  • [ ] REFUNDED / CLOSED 等终态有回退保护

  • [ ] 可重试 / 不可重试分开

补偿:
  • [ ] 5 分钟 modified 增量兜底

  • [ ] 死信队列有人看

  • [ ] 采购单/退款事件可重放

  • [ ] 对账任务每天跑一次“ERP vs 1688 实付/运费/退款”


九、和前几篇收口

  • 《1688 订单 API》:消息是触发器,trade.order.get 是真相

  • 《1688 采购单》:消息里说“已支付”≠ 已经成功下单,还要看 createOrder.preview + fastCreateOrder 结果

  • 《1688 跨境分销》:跨境单消息来了,先确认是否 boutiquefenxiao / 跨境宝已付,再决定下游动作

  • 《1688 物流 API》:logistics.trace.get 不靠消息推全量,消息只告诉你“有物流事件了,去拉”

  • 本文:消息系统本质 = “至少一次投递 + 你自己的幂等”,不是“平台保证只发生一次”


十、一句话收口

1688 消息服务的生产写法只有一句:
Webhook 快回 200,msgId 去重,order_id+status_version 推进状态,消息只触发拉单,拉到的才是真相,拉不到就 5 分钟兜底,重试有退避,非法进死信。
谁在回调里直接下单、直接退款、直接用 msgId 判断“已发货”,谁半夜被幽灵订单叫醒。


群贤毕至

访客