Python脚本如何统一多环境数据同步规则

wen python案例 32

本文目录导读:

Python脚本如何统一多环境数据同步规则

  1. 核心架构设计
  2. 核心同步引擎
  3. 数据转换器
  4. 冲突解决机制
  5. 完整使用示例
  6. 配置文件示例 (sync_rules.yaml)
  7. 最佳实践建议

我来为你介绍一个统一多环境数据同步规则的 Python 方案,这个系统可以让你在不同环境(开发、测试、生产)之间进行标准化的数据同步。

核心架构设计

import json
import yaml
import hashlib
from datetime import datetime
from typing import Dict, List, Optional
from dataclasses import dataclass, asdict
from enum import Enum
class SyncEnvironment(Enum):
    """同步环境枚举"""
    DEV = "development"
    TEST = "testing"
    STAGING = "staging"
    PRODUCTION = "production"
@dataclass
class SyncRule:
    """同步规则定义"""
    source_table: str
    target_table: str
    sync_type: str  # full, incremental, delta
    field_mapping: Dict[str, str]
    filter_condition: Optional[str] = None
    transform_rules: Optional[List[Dict]] = None
    conflict_resolution: str = "source_wins"  # source_wins, target_wins, merge
@dataclass
class SyncConfig:
    """同步配置"""
    rules: List[SyncRule]
    schedule: Optional[str] = None
    batch_size: int = 1000
    error_handling: str = "skip"  # skip, stop, log
    audit_enabled: bool = True

核心同步引擎

class DataSyncEngine:
    """统一数据同步引擎"""
    def __init__(self, source_config: Dict, target_config: Dict):
        self.source_config = source_config
        self.target_config = target_config
        self.sync_rules = self._load_sync_rules()
        self.audit_log = []
    def _load_sync_rules(self) -> List[SyncRule]:
        """加载同步规则"""
        # 从配置文件加载
        with open('sync_rules.yaml', 'r') as f:
            rules_data = yaml.safe_load(f)
        return [SyncRule(**rule) for rule in rules_data['rules']]
    def execute_sync(self, environment: SyncEnvironment) -> Dict:
        """执行同步任务"""
        sync_result = {
            'start_time': datetime.now(),
            'environment': environment.value,
            'rules_executed': [],
            'failed_rules': []
        }
        for rule in self.sync_rules:
            try:
                rule_result = self._execute_rule(rule, environment)
                sync_result['rules_executed'].append(rule_result)
            except Exception as e:
                sync_result['failed_rules'].append({
                    'rule': rule.source_table,
                    'error': str(e)
                })
        sync_result['end_time'] = datetime.now()
        self._log_audit(sync_result)
        return sync_result
    def _execute_rule(self, rule: SyncRule, environment: SyncEnvironment) -> Dict:
        """执行单个同步规则"""
        # 1. 提取源数据
        source_data = self._extract_data(rule, environment)
        # 2. 数据转换
        transformed_data = self._transform_data(source_data, rule)
        # 3. 加载到目标
        load_result = self._load_data(transformed_data, rule, environment)
        return {
            'rule': rule.source_table,
            'records_processed': len(source_data),
            'records_loaded': load_result['loaded'],
            'records_failed': load_result['failed']
        }

数据转换器

class DataTransformer:
    """数据转换器 - 处理不同环境间的数据差异"""
    def __init__(self, environment: SyncEnvironment):
        self.environment = environment
        self.transformers = self._register_transformers()
    def _register_transformers(self) -> Dict:
        """注册转换器"""
        return {
            'date_format': self._transform_date_format,
            'currency': self._transform_currency,
            'id_mapping': self._transform_id_mapping,
            'sensitive_data': self._mask_sensitive_data,
            'field_rename': self._rename_fields
        }
    def transform(self, data: List[Dict], rules: List[Dict]) -> List[Dict]:
        """执行数据转换"""
        transformed_data = data.copy()
        for rule in rules:
            transformer = self.transformers.get(rule['type'])
            if transformer:
                transformed_data = transformer(transformed_data, rule)
        return transformed_data
    def _transform_date_format(self, data: List[Dict], rule: Dict) -> List[Dict]:
        """转换日期格式"""
        from_format = rule.get('from', '%Y-%m-%d')
        to_format = rule.get('to', '%Y-%m-%d %H:%M:%S')
        for record in data:
            if rule['field'] in record:
                try:
                    date_obj = datetime.strptime(record[rule['field']], from_format)
                    record[rule['field']] = date_obj.strftime(to_format)
                except ValueError:
                    pass
        return data
    def _transform_currency(self, data: List[Dict], rule: Dict) -> List[Dict]:
        """转换货币单位和格式"""
        rate = rule.get('conversion_rate', 1)
        for record in data:
            if rule['field'] in record:
                record[rule['field']] = round(record[rule['field']] * rate, 2)
        return data
    def _mask_sensitive_data(self, data: List[Dict], rule: Dict) -> List[Dict]:
        """屏蔽敏感数据"""
        for record in data:
            if rule['field'] in record:
                value = str(record[rule['field']])
                mask_char = rule.get('mask_char', '*')
                visible_chars = rule.get('visible_chars', 4)
                if len(value) > visible_chars:
                    masked = value[:visible_chars] + mask_char * (len(value) - visible_chars)
                    record[rule['field']] = masked
        return data

冲突解决机制

class ConflictResolver:
    """冲突解决器"""
    def __init__(self, strategy: str = "source_wins"):
        self.strategy = strategy
        self.resolvers = {
            "source_wins": self._source_wins,
            "target_wins": self._target_wins,
            "merge": self._merge,
            "timestamp_based": self._timestamp_based
        }
    def resolve(self, source_record: Dict, target_record: Dict, 
                fields: List[str]) -> Dict:
        """解决数据冲突"""
        resolver = self.resolvers.get(self.strategy, self._source_wins)
        return resolver(source_record, target_record, fields)
    def _timestamp_based(self, source_record: Dict, target_record: Dict,
                        fields: List[str]) -> Dict:
        """基于时间戳的冲突解决"""
        source_ts = source_record.get('updated_at', datetime.min)
        target_ts = target_record.get('updated_at', datetime.min)
        if isinstance(source_ts, str):
            source_ts = datetime.fromisoformat(source_ts)
        if isinstance(target_ts, str):
            target_ts = datetime.fromisoformat(target_ts)
        return target_record if target_ts > source_ts else source_record

完整使用示例

def create_sync_rules_config():
    """创建同步规则配置文件"""
    rules = {
        'rules': [
            {
                'source_table': 'users',
                'target_table': 'users',
                'sync_type': 'incremental',
                'field_mapping': {
                    'id': 'user_id',
                    'name': 'full_name',
                    'email': 'email_address'
                },
                'filter_condition': "status = 'active'",
                'transform_rules': [
                    {
                        'type': 'date_format',
                        'field': 'created_at',
                        'from': '%Y-%m-%d',
                        'to': '%Y-%m-%d %H:%M:%S'
                    },
                    {
                        'type': 'sensitive_data',
                        'field': 'email',
                        'mask_char': '*',
                        'visible_chars': 3
                    }
                ],
                'conflict_resolution': 'timestamp_based'
            }
        ]
    }
    with open('sync_rules.yaml', 'w') as f:
        yaml.dump(rules, f, default_flow_style=False)
def main():
    """主函数示例"""
    # 配置源数据库和目标数据库
    source_config = {
        'host': 'localhost',
        'port': 5432,
        'database': 'source_db',
        'user': 'user',
        'password': 'password'
    }
    target_config = {
        'host': 'production-server',
        'port': 5432,
        'database': 'target_db',
        'user': 'user',
        'password': 'password'
    }
    # 创建同步引擎
    sync_engine = DataSyncEngine(source_config, target_config)
    # 执行同步
    result = sync_engine.execute_sync(SyncEnvironment.DEV)
    # 输出结果
    print(f"同步完成,处理记录数: {result['rules_executed']}")
    print(f"失败规则: {result['failed_rules']}")
    # 记录审计日志
    sync_engine._log_audit(result)
if __name__ == "__main__":
    main()

配置文件示例 (sync_rules.yaml)

rules:
  - source_table: users
    target_table: users_prod
    sync_type: incremental
    field_mapping:
      id: user_id
      name: username
      email: email
    filter_condition: "is_deleted = false"
    transform_rules:
      - type: date_format
        field: created_at
        from: "%Y-%m-%d %H:%M:%S"
        to: "%Y-%m-%d"
    conflict_resolution: source_wins
  - source_table: orders
    target_table: orders
    sync_type: full
    field_mapping:
      order_id: id
      total: amount
      currency: currency_code
    transform_rules:
      - type: currency
        field: amount
        conversion_rate: 0.85  # 汇率转换
    conflict_resolution: merge

最佳实践建议

  1. 使用数据版本控制:为每条记录添加版本号或时间戳
  2. 实施断点续传:处理大量数据时支持从中断点恢复
  3. 添加数据校验:在同步前后进行数据完整性检查
  4. 监控和告警:实现实时监控和异常告警机制
  5. 回滚机制:支持在同步失败时回滚到之前的状态

这个系统提供了一个灵活、可扩展的数据同步框架,可以根据具体需求进行调整和扩展。

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