×

《从0搭多店铺聚合中台:淘宝+京东+1688+拼多多+抖店API统一调度实战》(附python源码)

万邦科技Lex 万邦科技Lex 发表于2026-07-27 09:31:17 浏览24 评论0

抢沙发发表评论

从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落地顺序(避坑路径)

  1. 资质先行:按前文认证表,淘宝/抖店/拼多多订单必须企业自研应用,1688高级库存买包,京东商家JOS+联盟隔离;

  2. 部署着色:淘宝→聚石塔ECS,抖店→抖店云,拼多多→拼多多云,1688/京东→同主体阿里云/京东云VPC;

  3. 先接推送:每家开通订单推送/同步服务,写StandardOrder落库;

  4. 再补轮询:增量modified每5分钟兜底,Guard卡80%免额;

  5. 商品/库存后接:1688批发用高级包+Webhook,淘宝库存用skus.quantity.update回写,别反向硬拉;

  6. 监控面板:每AppKey日调用/剩余免额/拼多多余额/令牌桶等待长度 → 企微告警。


这套骨架把“五家收费模型”编译进了代码:云内强制校验、免额熔断、拼多多余额守卫、按Key令牌桶、统一DTO收口、增量重叠防漏。它不是最重的(无Kafka/Celery),但把多平台中台最易烂尾的“计费-限流-幂等”三件事在第一次启动时就焊死了。
要不要我接着把 PushConsumer(淘宝DSS/拼多多同步服务/抖店Webhook)PostgreSQL upsert + Redis幂等键 的落地代码补完整,让这套中台从“轮询骨架”升级成“推拉一体可上大促”的版本?


群贤毕至

访客