Python脚本如何优先同步核心业务数据

wen python案例 32

本文目录导读:

Python脚本如何优先同步核心业务数据

  1. 基于优先级队列的异步架构(推荐)
  2. 基于分片策略的批处理同步
  3. 基于Redis的分布式优先级同步(分布式场景)
  4. 基于依赖关系的DAG同步(复杂依赖场景)
  5. 通用实践建议

在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. 核心数据立即同步(同步/异步均可,但不可延迟)
  2. 非核心数据批量延迟(聚合多个小请求,减少连接开销)
  3. 失败重试分级(核心数据重试次数多,普通数据少)
  4. 监控与追踪(记录每个任务的优先级、耗时、状态)
  5. 可配置性(优先级、并发数、超时时间等应外部化配置)

根据你的实际场景选择最适合的方案:小规模单进程选方案1,大规模分布式选方案3,有复杂依赖关系选方案4。

抱歉,评论功能暂时关闭!