脚本如何补全迁移缺失的数据项

wen 实用脚本 35

本文目录导读:

脚本如何补全迁移缺失的数据项

  1. 第一步:识别缺失数据
  2. 第二步:编写增量补全脚本(Python + Pandas + SQLAlchemy)
  3. 第三步:处理特定字段补全(而非整行)
  4. 第四步:验证与修复外键约束
  5. 第五步:定时与自动化(可选)
  6. 注意事项与最佳实践

要补全数据迁移过程中缺失的数据项,通常需要结合日志分析数据比对分批回写的策略,以下是一个系统性的解决方案,涵盖从诊断到修复的完整流程,并附有示例脚本(以 Python 和 SQL 为例)。

第一步:识别缺失数据

首先需要明确“缺失”的定义:

  1. 记录缺失:源表有,目标表没有。
  2. 字段缺失:记录存在,但某些字段(如关键属性、时间戳)为空或为默认值。
  3. 不一致缺失:关联的外键在目标表不存在(孤儿记录)。
-- 示例1: 找出在源表存在但目标表不存在的记录(基于主键)
SELECT s.id, s.name, s.data
FROM source_db.source_table s
LEFT JOIN target_db.target_table t ON s.id = t.origin_id
WHERE t.origin_id IS NULL;
-- 示例2: 找出目标表中特定字段为空且应该非空的记录
SELECT id, origin_id, missing_field, created_at
FROM target_db.target_table
WHERE missing_field IS NULL;

第二步:编写增量补全脚本(Python + Pandas + SQLAlchemy)

以下脚本假设:

  • 你已经确定了缺失的 origin_id 列表。
  • 源数据库和目标数据库都是关系型数据库(如 MySQL、PostgreSQL)。
  • 如果数据量很大(例如超过 10 万条),建议使用 批处理 方式。
import pandas as pd
from sqlalchemy import create_engine
import logging
from datetime import datetime
# ---------- 配置数据库连接 ----------
SOURCE_CONN = "mysql+pymysql://user:pass@source_host/source_db?charset=utf8"
TARGET_CONN = "mysql+pymysql://user:pass@target_host/target_db?charset=utf8"
BATCH_SIZE = 5000  # 每批次处理条数
# ---------- 设置日志 ----------
logging.basicConfig(
    level=logging.INFO,
    format='%(asctime)s - %(levelname)s - %(message)s',
    handlers=[
        logging.FileHandler(f'migration_fix_{datetime.now().strftime("%Y%m%d_%H%M%S")}.log'),
        logging.StreamHandler()
    ]
)
def get_missing_ids(engine_source, engine_target, table_name="users"):
    """
    获取缺失记录的 ID 列表
    """
    query = f"""
        SELECT s.id
        FROM {table_name} s
        LEFT JOIN target_db.{table_name} t ON s.id = t.origin_id
        WHERE t.origin_id IS NULL
        LIMIT 1000000;  -- 防止一次性太大,可根据实际情况调整或分批
    """
    # 注意:这里直接在 source 数据库查询目标库,需要跨库权限;更适合的方式是分别在两端查询后比对
    # 更推荐在脚本中拉取两个表的主键集合进行集合运算(见下文替代方案)
    logging.info("正在查询缺失 ID 列表...")
    missing = pd.read_sql(query, engine_source)
    return missing['id'].tolist()
def batch_migrate_missing(engine_source, engine_target, missing_ids, source_table="source_users", target_table="target_users"):
    """
    分批从源表读取缺失记录并插入目标表,同时处理字段映射
    """
    total = len(missing_ids)
    if total == 0:
        logging.info("没有发现缺失记录,任务完成。")
        return
    logging.info(f"开始补全 {total} 条缺失记录...")
    processed = 0
    for i in range(0, total, BATCH_SIZE):
        batch_ids = missing_ids[i:i+BATCH_SIZE]
        try:
            # 1. 从源表提取数据
            query = f"""
                SELECT id, name, email, created_at, updated_at, extra_field
                FROM {source_table}
                WHERE id IN ({','.join(map(str, batch_ids))})
            """
            df_source = pd.read_sql(query, engine_source)
            if df_source.empty:
                logging.warning(f"批次 {i//BATCH_SIZE + 1}: 未找到源数据,跳过")
                continue
            # 2. 字段映射与清洗(示例)
            df_target = df_source.rename(columns={
                'id': 'origin_id',          # 源主键映射为目标表的外键
                'name': 'display_name',
                'extra_field': 'metadata'
            })
            # 添加审计字段
            df_target['imported_at'] = datetime.now()
            df_target['created_at'] = pd.to_datetime(df_target['created_at'], errors='coerce')
            df_target['updated_at'] = pd.to_datetime(df_target['updated_at'], errors='coerce')
            # 3. 写入目标表(选择模式:append / replace / if_exists)
            df_target.to_sql(
                name=target_table,
                con=engine_target,
                if_exists='append',     # 因为目标表已存在且只缺失记录,使用 append
                index=False,
                chunksize=1000,
                method='multi'          # 使用 multi 插入以提升速度
            )
            processed += len(df_target)
            logging.info(f"批次 {i//BATCH_SIZE + 1}: 成功写入 {len(df_target)} 条,已处理 {processed}/{total}")
        except Exception as e:
            logging.error(f"批次 {i//BATCH_SIZE + 1} 失败 (IDs: {batch_ids[:5]}...): {str(e)}")
            # 可以将失败的 ID 记录到文件以便后续重试
            with open("failed_ids.txt", "a") as f:
                f.write('\n'.join(map(str, batch_ids)) + '\n')
            continue
    logging.info(f"补全完成: 成功处理 {processed} 条,失败 {total - processed} 条")
# ---------- 主流程 ----------
if __name__ == "__main__":
    engine_source = create_engine(SOURCE_CONN)
    engine_target = create_engine(TARGET_CONN)
    # 方案A: 如果数据库支持跨库查询
    missing_ids = get_missing_ids(engine_source, engine_target, "users")
    # 方案B: 通用替代——拉取两个表的 ID 集合求差集(适合中小规模)
    # source_ids = set(pd.read_sql("SELECT id FROM source_users", engine_source)['id'])
    # target_ids = set(pd.read_sql("SELECT origin_id FROM target_users WHERE origin_id IS NOT NULL", engine_target)['origin_id'])
    # missing_ids = list(source_ids - target_ids)
    batch_migrate_missing(engine_source, engine_target, missing_ids)

第三步:处理特定字段补全(而非整行)

如果记录已存在,只是某些字段为空,使用 UPDATE 而不是 INSERT。

def fill_missing_fields(engine_source, engine_target, missing_ids):
    """
    更新目标表中指定字段为 NULL 的记录
    """
    columns_to_fill = ['full_address', 'phone', 'company']
    with engine_target.begin() as conn:
        for batch_start in range(0, len(missing_ids), BATCH_SIZE):
            batch_ids = missing_ids[batch_start:batch_start+BATCH_SIZE]
            # 从源表获取待更新的数据
            src_query = f"""
                SELECT id, full_address, phone, company
                FROM source_users
                WHERE id IN ({','.join(map(str, batch_ids))})
            """
            df_src = pd.read_sql(src_query, engine_source)
            # 逐行更新(也可以构建批量 CASE WHEN 语句,但复杂且容易出错)
            for _, row in df_src.iterrows():
                update_parts = []
                params = {}
                for col in columns_to_fill:
                    if pd.notna(row[col]):
                        update_parts.append(f"{col} = :{col}")
                        params[col] = row[col]
                if update_parts:
                    set_clause = ", ".join(update_parts)
                    sql = f"""
                        UPDATE target_users
                        SET {set_clause}
                        WHERE origin_id = :origin_id
                    """
                    params['origin_id'] = row['id']
                    conn.execute(sql, params)
    logging.info(f"已补全 {len(missing_ids)} 条记录的字段值")

第四步:验证与修复外键约束

如果缺失数据导致外键冲突(用户表缺少父级 manager_id 对应的记录):

def fix_orphan_references(engine_source, engine_target, relation_table="orders"):
    """
    找出目标表中引用到缺失父记录的子记录,并补全父记录
    """
    # 查找缺失的 parent_id
    query = f"""
        SELECT DISTINCT t.parent_id
        FROM target_db.{relation_table} t
        LEFT JOIN target_db.parent_table p ON t.parent_id = p.id
        WHERE p.id IS NULL AND t.parent_id IS NOT NULL
    """
    orphan_parent_ids = pd.read_sql(query, engine_target)['parent_id'].tolist()
    if orphan_parent_ids:
        logging.info(f"发现 {len(orphan_parent_ids)} 个缺失的父记录,准备补全...")
        # 复用 batch_migrate_missing 函数,但 source_table 改为 parent_table
        batch_migrate_missing(engine_source, engine_target, orphan_parent_ids, 
                              source_table="source_parent_table", 
                              target_table="parent_table")
    else:
        logging.info("没有孤儿记录,外键关系完整。")

第五步:定时与自动化(可选)

将以上脚本包装为可配置的 CLI 工具或定时任务:

# cron job 示例(每天凌晨2点执行补全)
0 2 * * * cd /data/migration && python3 fix_missing.py --config config.yaml --log-file fix_$(date +\%Y\%m\%d).log

注意事项与最佳实践

场景 建议
数据量极大(亿级) 使用 Spark 或数据库端硬解析(如 MySQL Event Scheduler + 存储过程)
需要实时性 基于日志解析(如 Debezium + Kafka)实现增量同步,脚本只做初次补全
源表数据在变化 先锁定源表快照(或使用变更日志表),确保补全时源数据一致
目标表有唯一约束 使用 ON DUPLICATE KEY UPDATE(MySQL)或 ON CONFLICT DO NOTHING(PG)
字段类型不匹配 在映射时显式转换(如 Python 中 astype(str) 或 SQL 中 CAST

脚本是基于通用流程的模板,实际使用时请根据你的数据库类型、表结构、数据量大小进行相应调整,如果缺失数据是因为迁移流程本身的 bug,建议同时修复迁移工具,而不是只靠补全脚本。

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