×

《跨境多平台中台:亚马逊+eBay+淘宝+抖店API统一调度实战》(附Python源码)

万邦科技Lex 万邦科技Lex 发表于2026-08-17 10:57:49 浏览10 评论0

抢沙发发表评论

结论先拍:四家平台基因完全不同——亚马逊SP-API(LWA+OAuth2+SigV4,2026.5起调用费$0但留"at this time"后手)、eBay(单网关+X-EBAY-C-MARKETPLACE-ID头,5000次/天/AppKey/API)、淘宝TOP(MD5签名+聚石塔内+0.02/百次)、抖店(HMAC-SHA256+抖店云内+0.018/百次+预充值)。 中台最优解不是"四个SDK拼起来",而是六边形架构:统一DTO + 每平台一个Adapter + 全局令牌桶/日配额守卫 + 推送优先/增量兜底 + 异常语义归一。 这一套跑通后,新增第5家平台只写1个Adapter,核心业务代码零改。


一、四家接入约束速查(2026.08现行)

维度
亚马逊SP-API
eBay
淘宝TOP
抖店
网关
na/eu/fe 三端点
api.ebay.com 单网关
gw.api.taobao.com聚石塔内
openapi-fxg.jinritemai.com抖店云内
认证
LWA client_credentials / auth_code + AWS SigV4
OAuth2(app/client_credentials)+ user token
AppKey+AppSecret MD5签名+SessionKey
AppKey+AppSecret HMAC-SHA256+AccessToken
计费
$0(原拟已撤)
$0(5000次/天/AppKey/API)
0.02/百次(塔内,免额内0)
0.018/百次(云内,预充值)
云强制
AWS自有(推荐同区)
聚石塔内,外调×10
抖店云内,外调×10或禁敏感
订单同步
Notification(ORDER_CHANGE)优先
Platform Notifications
DSS推送
消息订阅Webhook
限流模型
按操作leaky bucket(Orders 40/s突100)
日5000+15s短窗6000
按Key QPS+日免额
按Key QPS+预充值余额
PII红线
RDT+30天删
30天删
聚石塔御城河
云内解密

二、中台分层(六边形架构落地)

┌──────────────────────────────────────────────────────────┐
│ 业务层 OMS/WMS/商品域(只认 StandardOrder / StandardSku) │
└───────────────────────────┬──────────────────────────────┘
                           │ 业务方法 get_orders()/sync_stock()
┌───────────────────────────┴──────────────────────────────┐
│ UnifiedGateway 统一门面                                     │
│  · 异常归一 ECommerceRateLimitExceeded / AuthError         │
│  · 令牌桶+日配额守卫(每AppKey独立)                       │
│  · 重试退避(429/5xx)+ 死信队列                           │
└───────────────────────────┬──────────────────────────────┘
                           │
        ┌──────────┬──────────┬──────────┬──────────┐
        ▼          ▼          ▼          ▼          │
  SpApiAdapter  EbayAdapter TaobaoAdapter DouyinAdapter
  (LWA+SigV4)   (OAuth+头)  (MD5+聚石塔) (HMAC+抖店云)
        └──────────┴──────────┴──────────┴──────────┘
                           │
                  推送消费层(SQS/Webhook/DSS/MQ)
                  增量modified兜底(5~30min)
└──────────────────────────────────────────────────────────┘
        Redis:令牌桶/日计数/幂等键/RDT缓存
        PG:StandardOrder落库
设计铁律(来自前几篇倒推):
  1. 能推不拉:四家都开通知,轮询只作补偿

  2. 云内着色:淘宝→聚石塔、抖店→抖店云、亚马逊→同区域AWS、eBay随意但别暴拉

  3. 守卫编译进Client:令牌桶+日配额+拼多多/抖店余额熔断,发起前拦截

  4. 异常归一:各家isv.invalid-permission/InvalidInputException/429 全部收敛为统一异常


三、Python:FourPlatformMiddleware(可直跑骨架)

# four_platform_middleware.py
"""
跨境多平台中台:亚马逊 SP-API + eBay + 淘宝TOP + 抖店
统一调度骨架(单进程可启,生产换 Redis/Celery/Kafka)
- 统一 DTO:StandardOrder
- 四 Adapter 签名/认证隔离
- 全局令牌桶 + 日配额守卫
- 异常语义归一
"""
import time, hashlib, hmac, json, requests
from typing import Dict, List, Optional
from dataclasses import dataclass, field
from enum import Enum
from datetime import datetime, timedelta
from threading import Lock

# ===================== 统一异常 =====================
class ECommerceError(Exception): pass
class ECommerceAuthError(ECommerceError): pass
class ECommerceRateLimitExceeded(ECommerceError): pass
class ECommercePlatformDown(ECommerceError): pass

# ===================== 统一 DTO =====================
class StdStatus(str, Enum):
    CREATED="CREATED"; PAID="PAID"; SHIPPED="SHIPPED"
    SIGNED="SIGNED"; REFUNDING="REFUNDING"; CLOSED="CLOSED"

@dataclass
class StandardOrder:
    channel: str
    shop_id: str
    order_id: str
    status: StdStatus
    paid_amount: float = 0.0
    currency: str = ""
    recipient: dict = field(default_factory=dict)
    items: list = field(default_factory=list)
    raw: dict = field(default_factory=dict)
    pii_delete_after: Optional[str] = None

    @property
    def idempotency_key(self):
        return f"{self.channel}:{self.shop_id}:{self.order_id}"

# ===================== 令牌桶 =====================
class TokenBucket:
    def __init__(self, rate, burst):
        self.rate=rate; self.cap=burst; self.tokens=burst
        self.ts=time.monotonic(); self.lk=Lock()
    def acquire(self):
        with self.lk:
            now=time.monotonic()
            self.tokens=min(self.cap, self.tokens+(now-self.ts)*self.rate)
            self.ts=now
            if self.tokens<1:
                return (1-self.tokens)/self.rate+0.01
            self.tokens-=1
            return 0.0

# ===================== 基类 Adapter =====================
class BaseAdapter:
    CHANNEL="base"
    def __init__(self, shop_id, day_limit=5000):
        self.shop_id=shop_id
        self.bucket=TokenBucket(rate=day_limit/86400, burst=min(day_limit, 50))
        self.day_limit=day_limit
        self.day_used=0
        self.day_reset=self._next_local_midnight()
    def _next_local_midnight(self):
        now=datetime.now()
        return (now+timedelta(days=1)).replace(hour=0,minute=0,second=0,microsecond=0).timestamp()
    def _guard(self):
        if time.time()>=self.day_reset:
            self.day_used=0; self.day_reset=self._next_local_midnight()
        wait=self.bucket.acquire()
        if self.day_used>=self.day_limit:
            raise ECommerceRateLimitExceeded(f"{self.CHANNEL} 日配额{self.day_limit}耗尽")
        self.day_used+=1
        if wait>0: time.sleep(wait)
    def pull_orders(self, created_after: str) -> List[StandardOrder]:
        raise NotImplementedError
    def sync_stock(self, sku: str, qty: int) -> dict:
        raise NotImplementedError

# ===================== 亚马逊 SP-API =====================
class SpApiAdapter(BaseAdapter):
    CHANNEL="amazon"
    def __init__(self, shop_id, lwa_id, lwa_sec, refresh_token, region="NA"):
        super().__init__(shop_id, day_limit=2_500_000)  # Basic档警戒(当前$0也防429)
        self.lwa_id=lwa_id; self.lwa_sec=lwa_sec; self.rt=refresh_token
        self.region=region
        self.GW={"NA":"https://sellingpartnerapi-na.amazon.com",
                 "EU":"https://sellingpartnerapi-eu.amazon.com",
                 "FE":"https://sellingpartnerapi-fe.amazon.com"}[region]
        self._tok=None
    def _app_token(self):
        if self._tok and time.time()<self._tok[1]-300: return self._tok[0]
        r=requests.post("https://api.amazon.com/auth/o2/token",
            data={"grant_type":"refresh_token","refresh_token":self.rt,
                  "client_id":self.lwa_id,"client_secret":self.lwa_sec},timeout=10)
        d=r.json(); self._tok=(d["access_token"], time.time()+d["expires_in"])
        return self._tok[0]
    def pull_orders(self, created_after: str):
        self._guard()
        tok=self._app_token()
        # 演示头;生产加 AWS SigV4
        headers={"Authorization":f"Bearer {tok}","x-amz-access-token":tok,
                 "Content-Type":"application/json"}
        url=f"{self.GW}/orders/v0/orders"
        params={"MarketplaceIds":"ATVPDKIKX0DER","CreatedAfter":created_after}
        r=requests.get(url, params=params, headers=headers, timeout=15)
        if r.status_code==429: raise ECommerceRateLimitExceeded("SP-API 429")
        if r.status_code!=200: raise ECommercePlatformDown(f"SP-API {r.status_code}")
        payload=r.json().get("payload",{}).get("Orders",[])
        out=[]
        for o in payload:
            out.append(StandardOrder(
                channel="amazon", shop_id=self.shop_id,
                order_id=o.get("AmazonOrderId",""),
                status=StdStatus.PAID if o.get("OrderStatus")=="Unshipped" else StdStatus.CREATED,
                paid_amount=float(o.get("OrderTotal",{}).get("Amount",0) or 0),
                currency=o.get("OrderTotal",{}).get("CurrencyCode",""),
                pii_delete_after=(datetime.utcnow()+timedelta(days=30)).isoformat(),
                raw=o))
        return out

# ===================== eBay =====================
class EbayAdapter(BaseAdapter):
    CHANNEL="ebay"
    MP={"US":"EBAY_US","GB":"EBAY_GB","DE":"EBAY_DE","AU":"EBAY_AU"}
    def __init__(self, shop_id, cid, csec, site="US", user_refresh=None):
        super().__init__(shop_id, day_limit=5000)
        self.cid=cid; self.csec=csec; self.site=site; self.uref=user_refresh
        self._tok=None
    def _token(self):
        if self._tok and time.time()<self._tok[1]-300: return self._tok[0]
        if self.uref:
            r=requests.post("https://api.ebay.com/identity/v1/oauth2/token",
                data={"grant_type":"refresh_token","refresh_token":self.uref,
                      "scope":"https://api.ebay.com/oauth/api_scope/sell.fulfillment"},
                auth=(self.cid,self.csec),timeout=10)
            d=r.json(); self._tok=(d["access_token"],time.time()+d["expires_in"])
        else:
            r=requests.post("https://api.ebay.com/identity/v1/oauth2/token",
                data={"grant_type":"client_credentials",
                      "scope":"https://api.ebay.com/oauth/api_scope"},
                auth=(self.cid,self.csec),timeout=10)
            d=r.json(); self._tok=(d["access_token"],time.time()+d["expires_in"])
        return self._tok[0]
    def pull_orders(self, created_after: str):
        self._guard()
        tok=self._token()
        headers={"Authorization":f"Bearer {tok}",
                 "X-EBAY-C-MARKETPLACE-ID":self.MP[self.site],
                 "Content-Type":"application/json"}
        url="https://api.ebay.com/sell/fulfillment/v1/order"
        params={"create_date_from":created_after,"limit":50}
        r=requests.get(url,params=params,headers=headers,timeout=15)
        if r.status_code==429: raise ECommerceRateLimitExceeded("eBay 429/日配额")
        if r.status_code!=200: raise ECommercePlatformDown(f"eBay {r.status_code}")
        orders=r.json().get("orders",[])
        out=[]
        for o in orders:
            amt=0.0; cur=""
            for lg in o.get("paymentSummary",{}).get("payments",[]):
                amt+=float(lg.get("amount",{}).get("value",0) or 0)
                cur=lg.get("amount",{}).get("currency",cur)
            out.append(StandardOrder(
                channel="ebay", shop_id=self.shop_id,
                order_id=o.get("orderId",""),
                status=StdStatus.PAID if o.get("orderStatus")=="READY_FOR_DISPATCH" else StdStatus.CREATED,
                paid_amount=amt, currency=cur,
                pii_delete_after=(datetime.utcnow()+timedelta(days=30)).isoformat(),
                raw=o))
        return out

# ===================== 淘宝 TOP =====================
class TaobaoAdapter(BaseAdapter):
    CHANNEL="taobao"
    def __init__(self, shop_id, app_key, app_sec, session_key):
        super().__init__(shop_id, day_limit=80_000)  # 企业日免额近似
        self.app_key=app_key; self.app_sec=app_sec; self.session=session_key
    def _sign(self, params: dict):
        s="".join(f"{k}{params[k]}" for k in sorted(params))
        return hashlib.md5((self.app_sec+s+self.app_sec).encode()).hexdigest().upper()
    def pull_orders(self, created_after: str):
        self._guard()
        params={
            "method":"taobao.trades.sold.get",
            "app_key":self.app_key,"session":self.session,
            "timestamp":datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
            "format":"json","v":"2.0",
            "fields":"tid,status,payment,receiver_name,orders",
            "start_created":created_after,
        }
        params["sign"]=self._sign(params)
        r=requests.post("http://gw.api.taobao.com/router/rest",data=params,timeout=15)
        # 生产必须跑在聚石塔ECS内,外调×10
        d=r.json()
        if "error_response" in d:
            code=d["error_response"].get("sub_code","")
            if "invalid-permission" in code or "session" in code:
                raise ECommerceAuthError(f"淘宝 {code}")
            raise ECommercePlatformDown(f"淘宝 {d['error_response']}")
        out=[]
        for t in d.get("trades_sold_get_response",{}).get("trades",{}).get("trade",[]):
            out.append(StandardOrder(
                channel="taobao", shop_id=self.shop_id,
                order_id=str(t.get("tid","")),
                status=StdStatus.PAID if t.get("status")=="WAIT_SELLER_SEND_GOODS" else StdStatus.CREATED,
                paid_amount=float(t.get("payment",0) or 0),
                currency="CNY",
                recipient={"name":t.get("receiver_name")},
                pii_delete_after=(datetime.utcnow()+timedelta(days=30)).isoformat(),
                raw=t))
        return out

# ===================== 抖店 =====================
class DouyinAdapter(BaseAdapter):
    CHANNEL="douyin"
    def __init__(self, shop_id, app_key, app_sec, access_token):
        super().__init__(shop_id, day_limit=50_000)
        self.app_key=app_key; self.app_sec=app_sec; self.token=access_token
    def _sign(self, params: dict):
        s="".join(f"{k}={params[k]}" for k in sorted(params))
        return hmac.new(self.app_sec.encode(), s.encode(), hashlib.sha256).hexdigest()
    def pull_orders(self, created_after: str):
        self._guard()
        params={"app_key":self.app_key,"timestamp":int(time.time()),
                "order_status":"1,2,3","start_time":created_after,
                "access_token":self.token,"method":"order.listQuery"}
        params["sign"]=self._sign(params)
        # 生产必须抖店云内
        r=requests.post("https://openapi-fxg.jinritemai.com/order/listQuery",
                        json=params, timeout=15)
        d=r.json()
        if d.get("err_no",0)!=0:
            if "token" in d.get("message",""):
                raise ECommerceAuthError(f"抖店 {d['message']}")
            raise ECommercePlatformDown(f"抖店 {d}")
        out=[]
        for o in d.get("data",{}).get("order_list",[]) or []:
            out.append(StandardOrder(
                channel="douyin", shop_id=self.shop_id,
                order_id=str(o.get("order_id","")),
                status=StdStatus.PAID if o.get("order_status")=="1" else Std_:=StdStatus.CREATED,
                paid_amount=float(o.get("pay_amount",0) or 0)/100,
                currency="CNY",
                pii_delete_after=(datetime.utcnow()+timedelta(days=30)).isoformat(),
                raw=o))
        return out

# ===================== 统一调度门面 =====================
class CommerceMiddleware:
    def __init__(self):
        self.adapters: Dict[str, BaseAdapter] = {}
    def register(self, adapter: BaseAdapter):
        self.adapters[f"{adapter.CHANNEL}:{adapter.shop_id}"] = adapter
    def pull_all_orders(self, created_after: str) -> List[StandardOrder]:
        result=[]
        for key, ad in self.adapters.items():
            try:
                result.extend(ad.pull_orders(created_after))
            except ECommerceRateLimitExceeded as e:
                print(f"⏸ {key} 限流跳过: {e}")
            except ECommerceAuthError as e:
                print(f"🔑 {key} 授权失效需重刷: {k}")
            except ECommercePlatformDown as e:
                print(f"⚠️ {key} 平台异常: {e}")
        return result

# ===================== 演示 =====================
if __name__ == "__main__":
    mw = CommerceMiddleware()
    mw.register(SpApiAdapter("shop_amz_1", "LWA_ID", "LWA_SEC", "RT", "NA"))
    mw.register(EbayAdapter("shop_ebay_1", "CID", "CSEC", "US", "UREF"))
    mw.register(TaobaoAdapter("shop_tb_1", "APPKEY", "APPSEC", "SESSION"))
    mw.register(DouyinAdapter("shop_dy_1", "APPKEY", "APPSEC", "ACCESS_TOKEN"))
    orders = mw.pull_all_orders("2026-08-01T00:00:00Z")
    print(f"跨四平台拉到订单 {len(orders)} 条")
    for o in orders[:3]:
        print(f"  {o.channel:7} {o.shop_id} {o.order_id} {o.status.value} {o.paid_amount}{o.currency}")
跑通后:业务层只消费 StandardOrder,不知道背后是亚马逊还是抖店;新增快手/1688/拼多多=再加一个 *Adapter(BaseAdapter) 注册进 CommerceMiddleware,OMS代码零改。

四、四家调度避坑清单

  • 亚马逊:LWA refresh token最长18个月,多Worker用Redis集中存access token(7200s缓存),别每个容器各自刷新;RDT按orderId缓存60s;Notification(SQS/EventBridge)替代getOrders轮询,Basic 2.5M/月警戒线留着防429。

  • eBay:5000次/天是按AppKey+API名不是按站点,US+GB+DE同调Browse会共享配额;ReviseInventoryStatus额外受15s/6000短窗限制;OAuth应用级token用于Browse,用户级用于成交/发货。

  • 淘宝必须聚石塔内ECS,公网调0.02→0.20/百次×10;SessionKey会过期,定时刷新;DSS推送走MQ consumer,trades.sold.get只作5min补偿;fields只取需要的,别*全拉。

  • 抖店必须抖店云内,AccessToken 24h过期+refresh_token 30天;预充值余额<3天预估断非核心调用;2026.7起商品发布也收费,上新流程合并调用别循环单SKU发;Webhook消息订阅优先于order.listQuery


五、和前几篇的衔接

把本篇 CommerceMiddleware 与前篇 NinePlatformTCOMODE 开关并排:亚马逊栏 current=$0 / proposed敞口$7166/千卖家,淘宝/抖店栏走 in_cloud+预充值余额守卫,eBay栏走 5000/天桶,TCO测算器直接读各Adapter的 day_used 真实计数而不是拍脑袋。架构层统一、计费层双轨(国内按量+跨境$0留后手),这就是跨境中台在2026.08的稳态形态。
要不要我把上面骨架扩成 Redis Lua令牌桶(多容器共享)+ SQS/DSS/抖店Webhook统一消费进Kafka + StandardOrder落PG幂等表 + 每Adapter成本水位企微告警,直接合成你前八篇(国内5+亚马逊+eBay)的 commerce-mesh 单机可启版?


群贤毕至

访客