本文目录导读:

我来详细说明如何使用Python脚本修复多环境数据不一致的问题。
核心方案设计
1 数据差异检测工具
import hashlib
import json
from typing import Dict, List, Any
import pandas as pd
from sqlalchemy import create_engine, text
class DataConsistencyChecker:
"""数据一致性检查器"""
def __init__(self, db_configs: Dict[str, Dict]):
"""
:param db_configs: 环境数据库配置
config = {
'dev': {'host': 'localhost', 'port': 3306, 'user': 'root', 'password': '123', 'db': 'test'},
'staging': {'host': 'staging.com', ...},
'prod': {'host': 'prod.com', ...}
}
"""
self.engines = {}
for env, config in db_configs.items():
conn_str = f"mysql+pymysql://{config['user']}:{config['password']}@{config['host']}:{config['port']}/{config['db']}"
self.engines[env] = create_engine(conn_str)
def get_table_hash(self, env: str, table: str, key_columns: List[str] = None) -> Dict:
"""获取表数据的哈希值"""
engine = self.engines[env]
if key_columns:
# 按主键分组计算Hash
query = f"""
SELECT {','.join(key_columns)},
MD5(CONCAT_WS('|', {','.join([f'COALESCE({col}, "")' for col in key_columns])})) as row_hash
FROM {table}
ORDER BY {key_columns[0]}
"""
else:
query = f"SELECT * FROM {table}"
df = pd.read_sql(query, engine)
return {
'row_count': len(df),
'data_hash': hashlib.md5(df.to_json().encode()).hexdigest(),
'hash_by_key': df.set_index(key_columns)['row_hash'].to_dict() if key_columns else {}
}
def find_differences(self, env1: str, env2: str, table: str, key_columns: List[str]) -> Dict:
"""查找两个环境之间的数据差异"""
env1_data = self.get_table_hash(env1, table, key_columns)
env2_data = self.get_table_hash(env2, table, key_columns)
differences = {
'missing_in_env2': [],
'missing_in_env1': [],
'data_mismatch': [],
'row_count_diff': {
env1: env1_data['row_count'],
env2: env2_data['row_count']
}
}
# 找出差异
for key, hash1 in env1_data['hash_by_key'].items():
if key not in env2_data['hash_by_key']:
differences['missing_in_env2'].append(key)
elif hash1 != env2_data['hash_by_key'][key]:
differences['data_mismatch'].append(key)
for key in env2_data['hash_by_key']:
if key not in env1_data['hash_by_key']:
differences['missing_in_env1'].append(key)
return differences
2 自动化修复脚本
import logging
from datetime import datetime
from typing import List, Optional
class DataRepairTool:
"""数据修复工具"""
def __init__(self, source_env: str, target_env: str, db_configs: Dict):
self.source_env = source_env
self.target_env = target_env
self.checker = DataConsistencyChecker(db_configs)
self.logger = self._setup_logger()
def _setup_logger(self):
"""设置日志记录"""
logger = logging.getLogger('DataRepair')
logger.setLevel(logging.INFO)
# 文件日志
fh = logging.FileHandler(f'repair_{datetime.now().strftime("%Y%m%d_%H%M%S")}.log')
fh.setLevel(logging.INFO)
# 控制台日志
ch = logging.StreamHandler()
ch.setLevel(logging.INFO)
formatter = logging.Formatter('%(asctime)s - %(name)s - %(levelname)s - %(message)s')
fh.setFormatter(formatter)
ch.setFormatter(formatter)
logger.addHandler(fh)
logger.addHandler(ch)
return logger
def sync_data(self, table: str, key_columns: List[str], columns: Optional[List[str]] = None):
"""同步数据从源环境到目标环境"""
self.logger.info(f"开始同步表 {table} 从 {self.source_env} 到 {self.target_env}")
# 1. 检测差异
diffs = self.checker.find_differences(
self.source_env, self.target_env, table, key_columns
)
self.logger.info(f"差异统计:")
self.logger.info(f"- 在源环境但不在目标环境: {len(diffs['missing_in_target'])}")
self.logger.info(f"- 在目标环境但不在源环境: {len(diffs['missing_in_source'])}")
self.logger.info(f"- 数据不一致: {len(diffs['data_mismatch'])}")
# 2. 生成修复SQL
repair_sql = self._generate_repair_sql(
table, key_columns, columns, diffs
)
# 3. 执行修复(生产环境需审批)
if self.target_env == 'prod':
self.logger.warning("⚠️ 目标环境是生产环境,需要人工确认!")
confirm = input("是否继续执行修复? (yes/no): ")
if confirm.lower() != 'yes':
self.logger.info("修复已取消")
return
return self._execute_repair(repair_sql)
def _generate_repair_sql(self, table: str, key_columns: List[str],
columns: Optional[List[str]], diffs: Dict) -> str:
"""生成修复SQL"""
sql_statements = []
# 处理缺失记录
if diffs['missing_in_target'] and len(diffs['missing_in_target'][0]) == len(key_columns):
# 批量INSERT
insert_values = []
for key in diffs['missing_in_target']:
where_clause = " AND ".join([
f"{col} = '{val}'" for col, val in zip(key_columns, key)
] if isinstance(key, tuple) else [f"{key_columns[0]} = '{key}'"])
insert_values.append(where_clause)
sql_statements.append(f"""
INSERT INTO {table}
SELECT * FROM {self.source_env}.{table}
WHERE ({' OR '.join(insert_values)})
""")
# 处理数据不一致
for key in diffs['data_mismatch']:
if isinstance(key, tuple):
where = " AND ".join([f"{col} = '{val}'" for col, val in zip(key_columns, key)])
else:
where = f"{key_columns[0]} = '{key}'"
if columns:
update_cols = ", ".join([f"{col} = (SELECT {col} FROM {self.source_env}.{table} WHERE {where})"
for col in columns])
sql_statements.append(f"UPDATE {table} SET {update_cols} WHERE {where}")
else:
# 用源数据覆盖
sql_statements.append(f"""
UPDATE {table} t
JOIN {self.source_env}.{table} s ON {' AND '.join([f't.{col} = s.{col}' for col in key_columns])}
SET t.data = s.data
WHERE {where}
""")
return "\n".join(sql_statements)
def _execute_repair(self, sql: str):
"""执行修复SQL"""
self.logger.info(f"执行修复SQL: \n{sql}")
# 写入文件,供手动执行
with open(f'repair_sql_{datetime.now().strftime("%Y%m%d_%H%M%S")}.sql', 'w') as f:
f.write(sql)
return sql
高级数据同步方案
1 双向同步与冲突解决
class BiDirectionalSync:
"""双向数据同步"""
def __init__(self, env1: str, env2: str, db_configs: Dict):
self.env1 = env1
self.env2 = env2
self.repair_tools = {
'env1_to_env2': DataRepairTool(env1, env2, db_configs),
'env2_to_env1': DataRepairTool(env2, env1, db_configs)
}
def sync_with_conflict_resolution(self, table: str, key_columns: List[str],
priority_env: str, conflict_resolver: str = 'latest'):
"""
双向同步并解决冲突
:param table: 表名
:param key_columns: 主键列
:param priority_env: 发生冲突时的优先环境
:param conflict_resolver: 冲突解决策略 ('latest', 'priority', 'merge')
"""
checker = self.repair_tools['env1_to_env2'].checker
diffs = checker.find_differences(self.env1, self.env2, table, key_columns)
# 处理缺失记录
if conflict_resolver == 'latest':
# 以时间戳最新的为准
pass
elif conflict_resolver == 'priority':
# 优先环境覆盖
if priority_env == self.env1:
self.repair_tools['env1_to_env2'].sync_data(table, key_columns)
else:
self.repair_tools['env2_to_env1'].sync_data(table, key_columns)
elif conflict_resolver == 'merge':
# 合并策略
self._merge_data(table, key_columns, diffs)
return diffs
def _merge_data(self, table: str, key_columns: List[str], diffs: Dict):
"""合并数据策略"""
pass
2 基于校验和的批量修复
import zlib
class ChecksumRepair:
"""基于校验和的批量修复"""
def calculate_checksum(self, db_engine, table: str, chunk_size: int = 1000):
"""计算分块校验和"""
checksums = {}
offset = 0
while True:
query = f"SELECT * FROM {table} LIMIT {chunk_size} OFFSET {offset}"
df = pd.read_sql(query, db_engine)
if df.empty:
break
# 计算这一块的校验和
chunk_bytes = df.to_json().encode()
checksum = zlib.crc32(chunk_bytes)
checksums[offset] = checksum
offset += chunk_size
return checksums
def repair_by_chunks(self, source_engine, target_engine, table: str, chunk_size: int = 1000):
"""按块对比并修复"""
source_checksums = self.calculate_checksum(source_engine, table, chunk_size)
target_checksums = self.calculate_checksum(target_engine, table, chunk_size)
for offset, source_checksum in source_checksums.items():
if offset not in target_checksums or source_checksum != target_checksums[offset]:
# 只修复不一致的块
self._repair_chunk(source_engine, target_engine, table, offset, chunk_size)
print(f"已修复偏移 {offset} 的数据块")
def _repair_chunk(self, source_engine, target_engine, table: str, offset: int, chunk_size: int):
"""修复单个数据块"""
query = f"SELECT * FROM {table} LIMIT {chunk_size} OFFSET {offset}"
source_data = pd.read_sql(query, source_engine)
# 替换目标环境的数据块
source_data.to_sql(
f"{table}_temp",
target_engine,
if_exists='replace',
index=False
)
# 执行替换操作(需要根据具体数据库调整)
with target_engine.connect() as conn:
conn.execute(text(f"""
DELETE FROM {table}
WHERE id IN (SELECT id FROM {table}_temp)
"""))
conn.execute(text(f"""
INSERT INTO {table} SELECT * FROM {table}_temp
"""))
conn.commit()
实用工具函数
1 数据对比报告生成
def generate_difference_report(checker: DataConsistencyChecker,
env1: str, env2: str,
table_list: List[str]) -> Dict:
"""生成完整的数据差异报告"""
report = {
'generated_at': datetime.now().isoformat(),
'environments': [env1, env2],
'tables': {}
}
for table in table_list:
diffs = checker.find_differences(env1, env2, table, ['id'])
report['tables'][table] = {
'status': 'synced' if all(len(v) == 0 for v in diffs.values()) else 'diverged',
'differences': diffs
}
# 统计信息
synced_count = sum(1 for t in report['tables'].values() if t['status'] == 'synced')
diverged_count = len(table_list) - synced_count
report['summary'] = {
'total_tables': len(table_list),
'synced': synced_count,
'diverged': diverged_count
}
return report
2 自动修复触发器
def auto_repair_workflow():
"""自动修复工作流示例"""
# 配置
db_configs = {
'dev': {'host': 'localhost', 'port': 3306, 'user': 'root', 'password': '123', 'db': 'test'},
'staging': {'host': 'staging.example.com', 'port': 3306, 'user': 'root', 'password': '456', 'db': 'test'},
'prod': {'host': 'prod.example.com', 'port': 3306, 'user': 'root', 'password': '789', 'db': 'test'}
}
# 重要表列表
critical_tables = ['users', 'orders', 'products']
# 创建修复工具
repair = DataRepairTool('staging', 'prod', db_configs)
# 自动化检查并修复
for table in critical_tables:
try:
print(f"检查表 {table}...")
diffs = repair.checker.find_differences('staging', 'prod', table, ['id'])
if any(len(v) > 0 for v in diffs.values()):
print(f"发现差异,开始修复 {table}")
repair.sync_data(table, ['id'])
# 验证修复
new_diffs = repair.checker.find_differences('staging', 'prod', table, ['id'])
if all(len(v) == 0 for v in new_diffs.values()):
print(f"✅ {table} 修复成功")
else:
print(f"❌ {table} 修复失败,需要人工干预")
except Exception as e:
print(f"处理 {table} 失败: {e}")
if __name__ == "__main__":
auto_repair_workflow()
最佳实践建议
- 备份优先: 执行修复前务必备份目标环境数据
- 灰度发布: 先在非生产环境验证修复脚本
- 断点续传: 大表修复时支持中断后继续
- 审计追踪: 记录所有数据变更操作
- 自动回滚: 修复失败时自动回滚到之前状态
- 性能优化: 使用批量操作而非逐条更新
这个方案可以根据具体需求调整,比如支持更多数据库类型、添加数据校验规则、集成CI/CD流程等。