从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幂等键 的落地代码补完整,让这套中台从“轮询骨架”升级成“推拉一体可上大促”的版本?