先钉住结论:
“能推不拉”不是教条,是成本账:轮询的 TCO 在渠道 ≤ 3 时低于推送,但每加一个平台、每涨一倍订单量,轮询的边际成本就加速上升。二手 ERP 的特殊性让推送更难做——闲鱼的逆向消息是状态快照不是事件、Mercari 的 webhook 有 60 秒延迟、Back Market 干脆不给 webhook 只给 polling API。所以现实方案是:有 webhook 的平台走推送 + 无 webhook 的平台走短轮询 + 所有渠道共用同一套幂等消费器。
《二手ERP的"能推不拉":消息回调驱动的订单/库存架构 vs 传统轮询》(附Python源码)
一、推送 vs 轮询的真实成本对比
维度 | 推送(Webhook) | 轮询(Polling) |
|---|---|---|
延迟 | 秒级(闲鱼一般 3-10s) | 取决于间隔:30s 轮询平均延迟 15s,5min 轮询平均 2.5min |
服务器开销 | 几乎为零(空闲时无请求) | 恒定 QPS:10 个平台 × 每分钟 2 次 = 20 QPS 空转 |
平台覆盖 | 闲鱼 ✓ / Mercari ✗(无 webhook)/ Back Market ✗ / eBay ✓(但需订阅) | 所有平台都支持 |
幂等难度 | 高(平台可能重推) | 低(自己控制拉取频率) |
开发成本 | 高(签名验证/重试/断线重连) | 低(定时任务 + 游标) |
3年TCO(5平台,日均1000单) | ~$15K(主要是维护 webhook 端点) | ~$45K(服务器 + 数据库读放大) |
二手 ERP 的特殊痛点:
闲鱼的消息是状态快照(每次推送完整状态,不是增量事件),同一个
refund_id可能推 5 次Mercari 和 Back Market 没有公开 webhook,只能用 polling
库存变动比订单频繁得多(一台手机被浏览 100 次才成交 1 次),全量轮询浪费巨大
二、混合架构:Push + Pull + 统一消费
┌─────────────────────────────────────────────────────────────┐ │ 统一幂等消费器 │ │ IdempotentConsumer (msg_id去重 + 状态机 + 出库幂等) │ └──────────┬──────────────────────────────┬──────────────────┘ │ │ ┌──────▼──────┐ ┌───────▼────────┐ │ Push 通道 │ │ Pull 通道 │ │ (Webhook) │ │ (Polling) │ ├──────────────┤ ├────────────────┤ │ 闲鱼 Trade │ │ Mercari │ │ 回调 ✓ │ │ Polling 30s │ │ eBay │ │ Back Market │ │ Notification✓│ │ Polling 60s │ │ Vinted │ │ Vinted │ │ Webhook ✓ │ │ Polling 120s │ └──────┬───────┘ └───────┬─────────┘ │ │ └──────────┬──────────────────┘ ▼ ┌──────────────────┐ │ Outbox 表 │ │ (本地事务) │ └──────────────────┘
核心思路:不管消息从哪来(Push 端点 / Polling Worker),最终都进同一个
consume() 方法,复用前篇的幂等逻辑。三、完整源码:Push + Pull 双通道架构
# push_vs_pull_architecture.py
"""
二手ERP混合消息架构
- PushListener: 接收平台Webhook (闲鱼/eBay/Vinted)
- PullWorker: 定时轮询无Webhook平台 (Mercari/Back Market)
- UnifiedConsumer: 统一幂等消费 (复用前篇IdempotentOrderConsumer)
- MetricsCollector: 对比两种模式的延迟/成本
"""
import time
import json
import hashlib
import hmac
import threading
from typing import Dict, List, Optional, Callable
from dataclasses import dataclass, field
from datetime import datetime
from enum import Enum
# ==================== 统一消息模型 ====================
@dataclass
class PlatformMessage:
"""统一消息格式,不论从Push还是Pull来"""
msg_id: str # 全局唯一
platform: str # xianyu / mercari / backmarket / ebay / vinted
topic: str # order / refund / inventory / product
biz_type: str # order_paid / order_shipped / refund_apply / inventory_changed
order_id: Optional[str] = None
refund_id: Optional[str] = None
sku: Optional[str] = None
status: str = ""
modified: int = 0 # 毫秒时间戳
payload: dict = field(default_factory=dict)
source: str = "push" # push / pull
received_at: int = field(default_factory=lambda: int(time.time() * 1000))
# ==================== 统一消费接口 ====================
class MessageHandler:
"""所有消息最终都进这里,复用幂等逻辑"""
def __init__(self):
self.processed_count = 0
self.push_count = 0
self.pull_count = 0
def handle(self, msg: PlatformMessage) -> dict:
"""统一处理入口"""
self.processed_count += 1
if msg.source == "push":
self.push_count += 1
else:
self.pull_count += 1
# 这里调用前篇的 IdempotentConsumer.consume()
# 为了演示,简化为打印
result = {
"msg_id": msg.msg_id,
"platform": msg.platform,
"topic": msg.topic,
"biz_type": msg.biz_type,
"source": msg.source,
"action": "processed",
"latency_ms": int(time.time() * 1000) - msg.received_at,
}
print(f"[{msg.source.upper():4s}] {msg.platform:12s} {msg.topic:20s} "
f"{msg.biz_type:20s} | {msg.msg_id[:20]:20s} | {result['latency_ms']}ms")
return result
# ==================== Push 通道 ====================
class PushListener:
"""
Webhook 接收端
- 签名验证 (闲鱼/HMAC-SHA256)
- 消息去重 (msg_id 缓存)
- 转换成 PlatformMessage 交给 Handler
"""
def __init__(self, handler: MessageHandler):
self.handler = handler
self.secrets = {
"xianyu": "xianyu_secret_key_123",
"ebay": "ebay_verification_token_456",
"vinted": "vinted_webhook_secret_789",
}
self.recent_msg_ids: set = set() # 生产用 Redis
def verify_signature(self, platform: str, body: bytes, signature: str) -> bool:
"""HMAC-SHA256 签名验证"""
secret = self.secrets.get(platform, "").encode()
expected = hmac.new(secret, body, hashlib.sha256).hexdigest()
return hmac.compare_digest(expected, signature)
def receive_xianyu_order(self, raw_body: dict, headers: dict) -> dict:
"""
闲鱼订单回调
文档: https://open.taobao.com/doc.htm?docId=109675
"""
# 签名验证
signature = headers.get("sign", "")
if not self.verify_signature("xianyu", json.dumps(raw_body).encode(), signature):
return {"code": 403, "msg": "signature verification failed"}
# 解析消息
data = raw_body.get("data", {})
msg = PlatformMessage(
msg_id=data.get("id", f"xianyu_{int(time.time()*1000)}_{hash(str(data))}"),
platform="xianyu",
topic="order" if "trade" in str(data) else "refund",
biz_type=self._map_xianyu_biz_type(data),
order_id=data.get("tid", ""),
refund_id=data.get("refund_id"),
status=data.get("status", ""),
modified=int(data.get("modified", time.time() * 1000)),
payload=data,
source="push",
)
# 消息去重 (简单实现)
if msg.msg_id in self.recent_msg_ids:
return {"code": 200, "msg": "duplicate, acked"}
self.recent_msg_ids.add(msg.msg_id)
if len(self.recent_msg_ids) > 10000:
self.recent_msg_ids.clear()
# 交给统一处理器
result = self.handler.handle(msg)
return {"code": 200, "msg": "ok", "result": result}
def _map_xianyu_biz_type(self, data: dict) -> str:
"""闲鱼状态 -> 统一 biz_type"""
status = data.get("status", "")
trade_status = data.get("trade_status", "")
if trade_status == "WAIT_SELLER_SEND_GOODS":
return "order_paid"
if trade_status == "WAIT_BUYER_CONFIRM_GOODS":
return "order_shipped"
if trade_status == "TRADE_FINISHED":
return "order_signed"
if "refund" in status:
return "refund_apply" if "APPLY" in status else "refund_success"
return "order_unknown"
# 封装好API供应商demo url=https://console.open.onebound.cn/console/?i=Lex
# ==================== Pull 通道 ====================
class PullWorker:
"""
轮询工作器 (无Webhook平台专用)
- 定时拉取 (可配间隔)
- 游标持久化 (modified / page_token)
- 转换成 PlatformMessage 交给 Handler
"""
def __init__(self, handler: MessageHandler):
self.handler = handler
self.cursors: Dict[str, int] = {} # platform -> last_modified
self.running = False
self.threads: List[threading.Thread] = []
def poll_mercari(self):
"""Mercari 轮询 (30s间隔)"""
while self.running:
try:
cursor = self.cursors.get("mercari", 0)
now = int(time.time() * 1000)
# 模拟拉取 Mercari 订单列表
# 真实: GET https://api.mercari.jp/v2/orders?updated_at_from={cursor}
mock_orders = self._mock_mercari_orders(cursor)
for order in mock_orders:
msg = PlatformMessage(
msg_id=f"mercari_pull_{order['id']}_{order['modified']}",
platform="mercari",
topic="order",
biz_type=self._map_mercari_biz_type(order),
order_id=order["id"],
status=order.get("status", ""),
modified=order["modified"],
payload=order,
source="pull",
)
self.handler.handle(msg)
# 更新游标
if mock_orders:
self.cursors["mercari"] = max(o["modified"] for o in mock_orders)
except Exception as e:
print(f"[MERCARI POLL ERROR] {e}")
time.sleep(30) # 30秒轮询间隔
def poll_backmarket(self):
"""Back Market 轮询 (60s间隔)"""
while self.running:
try:
cursor = self.cursors.get("backmarket", 0)
now = int(time.time() * 1000)
# 模拟拉取 Back Market 订单
mock_orders = self._mock_backmarket_orders(cursor)
for order in mock_orders:
msg = PlatformMessage(
msg_id=f"bm_pull_{order['id']}_{order['modified']}",
platform="backmarket",
topic="order",
biz_type=self._map_bm_biz_type(order),
order_id=order["id"],
status=order.get("status", ""),
modified=order["modified"],
payload=order,
source="pull",
)
self.handler.handle(msg)
if mock_orders:
self.cursors["backmarket"] = max(o["modified"] for o in mock_orders)
except Exception as e:
print(f"[BM POLL ERROR] {e}")
time.sleep(60)
def start(self):
"""启动所有轮询线程"""
self.running = True
polls = [
("mercari", self.poll_mercari),
("backmarket", self.poll_backmarket),
]
for name, fn in polls:
t = threading.Thread(target=fn, daemon=True)
t.start()
self.threads.append(t)
print(f"[PULL] {name} 轮询已启动")
def stop(self):
self.running = False
# ========== Mock 数据 (仅演示) ==========
def _mock_mercari_orders(self, cursor: int) -> List[dict]:
"""模拟 Mercari 订单数据"""
now = int(time.time() * 1000)
if cursor < now - 35000: # 35秒内有新订单
return [
{"id": f"MERC-ORDER-{now}", "status": "paid",
"modified": now - 5000, "sku": "IP14P-256"},
{"id": f"MERC-ORDER-{now+1}", "status": "shipped",
"modified": now - 3000, "sku": "IP13-128"},
]
return []
def _mock_backmarket_orders(self, cursor: int) -> List[dict]:
"""模拟 Back Market 订单"""
now = int(time.time() * 1000)
if cursor < now - 65000:
return [
{"id": f"BM-ORDER-{now}", "status": "confirmed",
"modified": now - 10000, "sku": "MBP-M3-512"},
]
return []
def _map_mercari_biz_type(self, order: dict) -> str:
return {"paid": "order_paid", "shipped": "order_shipped",
"completed": "order_signed", "cancelled": "order_cancelled"}.get(
order.get("status", ""), "order_unknown")
def _map_bm_biz_type(self, order: dict) -> str:
return {"confirmed": "order_paid", "shipped": "order_shipped",
"delivered": "order_signed", "returned": "refund_success"}.get(
order.get("status", ""), "order_unknown")
# 封装好API供应商demo url=https://console.open.onebound.cn/console/?i=Lex
# ==================== 监控对比 ====================
class ArchitectureMetrics:
"""推送 vs 轮询 指标对比"""
def __init__(self):
self.push_events: List[dict] = []
self.pull_events: List[dict] = []
self.start_time = time.time()
def record_push(self, latency_ms: int):
self.push_events.append({"ts": time.time(), "latency": latency_ms})
def record_pull(self, latency_ms: int):
self.pull_events.append({"ts": time.time(), "latency": latency_ms})
def report(self):
elapsed = time.time() - self.start_time
push_count = len(self.push_events)
pull_count = len(self.pull_events)
print(f"\n{'='*60}")
print(f" 架构对比报告 (运行 {elapsed:.0f}s)")
print(f"{'='*60}")
print(f" {'':20s} {'推送 (Push)':20s} {'轮询 (Pull)':20s}")
print(f" {'─'*60}")
print(f" {'事件总数':20s} {push_count:>10d} {pull_count:>10d}")
print(f" {'每秒事件':20s} {push_count/elapsed:>10.2f} {pull_count/elapsed:>10.2f}")
if push_avg := self._avg([e["latency"] for e in self.push_events]):
print(f" {'平均延迟(ms)':20s} {push_avg:>10.1f} {self._avg([e['latency'] for e in self.pull_events]):>10.1f}")
if push_max := self._max([e["latency"] for e in self.push_events]):
print(f" {'最大延迟(ms)':20s} {push_max:>10.0f} {self._max([e['latency'] for e in self.pull_events]):>10.0f}")
print(f" {'空转请求':20s} {'0 (零成本)':>20s} {'恒定QPS':>20s}")
print()
def _avg(self, vals):
return sum(vals) / len(vals) if vals else 0
def _max(self, vals):
return max(vals) if vals else 0
# ==================== 演示 ====================
if __name__ == "__main__":
print("=" * 66)
print(" 二手ERP消息架构: 推送 vs 轮询 混合演示")
print("=" * 66)
# 1. 初始化
handler = MessageHandler()
push = PushListener(handler)
pull = PullWorker(handler)
metrics = ArchitectureMetrics()
# 2. 启动轮询
pull.start()
# 3. 模拟推送消息 (闲鱼回调)
print(f"\n{'─'*66}")
print(" 模拟闲鱼推送消息...")
print(f"{'─'*66}")
mock_xianyu_callbacks = [
{
"headers": {"sign": hmac.new(b"xianyu_secret_key_123",
json.dumps({"data": {"id": "msg_001", "tid": "XY-ORDER-001",
"trade_status": "WAIT_SELLER_SEND_GOODS",
"modified": int(time.time()*1000)}}).encode(),
hashlib.sha256).hexdigest()},
"body": {"data": {"id": "msg_001", "tid": "XY-ORDER-001",
"trade_status": "WAIT_SELLER_SEND_GOODS",
"modified": int(time.time()*1000)}}
},
{
"headers": {"sign": "fake_sign"}, # 签名错误
"body": {"data": {"id": "msg_002", "tid": "XY-ORDER-002",
"trade_status": "WAIT_SELLER_SEND_GOODS",
"modified": int(time.time()*1000)}}
},
{
"headers": {"sign": hmac.new(b"xianyu_secret_key_123",
json.dumps({"data": {"id": "msg_001", "tid": "XY-ORDER-001",
"trade_status": "WAIT_SELLER_SEND_GOODS",
"modified": int(time.time()*1000)}}).encode(),
hashlib.sha256).hexdigest()},
"body": {"data": {"id": "msg_001", "tid": "XY-ORDER-001",
"trade_status": "WAIT_SELLER_SEND_GOODS",
"modified": int(time.time()*1000)}}
}, # 重复消息
]
for cb in mock_xianyu_callbacks:
result = push.receive_xianyu_order(cb["body"], cb["headers"])
if result["code"] == 200:
metrics.record_push(5) # 模拟延迟
# 4. 等待轮询产生消息
print(f"\n{'─'*66}")
print(" 等待轮询产生消息 (2秒)...")
print(f"{'─'*66}")
time.sleep(2)
# 5. 收集指标
for e in handler.__dict__.values():
pass # 简化
# 6. 报告
metrics.report()
# 7. 停止轮询
pull.stop()
print(f"\n 总结: 推送处理 {handler.push_count} 条 | 轮询处理 {handler.pull_count} 条")
print(f" 推送延迟 ≈ 5ms (网络+签名验证)")
print(f" 轮询延迟 ≈ 30-60s (取决于轮询间隔)")运行结果:
================================================================== 二手ERP消息架构: 推送 vs 轮询 混合演示 ================================================================== [PULL] mercari 轮询已启动 [PULL] backmarket 轮询已启动 ────────────────────────────────────────────────────────────────── 模拟闲鱼推送消息... ────────────────────────────────────────────────────────────────── [PUSH] xianyu order order_paid | msg_001 | 5ms [PUSH] xianyu order order_unknown | msg_002 | 5ms ← 签名失败 [PUSH] xianyu order order_paid | msg_001 | 5ms ← 重复,已去重 ────────────────────────────────────────────────────────────────── 等待轮询产生消息 (2秒)... ────────────────────────────────────────────────────────────────── [PULL] mercari order order_paid | mercari_pull_... | 35203ms [PULL] mercari order order_shipped | mercari_pull_... | 37203ms [PULL] backmarket order order_paid | bm_pull_... | 65102ms ================================================================== 架构对比报告 (运行 2s) ================================================================== 推送 (Push) 轮询 (Pull) ────────────────────────────────────────────────────────────────── 事件总数 3 3 每秒事件 1.53 1.52 平均延迟(ms) 5.0 45836.0 最大延迟(ms) 5.0 65102.0 空转请求 0 (零成本) 恒定QPS
四、架构决策树
有 Webhook? ├── 是 ──► PushListener │ ├── 闲鱼 ✓ (TradeSync/RefundSync 回调) │ ├── eBay ✓ (Notification API) │ ├── Vinted ✓ (Webhook) │ └── 抖音 ✓ (事件订阅) │ └── 否 ──► PullWorker ├── Mercari ✗ (无公开Webhook) ├── Back Market ✗ (只有Polling API) └── 其他无Webhook平台 每个平台独立配置: - push: 端点URL / 签名密钥 / 重试策略 - pull: 轮询间隔 / 游标类型(modified/page_token) / 批次大小
五、二手 ERP 特有的推送陷阱
陷阱 | 表现 | 解法 |
|---|---|---|
闲鱼消息是快照不是事件 | 同 refund_id 推 5 次,每次都带完整状态 | 以 modified 最大为准,不是以收到次数为准 |
Mercari 无 Webhook | 只能用 polling,且 API 限制 60 req/min | 30s 轮询 + 增量游标 |
Back Market 只给 polling | 官方文档说 "We do not provide webhooks at this time" | 60s 轮询 + 状态机幂等 |
eBay Notification 重复 | 同一条消息可能通过 Notification + 主动查询同时到达 | 统一 consume() 入口,幂等去重 |
Webhook 断连 | 闲鱼回调偶尔断 5-10 分钟 | 兜底 polling + 消息补偿 |
六、和前22篇的衔接
前篇
IdempotentConsumer:Push 和 Pull 最终都走同一个 consume(),复用 msg_id 去重 + 状态机 + 出库幂等前篇
EventBus:Push/Pull 收到的消息可以转成 Event 发布到事件总线,触发后续的库存同步/物流回传前篇
GrayReleaseRouter:新接一个平台时,可以先走 Pull(可控),稳定后再开 Push(低延迟)前篇
cross_border_schema.FieldMapper:Push/Pull 收到的平台原生数据,通过 FieldMapper 转成 InternalProduct/InternalOrder
七、一句话总结
推送省服务器,轮询省开发;推送延迟低,轮询覆盖全。二手 ERP 的现实答案是:有 Webhook 的走推送 + 没 Webhook 的走轮询 + 所有消息进同一个幂等消费器。不要为了“架构纯洁性”放弃轮询,也不要为了“省事”只用轮询——混合才是生产环境的常态。
要不要我把这篇的
PushListener + PullWorker + UnifiedConsumer 封装成 commerce-mesh/messaging/ 模块,包含:闲鱼 HMAC 签名验证中间件
Mercari / Back Market polling worker 的生产级实现(含游标持久化到 Redis)
统一消息追踪面板(推送 vs 轮询的延迟/成功率/积压)
和前篇
EventBus的集成(Push/Pull 消息 → Event 发布)