本文目录导读:

在Python中实现数据同步时,优先同步核心业务数据的关键在于优先级队列、异步分片和依赖管理,以下是几种成熟的实现方案,按推荐程度从高到低排列:
基于优先级队列的异步架构(推荐)
使用 asyncio + PriorityQueue 实现精细化调度。
import asyncio
from dataclasses import dataclass, field
from typing import Any, Callable
from enum import IntEnum
class Priority(IntEnum):
CRITICAL = 1 # 核心业务
HIGH = 2 # 重要辅助
NORMAL = 3 # 常规数据
LOW = 4 # 日志/归档
@dataclass(order=True)
class SyncTask:
priority: int
data: Any = field(compare=False)
sync_fn: Callable = field(compare=False)
class PrioritySyncEngine:
def __init__(self, max_concurrent=5):
self.queue = asyncio.PriorityQueue()
self.semaphore = asyncio.Semaphore(max_concurrent)
async def add_task(self, task: SyncTask):
await self.queue.put(task)
async def worker(self, worker_id):
while True:
task = await self.queue.get()
async with self.semaphore: # 控制并发数
try:
await task.sync_fn(task.data)
except Exception as e:
print(f"Worker {worker_id} failed: {e}")
finally:
self.queue.task_done()
async def run(self, num_workers=3):
workers = [asyncio.create_task(self.worker(i)) for i in range(num_workers)]
await self.queue.join()
for w in workers:
w.cancel()
# 使用示例
async def sync_orders(orders): # 核心业务
print(f"Syncing orders: {len(orders)}")
# 实际的同步逻辑
async def main():
engine = PrioritySyncEngine(max_concurrent=3)
# 优先添加核心业务任务
await engine.add_task(SyncTask(Priority.CRITICAL, {"order_id": 1}, sync_orders))
await engine.add_task(SyncTask(Priority.HIGH, ["user_updates"], sync_users))
await engine.add_task(SyncTask(Priority.NORMAL, {"logs": []}, sync_logs))
await engine.run()
asyncio.run(main())
基于分片策略的批处理同步
当数据量大时,核心数据应立即同步,非核心数据可分批处理。
import threading
from collections import deque
from typing import List, Dict
class TieredSyncManager:
def __init__(self):
self.tiers = {
'tier1': deque(), # 核心业务
'tier2': deque(), # 重要
'tier3': deque(), # 普通
}
self.lock = threading.Lock()
self.sync_interval = {'tier1': 5, 'tier2': 30, 'tier3': 120}
def add_data(self, data: Dict, tier: str):
with self.lock:
self.tiers[tier].append(data)
def _sync_tier(self, tier: str):
"""同步指定层级的数据"""
batch = []
with self.lock:
while self.tiers[tier] and len(batch) < 100:
batch.append(self.tiers[tier].popleft())
if batch:
if tier == 'tier1':
self._critical_sync(batch) # 实时同步
else:
threading.Timer(self.sync_interval[tier], self._sync_tier, [tier]).start()
def _critical_sync(self, batch: List):
"""核心数据立即同步(阻塞或异步)"""
print(f"Syncing {len(batch)} critical items immediately")
# 实际同步代码
def run(self):
# 核心业务立即同步,其他定时同步
self._sync_tier('tier1')
for tier in ['tier2', 'tier3']:
self._sync_tier(tier)
# 使用
manager = TieredSyncManager()
manager.add_data({"order": 123}, 'tier1')
manager.run()
基于Redis的分布式优先级同步(分布式场景)
适用于微服务架构,多个同步进程共享优先级。
import redis
import json
from datetime import datetime
class RedisPrioritySync:
def __init__(self, redis_url='redis://localhost:6379/0'):
self.client = redis.from_url(redis_url)
self.priority_keys = {
1: 'sync:critical', # 核心
2: 'sync:high', # 重要
3: 'sync:normal', # 普通
4: 'sync:low' # 低优
}
def enqueue(self, data: dict, priority: int):
"""将数据按优先级推入Redis列表"""
key = self.priority_keys.get(priority, 'sync:normal')
payload = {
'data': data,
'timestamp': datetime.now().isoformat(),
'priority': priority
}
self.client.lpush(key, json.dumps(payload))
def dequeue(self, batch_size=10):
"""按优先级从高到低获取数据"""
for priority in sorted(self.priority_keys.keys()):
key = self.priority_keys[priority]
items = []
for _ in range(batch_size):
item = self.client.rpop(key)
if item:
items.append(json.loads(item))
else:
break
if items:
return items, priority
return [], None
def sync_loop(self):
"""轮询同步"""
while True:
items, priority = self.dequeue()
if items:
print(f"Syncing {len(items)} items with priority {priority}")
# 实际同步逻辑
self._process_batch(items)
# 生产者
redis_sync = RedisPrioritySync()
redis_sync.enqueue({"order_id": 1}, priority=1) # 核心数据
redis_sync.enqueue({"log": "info"}, priority=4) # 低优数据
# 消费者(可多实例)
redis_sync.sync_loop()
基于依赖关系的DAG同步(复杂依赖场景)
当核心业务数据依赖其他数据时,使用有向无环图(DAG)调度。
from collections import defaultdict
from typing import Set, Dict, List
class DAGSyncScheduler:
def __init__(self):
self.graph = defaultdict(set) # 依赖关系
self.data = {} # 实际数据
self.priority = {} # 优先级权重
def add_node(self, node_id: str, data: any, priority: int, dependencies: List[str] = None):
self.data[node_id] = data
self.priority[node_id] = priority
if dependencies:
for dep in dependencies:
self.graph[node_id].add(dep)
def _get_execution_order(self) -> List[str]:
"""拓扑排序,高优先级优先"""
in_degree = {node: 0 for node in self.data}
for node, deps in self.graph.items():
for dep in deps:
in_degree[node] += 1
# 按优先级排序的队列
from heapq import heappush, heappop
queue = []
for node, degree in in_degree.items():
if degree == 0:
# 优先级越高(数值越小),越先执行
heappush(queue, (self.priority[node], node))
order = []
while queue:
_, node = heappop(queue)
order.append(node)
# 减少依赖节点的入度
for other, deps in self.graph.items():
if node in deps:
in_degree[other] -= 1
if in_degree[other] == 0:
heappush(queue, (self.priority[other], other))
return order
def execute(self):
order = self._get_execution_order()
print(f"Execution order: {order}")
for node in order:
print(f"Syncing {node} (priority {self.priority[node]})")
# 实际同步逻辑
# 使用
scheduler = DAGSyncScheduler()
scheduler.add_node("user_orders", {"orders": []}, priority=1) # 核心,无依赖
scheduler.add_node("payment_info", {"payments": []}, priority=1, dependencies=["user_orders"])
scheduler.add_node("product_catalog", {"products": []}, priority=2)
scheduler.add_node("analytics", {"stats": []}, priority=3, dependencies=["user_orders", "payment_info"])
scheduler.execute()
# 输出: user_orders -> product_catalog -> payment_info -> analytics
通用实践建议
配置化优先级管理
# config/sync_config.yaml
sync_tiers:
critical: # 核心业务
tables: ['orders', 'payments', 'inventory']
strategy: 'immediate'
retry: 3
timeout: 5
high:
tables: ['users', 'products']
strategy: 'batch'
batch_size: 50
normal:
tables: ['logs', 'analytics']
strategy: 'scheduled'
interval: 300 # 5分钟
监控与告警
import logging
from prometheus_client import Counter, Histogram
sync_counter = Counter('sync_total', 'Total sync operations', ['priority'])
sync_duration = Histogram('sync_duration_seconds', 'Sync duration', ['priority'])
def monitored_sync(sync_fn, data, priority):
sync_counter.labels(priority=priority).inc()
with sync_duration.labels(priority=priority).time():
return sync_fn(data)
- 核心数据立即同步(同步/异步均可,但不可延迟)
- 非核心数据批量延迟(聚合多个小请求,减少连接开销)
- 失败重试分级(核心数据重试次数多,普通数据少)
- 监控与追踪(记录每个任务的优先级、耗时、状态)
- 可配置性(优先级、并发数、超时时间等应外部化配置)
根据你的实际场景选择最适合的方案:小规模单进程选方案1,大规模分布式选方案3,有复杂依赖关系选方案4。