从0搭一个淘宝+京东+1688+拼多多+抖店五平台聚合中台,核心不是“把五个SDK凑一起”,而是用前几篇拆出的计费/限流/入塔规则倒推架构:订单能推不拉、按平台着色部署、统一DTO收口、令牌桶按Key隔离、配额/余额守卫编进Client。
下面给一套可直跑的轻量中台骨架(Python,单进程可启,生产换Redis/Celery/Kafka即可)。
一、中台分层(倒推出来的形态)
统一调度入口 Scheduler │ ├─ 平台Adapter层(每平台一个Client,签名/网关/Token刷新隔离) │ TaobaoAdapter(聚石塔内) JdAdapter Ali1688Adapter │ PddAdapter(云内+余额守卫) DyAdapter(云内) │ ├─ 统一模型层 DTO(StandardOrder / StandardSku / StandardStock) │ ├─ 限流守卫层(每AppKey独立令牌桶 + 日配额 + 拼多多余额熔断) │ ├─ 同步策略层(推送消费为主 + 增量modified兜底 + 失败死信补偿) │ └─ 存储层(PostgreSQL业务表 + Redis幂等/计数/令牌桶)
设计铁律(来自前五篇):
- 淘宝/抖店/拼多多必须云内,否则×10倍或禁调;
- 订单DSS/Webhook/订单同步服务为主,API增量仅兜底;
- 1688批发别硬轮询高级库存,爆款走高级包或Webhook;
- 京东联盟Key与商家Key物理隔离;
- 拼多多欠费硬切断,本地计数器兜底余额。
二、统一DTO(先把五家订单归一)
# dto.py
from dataclasses import dataclass, field
from enum import Enum
from datetime import datetime
class StdOrderStatus(str, Enum):
CREATED = "CREATED"
PAID = "PAID"
SHIPPED = "SHIPPED"
SIGNED = "SIGNED"
REFUNDING = "REFUNDING"
CLOSED = "CLOSED"
@dataclass
class StandardOrder:
channel: str # taobao/jd/ali1688/pdd/douyin
shop_id: str
order_id: str # 平台原始订单号
idempotency_key: str = "" # channel+order_id
status: StdOrderStatus = StdOrderStatus.CREATED
pay_amount: float = 0.0
post_fee: float = 0.0
item_count: int = 0
buyer_remark: str = ""
created_at: datetime = None
modified_at: datetime = None
raw: dict = field(default_factory=dict) # 原始报文留存溯源
def __post_init__(self):
if not self.idempotency_key:
self.idempotency_key = f"{self.channel}:{self.shop_id}:{self.order_id}"状态映射表(各Adapter转换时查这张表):
STATUS_MAP = {
"taobao": {"WAIT_BUYER_PAY":"CREATED","TRADE_PAID":"PAID",
"WAIT_SELLER_SEND_GOODS":"PAID","TRADE_BUYER_SIGNED":"SIGNED",
"TRADE_CLOSED":"CLOSED"},
"jd": {"10":"PAID","20":"PAID","30":"SHIPPED","40":"SIGNED","60":"CLOSED"},
"pdd": {"0":"CREATED","1":"PAID","2":"SHIPPED","3":"SIGNED","5":"REFUNDING"},
"douyin": {"1":"CREATED","2":"PAID","3":"SHIPPED","4":"SIGNED","5":"CLOSED"},
"ali1688": {"waitbuyerpay":"CREATED","waitsellersend":"PAID",
"waitbuyerreceive":"SHIPPED","confirm_send":"SIGNED","cancel":"CLOSED"},
}三、按Key隔离的令牌桶 + 配额守卫(核心)
# guard.py
import time, hashlib, json, requests
from datetime import datetime
from threading import Lock
class KeyRateGuard:
"""每个AppKey独立:令牌桶限速 + 日调用计数 + 拼多多余额熔断"""
def __init__(self, platform, app_key, qps, daily_free, in_cloud=True):
self.platform = platform
self.app_key = app_key
self.qps = qps
self.tokens = qps
self.ts = time.monotonic()
self.lk = Lock()
self.day = datetime.now().date()
self.today_calls = 0
self.daily_free = daily_free
self.in_cloud = in_cloud
self.pdd_balance = None # 拼多多外部注入
def _roll_day(self):
if datetime.now().date() != self.day:
with self.lk:
self.day = datetime.now().date()
self.today_calls = 0
def acquire(self, is_value=False):
self._roll_day()
# 1. 增值接口云外禁调
if is_value and not self.in_cloud and self.platform in ("taobao","pdd","douyin"):
朋 raise PermissionError(f"{self.platform} 增值接口必须云内")
# 2. 日免额80%预警,100%熔断非核心
if self.today_calls >= self.daily_free:
raise RuntimeError(f"{self.app_key} 日免额{self.daily_free}耗尽,停调防扣费")
elif self.today_calls == int(self.daily_free*0.8):
print(f"⚠️ {self.app_key} 达免额80%,切纯增量")
# 3. 拼多多余额守卫
if self.platform=="pdd" and self.pdd_balance is not None:
unit = 0.01/100 if self.in_cloud else 0.10/100
if self.pdd_balance <= (self.today_calls+1)*unit*3:
raise RuntimeError("pdd 余额<3天预估,熔断")
# 4. 令牌桶
with self.lk:
now = time.monotonic()
self.tokens = min(self.qps, self.tokens + (now-self.ts)*self.qps)
self.ts = now
if self.tokens < 1:
time.sleep((1-self.tokens)/self.qps + 0.005)
self.tokens = 0
else:
self.tokens -= 1
self.today_calls += 1四、五平台Adapter(统一接口,签名各异)
# adapters.py
from abc import ABC, abstractmethod
import hashlib, time, json, requests
from dto import StandardOrder, STATUS_MAP
class BaseAdapter(ABC):
def __init__(self, guard: KeyRateGuard, app_key, app_secret):
self.g = guard
self.ak = app_key
self.ask = app_secret
@abstractmethod
def pull_increment_orders(self, shop_id, token, start_mod, end_mod, page=1) -> list[StandardOrder]:
...
def _sign_top_like(self, params):
f = sorted((k,v) for k,v in params.items() if k!="sign" and v is not None and str(v)!="")
qs = "".join(f"{k}{v}" for k,v in f)
return hashlib.md5(f"{self.ask}{qs}{self.ask}".encode()).hexdigest().upper()
def _safe_req(self, url, params, is_value=False, max_retry=4):
self.g.acquire(is_value)
params["sign"] = self._sign_top_like(params)
for att in range(max_retry):
try:
r = requests.post(url, data=params, timeout=15)
d = r.json()
if "error_response" in d or "errorResponse" in d:
blob = json.dumps(d)
if any(k in blob for k in ("FLOW_CONTROL","limited-by","50001","no permission")):
time.sleep(min(2**att,8)); continue
raise Exception(blob)
return d
except requests.RequestException:
time.sleep(2**att); continue
raise RuntimeError("retry exhausted")
class TaobaoAdapter(BaseAdapter):
GW = "https://gw.api.taobao.com/router/rest"
def pull_increment_orders(self, shop_id, token, start_mod, end_mod, page=1):
biz = {"start_modified":start_mod,"end_modified":end_mod,
"page_no":page,"page_size":50,"fields":"tid,status,payment,post_fee,modified"}
p = {"method":"taobao.trades.sold.increment.get","app_key":self.ak,
"timestamp":str(int(time.time()*1000)),"format":"json","v":"2.0",
"sign_method":"md5","access_token":token}
p.update(biz)
d = self._safe_req(self.GW, p)
out=[]
for t in d.get("trades_sold_increment_get_response",{}).get("trades",{}).get("trade",[]):
out.append(StandardOrder(
channel="taobao", shop_id=shop_id, order_id=str(t["tid"]),
status=STATUS_MAP["taobao"].get(t["status"],"CREATED"),
pay_amount=float(t.get("payment",0)), post_fee=float(t.get("post_fee",0)),
modified_at=t.get("modified"), raw=t))
return out
class PddAdapter(BaseAdapter):
GW = "https://gw-api.pinduoduo.com/api/router"
def pull_increment_orders(self, shop_id, token, start_mod, end_mod, page=1):
p = {"client_id":self.ak,"method":"pdd.order.number.list.increment.get",
"timestamp":str(int(time.time())),"data_type":"JSON","v":"V1.0",
"start_updated_at":int(start_mod),"end_updated_at":int(end_mod),
"page":page,"page_size":50,"access_token":token}
d = self._safe_req(self.GW, p)
out=[]
for o in d.get("order_number_list_increment_get_response",{}).get("order_list",[]):
out.append(StandardOrder(
channel="pdd", shop_id=shop_id, order_id=o["order_sn"],
status=STATUS_MAP["pdd"].get(str(o["order_status"]),"CREATED"),
pay_amount=float(o.get("pay_amount",0)), modified_at=o.get("updated_at"), raw=o))
return out
# JdAdapter / Ali1688Adapter / DyAdapter 同构,略(方法名一致,签名换秒级/毫秒、method命名不同)生产里把TaobaoAdapter/PddAdapter/JdAdapter/Ali1688Adapter/DyAdapter都实现同一抽象,Scheduler不感知平台。
五、统一调度器(增量时间窗 + 多店轮转)
# scheduler.py
import time
from datetime import datetime, timedelta
from adapters import TaobaoAdapter, PddAdapter
from guard import KeyRateGuard
class ShopBinding:
def __init__(self, channel, shop_id, adapter, app_key, token,
qps, daily_free, in_cloud=True):
self.channel = channel
self.shop_id = shop_id
self.adapter = adapter
self.token = token
self.guard = KeyRateGuard(channel, app_key, qps, daily_free, in_cloud)
# 注册中心(实际从DB载)
SHOPS = [
ShopBinding("taobao","shopA",TaobaoAdapter(KeyRateGuard("taobao","AK_TB",8,80000), "AK_TB","AS_TB"),
"TB_TOKEN", 8, 80000, in_cloud=True),
ShopBinding("pdd","shopB",PddAdapter(KeyRateGuard("pdd","AK_PDD",8,50000), "AK_PDD","AS_PDD"),
"PDD_TOKEN", 8, 50000, in_cloud=True),
]
def sync_loop():
while True:
end = datetime.now()
start = end - timedelta(minutes=5) # 5分钟增量窗
for sb in SHOPS:
try:
orders = sb.adapter.pull_increment_orders(
sb.shop_id, sb.token,
start.strftime("%Y-%m-%d %H:%M:%S"),
end.strftime("%Y-%m-%d %H:%M:%S"))
for o in orders:
# 1. Redis幂等:key存在则跳
# 2. 写PG standard_order(upsert by idempotency_key)
# 3. 发Kafka事件 order.updated
print(f"✔ {o.channel}/{o.shop_id}/{o.order_id} -> {o.status}")
except (RuntimeError, PermissionError) as e:
print(f"⚠️ {sb.channel}/{sb.shop_id} 守卫拦截: {e}")
except Exception as e:
print(f"❌ {sb.channel}/{sb.shop_id} 异常: {e}")
time.sleep(60) # 主控节拍1分钟,内部增量5分钟窗
if __name__ == "__main__":
sync_loop()关键点:
- 主控1分钟心跳,拉取窗5分钟,重叠防漏(平台modified有秒级延迟);
- 每店独立Guard,店铺A限流不影响店铺B;
- 守卫抛错不进DB,只告警,避免把限流当业务异常处理。
六、推送为主的可插拔扩展点
上面是“增量轮询兜底”版,生产建议把各平台推送接进来:
- 淘宝:聚石塔DSS订单推送 → 消费RDS Binlog/推送服务,省API费;
- 拼多多:订单同步服务(多多云DB推送)替代
order.list.get; - 抖店/1688:消息订阅Webhook → MQ消费;
- 京东:宙斯能力中心数据推送(云鼎)。
调度器里加一个
PushConsumer 把消息转成 StandardOrder 走同一套幂等写,轮询只作“每30分钟全量校对”的补偿任务。七、从0到1落地顺序(避坑路径)
- 资质先行:按前文认证表,淘宝/抖店/拼多多订单必须企业自研应用,1688高级库存买包,京东商家JOS+联盟隔离;
- 部署着色:淘宝→聚石塔ECS,抖店→抖店云,拼多多→拼多多云,1688/京东→同主体阿里云/京东云VPC;
- 先接推送:每家开通订单推送/同步服务,写StandardOrder落库;
- 再补轮询:增量
modified每5分钟兜底,Guard卡80%免额; - 商品/库存后接:1688批发用高级包+Webhook,淘宝库存用
skus.quantity.update回写,别反向硬拉; - 监控面板:每AppKey日调用/剩余免额/拼多多余额/令牌桶等待长度 → 企微告警。
这套骨架把“五家收费模型”编译进了代码:云内强制校验、免额熔断、拼多多余额守卫、按Key令牌桶、统一DTO收口、增量重叠防漏。它不是最重的(无Kafka/Celery),但把多平台中台最易烂尾的“计费-限流-幂等”三件事在第一次启动时就焊死了。
要不要我接着把 PushConsumer(淘宝DSS/拼多多同步服务/抖店Webhook) 和 PostgreSQL upsert + Redis幂等键 的落地代码补完整,让这套中台从“轮询骨架”升级成“推拉一体可上大促”的版本?