本文目录导读:

我来为你介绍一个统一多环境数据同步规则的 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
最佳实践建议
- 使用数据版本控制:为每条记录添加版本号或时间戳
- 实施断点续传:处理大量数据时支持从中断点恢复
- 添加数据校验:在同步前后进行数据完整性检查
- 监控和告警:实现实时监控和异常告警机制
- 回滚机制:支持在同步失败时回滚到之前的状态
这个系统提供了一个灵活、可扩展的数据同步框架,可以根据具体需求进行调整和扩展。