本文目录导读:

我来分享几种Python自动化修复同步异常数据的方法:
基础检测与修复脚本
import pandas as pd
import numpy as np
from datetime import datetime, timedelta
import logging
# 配置日志
logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s')
logger = logging.getLogger(__name__)
class SyncDataRepairer:
def __init__(self, source_data, target_data, sync_key='id'):
"""
初始化同步修复器
:param source_data: 源数据 DataFrame
:param target_data: 目标数据 DataFrame
:param sync_key: 同步关键字段
"""
self.source = source_data
self.target = target_data
self.sync_key = sync_key
def detect_missing_data(self):
"""检测缺失数据"""
source_ids = set(self.source[self.sync_key])
target_ids = set(self.target[self.sync_key])
missing_in_target = source_ids - target_ids
missing_in_source = target_ids - source_ids
return {
'missing_in_target': missing_in_target,
'missing_in_source': missing_in_source
}
def detect_inconsistent_data(self, compare_fields=None):
"""检测不一致数据"""
if compare_fields is None:
compare_fields = self.source.columns.tolist()
merged = pd.merge(
self.source,
self.target,
on=self.sync_key,
how='inner',
suffixes=('_source', '_target')
)
inconsistencies = []
for field in compare_fields:
if field != self.sync_key:
source_field = f'{field}_source'
target_field = f'{field}_target'
if source_field in merged.columns and target_field in merged.columns:
# 忽略NaN值的比较
mask = merged[source_field] != merged[target_field]
mask &= ~(merged[source_field].isna() & merged[target_field].isna())
inconsistent = merged[mask][[self.sync_key, source_field, target_field]]
if not inconsistent.empty:
inconsistencies.append((field, inconsistent))
return inconsistencies
def repair_missing_data(self, direction='target'):
"""修复缺失数据"""
missing = self.detect_missing_data()
if direction == 'target':
missing_ids = missing['missing_in_target']
missing_data = self.source[self.source[self.sync_key].isin(missing_ids)]
repaired_data = missing_data.copy()
logger.info(f"修复了 {len(missing_ids)} 条目标数据缺失")
return repaired_data
elif direction == 'source':
missing_ids = missing['missing_in_source']
missing_data = self.target[self.target[self.sync_key].isin(missing_ids)]
repaired_data = missing_data.copy()
logger.info(f"修复了 {len(missing_ids)} 条源数据缺失")
return repaired_data
return None
def repair_inconsistent_data(self, compare_fields=None, strategy='source_override'):
"""修复不一致数据"""
inconsistencies = self.detect_inconsistent_data(compare_fields)
repairs = []
for field, inconsistent_data in inconsistencies:
if strategy == 'source_override':
# 以源数据为准覆盖目标数据
for _, row in inconsistent_data.iterrows():
repair = {
self.sync_key: row[self.sync_key],
'field': field,
'old_value': row[f'{field}_target'],
'new_value': row[f'{field}_source'],
'strategy': 'source_override'
}
repairs.append(repair)
elif strategy == 'target_override':
# 以目标数据为准
for _, row in inconsistent_data.iterrows():
repair = {
self.sync_key: row[self.sync_key],
'field': field,
'old_value': row[f'{field}_source'],
'new_value': row[f'{field}_target'],
'strategy': 'target_override'
}
repairs.append(repair)
logger.info(f"发现并记录了 {len(repairs)} 条不一致数据")
return repairs
def auto_repair(self, strategy='source_override'):
"""自动修复所有异常"""
# 1. 修复缺失数据
repaired_missing = self.repair_missing_data()
# 2. 修复不一致数据
repairs = self.repair_inconsistent_data(strategy=strategy)
# 3. 应用修复
if repairs:
for repair in repairs:
if repair['strategy'] in ['source_override', 'latest_timestamp']:
# 这里实现具体的修复逻辑
logger.info(f"修复记录: {repair}")
return {
'missing_data_repaired': len(repaired_missing) if repaired_missing is not None else 0,
'inconsistencies_repaired': len(repairs)
}
带时间戳的高级修复方案
from datetime import datetime
import hashlib
class AdvancedSyncRepairer(SyncDataRepairer):
def __init__(self, source_data, target_data, sync_key='id', timestamp_field='update_time'):
super().__init__(source_data, target_data, sync_key)
self.timestamp_field = timestamp_field
def timestamp_based_repair(self):
"""基于时间戳的智能修复"""
merged = pd.merge(
self.source,
self.target,
on=self.sync_key,
how='outer',
suffixes=('_source', '_target')
)
repairs = []
for _, row in merged.iterrows():
# 检查是否需要修复
if pd.isna(row.get(f'{self.timestamp_field}_source')):
continue
source_time = row[f'{self.timestamp_field}_source']
target_time = row.get(f'{self.timestamp_field}_target')
# 如果目标没有该记录或源数据更新
if pd.isna(target_time) or source_time > target_time:
repair = self._create_timestamp_repair(row)
repairs.append(repair)
return repairs
def _create_timestamp_repair(self, row):
"""创建时间戳修复记录"""
repair = {
self.sync_key: row[self.sync_key],
'source_timestamp': row.get(f'{self.timestamp_field}_source'),
'target_timestamp': row.get(f'{self.timestamp_field}_target'),
'action': 'update' if not pd.isna(row.get(f'{self.timestamp_field}_target')) else 'insert'
}
return repair
def hash_based_comparison(self, fields_to_check):
"""基于哈希值的比较"""
def create_hash(row, fields):
data = '-'.join([str(row[f]) for f in fields])
return hashlib.md5(data.encode()).hexdigest()
# 为源数据和目标数据创建哈希
self.source['_hash'] = self.source.apply(
lambda row: create_hash(row, fields_to_check), axis=1
)
self.target['_hash'] = self.target.apply(
lambda row: create_hash(row, fields_to_check), axis=1
)
# 比较哈希值
merged = pd.merge(
self.source[[self.sync_key, '_hash']],
self.target[[self.sync_key, '_hash']],
on=self.sync_key,
suffixes=('_source', '_target'),
how='outer'
)
# 找出不一致的数据
inconsistent = merged[
(merged['_hash_source'] != merged['_hash_target']) |
(merged['_hash_source'].isna()) |
(merged['_hash_target'].isna())
]
return inconsistent
数据库同步修复脚本
import mysql.connector
from sqlalchemy import create_engine
class DatabaseSyncRepairer:
def __init__(self, source_db_config, target_db_config):
self.source_engine = create_engine(
f"mysql+mysqlconnector://{source_db_config['user']}:{source_db_config['password']}@"
f"{source_db_config['host']}/{source_db_config['database']}"
)
self.target_engine = create_engine(
f"mysql+mysqlconnector://{target_db_config['user']}:{target_db_config['password']}@"
f"{target_db_config['host']}/{target_db_config['database']}"
)
def sync_table(self, table_name, sync_key='id'):
"""同步特定表"""
# 读取数据
source_df = pd.read_sql(f"SELECT * FROM {table_name}", self.source_engine)
target_df = pd.read_sql(f"SELECT * FROM {table_name}", self.target_engine)
# 使用修复器
repairer = SyncDataRepairer(source_df, target_df, sync_key)
result = repairer.auto_repair()
# 应用修复到数据库
if result['missing_data_repaired'] > 0:
missing_data = repairer.repair_missing_data()
missing_data.to_sql(
table_name,
self.target_engine,
if_exists='append',
index=False
)
return result
def batch_sync_tables(self, tables_config):
"""批量同步多个表"""
results = {}
for table_name, config in tables_config.items():
logger.info(f"开始同步表: {table_name}")
results[table_name] = self.sync_table(
table_name,
config.get('sync_key', 'id')
)
return results
使用示例
# 示例1:基础修复
source_data = pd.DataFrame({
'id': [1, 2, 3, 4, 5],
'name': ['A', 'B', 'C', 'D', 'E'],
'value': [100, 200, 300, 400, 500]
})
target_data = pd.DataFrame({
'id': [1, 2, 3, 4],
'name': ['A', 'B', 'C', 'D'],
'value': [100, 250, 300, 400] # id=2 的值不同
})
repairer = SyncDataRepairer(source_data, target_data)
result = repairer.auto_repair()
print(f"修复结果: {result}")
# 示例2:带时间戳的修复
source_with_ts = source_data.copy()
source_with_ts['update_time'] = datetime.now()
target_with_ts = target_data.copy()
target_with_ts['update_time'] = datetime.now() - timedelta(hours=1)
advanced_repairer = AdvancedSyncRepairer(
source_with_ts,
target_with_ts,
timestamp_field='update_time'
)
repairs = advanced_repairer.timestamp_based_repair()
print(f"时间戳修复记录: {repairs}")
最佳实践建议
- 定期执行:设置cron job定期运行同步脚本
- 日志记录:详细记录每次修复操作
- 备份数据:修改前备份目标数据
- 增量同步:只同步变化的数据
- 异常处理:添加重试机制和报警
这个方案可以根据你的具体需求进行调整和扩展。