收官篇终于来了。前面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% 切回
六、部署建议
第一阶段:统一数据模型先落地,把国内ERP的
product表通过to_internal()转成InternalProduct第二阶段:事件总线替换现有的同步调用,先从
order.confirmed和inventory.changed开始第三阶段:灰度上线框架搭好,每次新接一个平台走完6原则
第四阶段:适配器层逐步替换为事件驱动版本
要不要我把这套方案做成一个可运行的脚手架项目:包含 FastAPI 服务、SQLite/PostgreSQL 持久化、RabbitMQ 事件总线集成、以及一个 Web 管理后台(查看灰度状态/手动回滚/事件追踪)?这样可以直接 clone 下来跑,不用从零搭。