Python脚本如何修复多环境数据不一致

wen python案例 33

本文目录导读:

Python脚本如何修复多环境数据不一致

  1. 核心方案设计
  2. 高级数据同步方案
  3. 实用工具函数
  4. 最佳实践建议

我来详细说明如何使用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()

最佳实践建议

  1. 备份优先: 执行修复前务必备份目标环境数据
  2. 灰度发布: 先在非生产环境验证修复脚本
  3. 断点续传: 大表修复时支持中断后继续
  4. 审计追踪: 记录所有数据变更操作
  5. 自动回滚: 修复失败时自动回滚到之前状态
  6. 性能优化: 使用批量操作而非逐条更新

这个方案可以根据具体需求调整,比如支持更多数据库类型、添加数据校验规则、集成CI/CD流程等。

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