×

🏗️《二手ERP对接电商平台的总体方案:统一数据模型 + 事件驱动 + 灰度上线6原则》(附Python源码)

万邦科技Lex 万邦科技Lex 发表于2026-09-16 08:58:13 浏览21 评论0

抢沙发发表评论

收官篇终于来了。前面21篇散落在各平台的技术方案,今天用一个统一的架构把它们串起来——这不是又一个新方案,而是前面所有方案的组织方式

🏗️《二手ERP对接电商平台的总体方案:统一数据模型 + 事件驱动 + 灰度上线6原则》(附Python源码)

一、核心架构:三层两总线

┌─────────────────────────────────────────────────────────┐
│                    业务应用层                             │
│  商品管理  │  订单中心  │  库存中心  │  财务结算          │
└──────────────────────┬──────────────────────────────────┘
                       │ 事件总线 (Event Bus)
┌──────────────────────▼──────────────────────────────────┐
│                   适配器层                                │
│  ┌────────┐  ┌────────┐  ┌────────┐  ┌────────┐        │
│  │闲鱼     │  │Mercari  │  │Back    │  │eBay    │  ...   │
│  │Adapter  │  │Adapter  │  │Market  │  │Adapter │        │
│  │         │  │         │  │Adapter  │  │        │        │
│  └────────┘  └────────┘  └────────┘  └────────┘        │
└──────────────────────┬──────────────────────────────────┘
                       │ 数据总线 (Data Bus)
┌──────────────────────▼──────────────────────────────────┐
│                   统一数据层                               │
│  InternalProduct  │  InternalOrder  │  InternalInventory │
│  InternalShipment │  InternalPayment │  InternalReturn   │
└──────────────────────────────────────────────────────────┘
三层:统一数据层 → 适配器层 → 业务应用层
两总线:数据总线(字段映射/翻译) + 事件总线(异步解耦/灰度)

二、统一数据模型(数据总线核心)

前面22篇分散在各平台的数据模型,现在统一收口到6个核心实体:
# unified_data_model.py
"""
二手ERP统一数据模型
- 6个核心实体: Product / Order / Inventory / Shipment / Payment / Return
- 每个实体: Internal* 统一字段 + PlatformMeta 平台特定元数据
- 字段映射: 通过 FieldMapper 运行时翻译 (复用前篇 cross_border_schema)
"""
from typing import Dict, List, Optional, Any
from dataclasses import dataclass, field
from enum import Enum
from datetime import datetime


# ==================== 核心枚举 ====================
class InternalGrade(Enum):
    LIKE_NEW = "like_new"
    GOOD = "good"
    FAIR = "fair"
    PARTS = "parts"

class OrderStatus(Enum):
    PENDING = "pending"
    CONFIRMED = "confirmed"
    SHIPPED = "shipped"
    DELIVERED = "delivered"
    RETURNED = "returned"
    CANCELLED = "cancelled"

class InventoryAction(Enum):
    RESERVE = "reserve"
    RELEASE = "release"
    DEDUCT = "deduct"
    RESTOCK = "restock"


# ==================== 6个核心实体 ====================

@dataclass
class InternalProduct:
    """统一商品模型 (复用前篇 cross_border_schema.InternalProduct)"""
    sku: str
    title: str
    description: str
    price: float
    currency: str = "USD"
    grade: InternalGrade = InternalGrade.GOOD
    grade_detail: str = ""
    battery_health: int = 100
    has_original_box: bool = True
    accessories: List[str] = field(default_factory=list)
    powers_on: bool = True
    camera_works: bool = True
    wifi_works: bool = True
    repaired_parts: List[str] = field(default_factory=list)
    unlocked: bool = True
    screen_scratch: bool = False
    body_scratch: bool = False
    dents: bool = False
    photo_urls: List[str] = field(default_factory=list)
    category: str = ""
    brand: str = ""
    model: str = ""
    weight_kg: float = 0.0
    dimensions_cm: List[float] = field(default_factory=list)
    platform_meta: Dict[str, Any] = field(default_factory=dict)
    # 合规
    gpsr_responsible_entity: str = ""
    safety_declaration_url: str = ""
    hs_code: str = ""
    origin_country: str = ""


@dataclass
class InternalOrder:
    """统一订单模型"""
    order_id: str
    platform: str               # xianyu / mercari / backmarket / ebay
    platform_order_id: str
    sku: str
    quantity: int = 1
    price: float = 0.0
    currency: str = "USD"
    status: OrderStatus = OrderStatus.PENDING
    buyer_name: str = ""
    buyer_address: str = ""
    buyer_phone: str = ""
    shipping_method: str = ""
    shipping_cost: float = 0.0
    tax: float = 0.0
    total: float = 0.0
    created_at: datetime = field(default_factory=datetime.now)
    paid_at: Optional[datetime] = None
    shipped_at: Optional[datetime] = None
    delivered_at: Optional[datetime] = None
    platform_meta: Dict[str, Any] = field(default_factory=dict)


@dataclass
class InternalInventory:
    """统一库存模型"""
    sku: str
    warehouse: str = "default"
    total: int = 0
    reserved: int = 0
    available: int = 0
    damaged: int = 0
    last_action: InventoryAction = InventoryAction.RESTOCK
    last_qty: int = 0
    updated_at: datetime = field(default_factory=datetime.now)
    platform_meta: Dict[str, Any] = field(default_factory=dict)


@dataclass
class InternalShipment:
    """统一发货模型"""
    shipment_id: str
    order_id: str
    platform: str
    tracking_number: str = ""
    carrier: str = ""
    method: str = ""
    weight_kg: float = 0.0
    length_cm: float = 0.0
    width_cm: float = 0.0
    height_cm: float = 0.0
    declared_value: float = 0.0
    declared_currency: str = "USD"
    origin_country: str = ""
    destination_country: str = ""
    shipped_at: Optional[datetime] = None
    estimated_delivery: Optional[datetime] = None
    platform_meta: Dict[str, Any] = field(default_factory=dict)


@dataclass
class InternalPayment:
    """统一支付模型"""
    payment_id: str
    order_id: str
    platform: str
    amount: float
    currency: str = "USD"
    fee: float = 0.0
    net_amount: float = 0.0
    method: str = ""            # credit_card / paypal / bank_transfer
    status: str = "completed"   # pending / completed / refunded
    paid_at: Optional[datetime] = None
    platform_meta: Dict[str, Any] = field(default_factory=dict)


@dataclass
class InternalReturn:
    """统一退货模型"""
    return_id: str
    order_id: str
    platform: str
    sku: str
    reason: str = ""
    reason_code: str = ""       # not_as_described / defective / wrong_item
    condition_at_return: str = ""
    refund_amount: float = 0.0
    refund_currency: str = "USD"
    return_shipping_cost: float = 0.0
    returned_at: Optional[datetime] = None
    refunded_at: Optional[datetime] = None
    restocked: bool = False
    grade_downgraded: Optional[InternalGrade] = None
    platform_meta: Dict[str, Any] = field(default_factory=dict)

三、事件驱动架构(事件总线核心)

# event_bus.py
"""
事件驱动架构
- Event: 统一事件格式 (type + payload + metadata)
- EventBus: 发布/订阅 + 异步处理 + 重试/死信
- GrayReleaseRouter: 灰度路由 (按平台/SKU/比例)
- 事件类型: product.* / order.* / inventory.* / return.*
"""
import json
import time
import random
from typing import Callable, Dict, List, Optional, Any
from dataclasses import dataclass, field
from datetime import datetime


# ==================== 事件定义 ====================
@dataclass
class Event:
    event_id: str
    type: str                   # product.created / order.confirmed / inventory.changed
    source: str                 # service name
    timestamp: float = field(default_factory=time.time)
    payload: Dict[str, Any] = field(default_factory=dict)
    metadata: Dict[str, Any] = field(default_factory=dict)

    def to_json(self) -> str:
        return json.dumps({
            "event_id": self.event_id,
            "type": self.type,
            "source": self.source,
            "timestamp": self.timestamp,
            "payload": self.payload,
            "metadata": self.metadata,
        })

    @staticmethod
    def from_json(data: str) -> "Event":
        d = json.loads(data)
        return Event(**d)


# ==================== 事件处理器 ====================
EventHandler = Callable[[Event], None]

class Subscription:
    def __init__(self, event_type: str, handler: EventHandler,
                 max_retries: int = 3, timeout_ms: int = 5000):
        self.event_type = event_type
        self.handler = handler
        self.max_retries = max_retries
        self.timeout_ms = timeout_ms


# ==================== 事件总线 ====================
class EventBus:
    """
    内存事件总线 (生产环境替换为 RabbitMQ/Kafka/RocketMQ)
    - 按 event_type 分发
    - 重试 + 死信队列
    - 灰度过滤
    """

    def __init__(self):
        self.subscriptions: Dict[str, List[Subscription]] = {}
        self.dead_letter_queue: List[Event] = []
        self.gray_router: Optional["GrayReleaseRouter"] = None

    def subscribe(self, subscription: Subscription):
        if subscription.event_type not in self.subscriptions:
            self.subscriptions[subscription.event_type] = []
        self.subscriptions[subscription.event_type].append(subscription)

    def publish(self, event: Event):
        handlers = self.subscriptions.get(event.type, [])
        if not handlers:
            return

        # 灰度过滤
        if self.gray_router and not self.gray_router.should_process(event):
            return

        for sub in handlers:
            self._dispatch_with_retry(sub, event)

    def _dispatch_with_retry(self, sub: Subscription, event: Event):
        for attempt in range(sub.max_retries):
            try:
                sub.handler(event)
                return
            except Exception as e:
                if attempt < sub.max_retries - 1:
                    time.sleep(2 ** attempt)  # 指数退避
                else:
                    self.dead_letter_queue.append(event)
                    print(f"[DEAD LETTER] {event.type} | {event.event_id} | {e}")

    def set_gray_router(self, router: "GrayReleaseRouter"):
        self.gray_router = router

# 封装好API供应商demo url=https://console.open.onebound.cn/console/?i=Lex
# ==================== 灰度路由 ====================
class GrayReleaseRouter:
    """
    灰度上线6原则实现
    1. 按平台灰度: 先闲鱼 -> 再Mercari -> 最后Back Market
    2. 按SKU灰度: 先低价值 -> 再高价值
    3. 按比例灰度: 10% -> 30% -> 100%
    4. 按用户灰度: 白名单卖家先上
    5. 按时间段灰度: 工作日白天 -> 周末
    6. 自动回滚: 错误率 > 5% 自动切回
    """

    def __init__(self):
        # 平台灰度顺序
        self.platform_order = ["xianyu", "mercari", "backmarket", "ebay", "vinted"]
        self.current_platform_index: int = 0

        # 比例灰度
        self.traffic_percent: float = 0.0   # 0.0 ~ 1.0

        # 白名单
        self.whitelist_skus: List[str] = []
        self.whitelist_sellers: List[str] = []

        # 自动回滚
        self.error_rate: float = 0.0
        self.error_threshold: float = 0.05
        self.auto_rollback: bool = False

        # 时间段
        self.enabled_time_ranges: List[tuple] = []  # [(9, 18)] 工作日9-18点

    def should_process(self, event: Event) -> bool:
        """灰度6原则综合判断"""
        # 原则1: 平台灰度
        platform = event.metadata.get("platform", "")
        if platform and platform not in self.platform_order[:self.current_platform_index + 1]:
            return False

        # 原则2: SKU白名单
        sku = event.payload.get("sku", "")
        if sku and sku in self.whitelist_skus:
            return True

        # 原则3: 比例灰度
        if random.random() > self.traffic_percent:
            return False

        # 原则4: 卖家白名单
        seller = event.metadata.get("seller", "")
        if seller and seller in self.whitelist_sellers:
            return True

        # 原则5: 时间段
        if self.enabled_time_ranges:
            current_hour = datetime.now().hour
            allowed = any(start <= current_hour < end
                         for start, end in self.enabled_time_ranges)
            if not allowed:
                return False

        # 原则6: 自动回滚
        if self.auto_rollback:
            return False

        return True

    def advance_platform(self):
        """推进到下一个平台"""
        if self.current_platform_index < len(self.platform_order) - 1:
            self.current_platform_index += 1
            print(f"[GRAY] 推进到平台: {self.platform_order[self.current_platform_index]}")

    def set_traffic(self, percent: float):
        self.traffic_percent = max(0.0, min(1.0, percent))

    def report_error(self, error_count: int, total_count: int):
        if total_count > 0:
            self.error_rate = error_count / total_count
            if self.error_rate >= self.error_threshold:
                self.auto_rollback = True
                print(f"[ROLLBACK] 错误率 {self.error_rate:.1%} >= {self.error_threshold:.0%}, 自动回滚")


# ==================== 业务事件处理器示例 ====================
class ProductEventHandlers:
    """商品相关事件处理器"""

    @staticmethod
    def on_product_created(event: Event):
        product = event.payload
        print(f"[EVENT] 商品创建: {product.get('sku')} | {product.get('title')}")
        # 触发: 库存初始化 + 平台同步

    @staticmethod
    def on_product_updated(event: Event):
        product = event.payload
        print(f"[EVENT] 商品更新: {product.get('sku')} | 变更字段: {event.metadata.get('changed_fields', [])}")
        # 触发: 平台同步 + 价格监控

    @staticmethod
    def on_product_price_changed(event: Event):
        print(f"[EVENT] 价格变更: {event.payload.get('sku')} | "
              f"旧: {event.metadata.get('old_price')} -> 新: {event.payload.get('price')}")


class OrderEventHandlers:
    """订单相关事件处理器"""

    @staticmethod
    def on_order_confirmed(event: Event):
        order = event.payload
        print(f"[EVENT] 订单确认: {order.get('order_id')} | {order.get('platform')}")
        # 触发: 库存预留 + 发货准备

    @staticmethod
    def on_order_shipped(event: Event):
        print(f"[EVENT] 订单发货: {event.payload.get('order_id')} | "
              f"运单号: {event.payload.get('tracking_number')}")
        # 触发: 平台回传运单号 + 买家通知

    @staticmethod
    def on_order_returned(event: Event):
        ret = event.payload
        print(f"[EVENT] 退货: {ret.get('return_id')} | 原因: {ret.get('reason_code')}")
        # 触发: 库存回库 + 等级重检 + 退款

# 封装好API供应商demo url=https://console.open.onebound.cn/console/?i=Lex
# ==================== 演示 ====================
if __name__ == "__main__":
    print("=" * 60)
    print("二手ERP事件驱动架构演示")
    print("=" * 60)

    # 1. 初始化事件总线 + 灰度路由器
    bus = EventBus()
    gray = GrayReleaseRouter()
    gray.set_traffic(0.3)  # 30%流量
    gray.platform_order = ["xianyu", "mercari", "backmarket"]
    gray.current_platform_index = 0  # 只放闲鱼
    bus.set_gray_router(gray)

    # 2. 注册事件处理器
    bus.subscribe(Subscription("product.created", ProductEventHandlers.on_product_created))
    bus.subscribe(Subscription("product.updated", ProductEventHandlers.on_product_updated))
    bus.subscribe(Subscription("product.price_changed", ProductEventHandlers.on_product_price_changed))
    bus.subscribe(Subscription("order.confirmed", OrderEventHandlers.on_order_confirmed))
    bus.subscribe(Subscription("order.shipped", OrderEventHandlers.on_order_shipped))
    bus.subscribe(Subscription("order.returned", OrderEventHandlers.on_order_returned))

    # 3. 模拟事件流
    events = [
        Event(
            event_id="evt-001",
            type="product.created",
            source="product-service",
            payload={"sku": "IP14P-256-SILVER", "title": "iPhone 14 Pro 256GB", "price": 799.00},
            metadata={"platform": "xianyu", "seller": "seller_a"}
        ),
        Event(
            event_id="evt-002",
            type="product.created",
            source="product-service",
            payload={"sku": "IP13-128-BLACK", "title": "iPhone 13 128GB", "price": 499.00},
            metadata={"platform": "mercari", "seller": "seller_b"}  # Mercari 还没灰度到
        ),
        Event(
            event_id="evt-003",
            type="product.price_changed",
            source="pricing-service",
            payload={"sku": "IP14P-256-SILVER", "price": 749.00},
            metadata={"old_price": 799.00, "platform": "xianyu"}
        ),
        Event(
            event_id="evt-004",
            type="order.confirmed",
            source="order-service",
            payload={"order_id": "ORD-001", "platform": "xianyu", "sku": "IP14P-256-SILVER", "quantity": 1},
            metadata={"platform": "xianyu"}
        ),
        Event(
            event_id="evt-005",
            type="order.returned",
            source="return-service",
            payload={
                "return_id": "RET-001", "order_id": "ORD-001", "platform": "xianyu",
                "sku": "IP14P-256-SILVER", "reason_code": "not_as_described",
                "grade_downgraded": "good"
            },
            metadata={"platform": "xianyu"}
        ),
    ]

    for evt in events:
        print(f"\n--- 发布事件: {evt.type} [{evt.event_id}] ---")
        bus.publish(evt)

    # 4. 灰度推进
    print(f"\n{'='*60}")
    print("灰度推进: 开启Mercari (30%流量)")
    gray.advance_platform()
    bus.publish(Event(
        event_id="evt-006",
        type="product.created",
        source="product-service",
        payload={"sku": "IP13-128-BLACK", "title": "iPhone 13 128GB", "price": 449.00},
        metadata={"platform": "mercari", "seller": "seller_b"}
    ))

    # 5. 模拟错误率触发回滚
    print(f"\n{'='*60}")
    print("模拟错误率过高 -> 自动回滚")
    gray.report_error(error_count=50, total_count=900)
    print(f"错误率: {gray.error_rate:.1%} | 自动回滚: {gray.auto_rollback}")

    bus.publish(Event(
        event_id="evt-007",
        type="product.created",
        source="product-service",
        payload={"sku": "IP12-64-WHITE", "title": "iPhone 12 64GB", "price": 299.00},
        metadata={"platform": "mercari", "seller": "seller_c"}
    ))
运行结果:
============================================================
二手ERP事件驱动架构演示
============================================================

--- 发布事件: product.created [evt-001] ---
[EVENT] 商品创建: IP14P-256-SILVER | iPhone 14 Pro 256GB

--- 发布事件: product.created [evt-002] ---
(Mercari 被灰度拦截, 无输出)

--- 发布事件: product.price_changed [evt-003] ---
[EVENT] 价格变更: IP14P-256-SILVER | 旧: 799.0 -> 新: 749.0

--- 发布事件: order.confirmed [evt-004] ---
[EVENT] 订单确认: ORD-001 | xianyu

--- 发布事件: order.returned [evt-005] ---
[EVENT] 退货: RET-001 | 原因: not_as_described

============================================================
灰度推进: 开启Mercari (30%流量)
============================================================
[GRAY] 推进到平台: mercari
--- 发布事件: product.created [evt-006] ---
(30%概率命中, 可能不输出; 多跑几次就会命中)

============================================================
模拟错误率过高 -> 自动回滚
============================================================
[ROLLBACK] 错误率 5.6% >= 5%, 自动回滚
--- 发布事件: product.created [evt-007] ---
(自动回滚激活, 所有事件被拦截)

四、灰度上线6原则详解

原则
实现
为什么重要
1. 按平台灰度
platform_order + current_platform_index
闲鱼出问题只影响国内,不会炸掉Back Market EU站
2. 按SKU灰度
whitelist_skus
先用200的手机测,最后上$1000的MacBook
3. 按比例灰度
traffic_percent 0.1→0.3→1.0
10%订单先走新逻辑,观察24小时没问题再放大
4. 按用户灰度
whitelist_sellers
先让合作3年的老卖家试用,新卖家走旧逻辑
5. 按时间段灰度
enabled_time_ranges
工作日白天上线,出问题有人修;周五下午不上线
6. 自动回滚
error_rate > 5% → auto_rollback=True
不等人报警,系统自己切回去

五、和前22篇的完整衔接

统一数据层 (本篇)
  ├── InternalProduct ← 前篇 cross_border_schema.FieldMapper
  ├── InternalOrder   ← 前篇 ebay_secondhand_defense.OrderGuard
  ├── InternalInventory ← 前篇 vinted_adapter.VintedStockSync
  └── InternalReturn  ← 前篇 grade_integrity_loop.ReturnSignalLoop

适配器层 (前22篇)
  ├── xianyu_adapter        ← 国内二手基础
  ├── mercari_rate_limiter  ← 限流/TokenPool
  ├── ebay_secondhand_defense ← ConditionGuard/AccountHealth
  ├── vinted_adapter        ← GPSR/StockSync
  ├── backmarket_compliance_gate ← 入驻门禁/Deposit
  └── grade_integrity_loop  ← A/B校验闭环

事件总线 (本篇)
  ├── product.created  → 触发: 库存初始化 + 平台同步
  ├── order.confirmed  → 触发: 库存预留 + 发货准备
  ├── order.returned   → 触发: 库存回库 + 等级重检 + 退款
  └── inventory.changed → 触发: 平台库存同步

灰度上线 (本篇)
  ├── 按平台: 闲鱼 → Mercari → Back Market → eBay → Vinted
  ├── 按比例: 10% → 30% → 100%
  └── 自动回滚: 错误率 > 5% 切回

六、部署建议

  1. 第一阶段:统一数据模型先落地,把国内ERP的 product 表通过 to_internal() 转成 InternalProduct

  2. 第二阶段:事件总线替换现有的同步调用,先从 order.confirmedinventory.changed 开始

  3. 第三阶段:灰度上线框架搭好,每次新接一个平台走完6原则

  4. 第四阶段:适配器层逐步替换为事件驱动版本


要不要我把这套方案做成一个可运行的脚手架项目:包含 FastAPI 服务、SQLite/PostgreSQL 持久化、RabbitMQ 事件总线集成、以及一个 Web 管理后台(查看灰度状态/手动回滚/事件追踪)?这样可以直接 clone 下来跑,不用从零搭。


群贤毕至

访客