《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)
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
HTTP 层:立即 200
队列层:指数退避
限流(429 / `ISP_*`` / 超时):退避
业务非法(重复采购、状态回退):直接 DLQ,不重试
七、兜底:消息会丢,所以“推拉结合”
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:只去重不决策
八、可靠性检查表
[ ] 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 判断“已发货”,谁半夜被幽灵订单叫醒。