Python脚本如何自动化修复同步异常数据

wen python案例 27

本文目录导读:

Python脚本如何自动化修复同步异常数据

  1. 基础检测与修复脚本
  2. 带时间戳的高级修复方案
  3. 数据库同步修复脚本
  4. 使用示例
  5. 最佳实践建议

我来分享几种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}")

最佳实践建议

  1. 定期执行:设置cron job定期运行同步脚本
  2. 日志记录:详细记录每次修复操作
  3. 备份数据:修改前备份目标数据
  4. 增量同步:只同步变化的数据
  5. 异常处理:添加重试机制和报警

这个方案可以根据你的具体需求进行调整和扩展。

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