You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何使用Python实现跨表并行数据迁移

表间数据并行迁移问题及解决方案

问题描述

需要实现表间数据的并行迁移:比如城市A、B在凌晨1点启动迁移,C在1点10分启动,单城市迁移耗时约30分钟。但现有代码无法实现并行,C的迁移必须等A、B完成后才会开始,需解决该问题。

以下是测试数据生成代码、city_flagging表结构及现有迁移代码:

测试数据生成代码

import pandas as pd 
import random

# 定义列名
columns = ['ID', 'Name', 'City', 'City_ID', 'Revenue']

# 城市及对应ID
cities = {
    'New York': 1,
    'Los Angeles': 2,
    'Chicago': 3,
    'Houston': 4,
    'Phoenix': 5
}

# 生成100万条测试数据
data = {
    'ID': range(1, 1000001),
    'Name': [f"Customer_{i}" for i in range(1, 1000001)],
    'City': [random.choice(list(cities.keys())) for _ in range(1000000)],
    'Revenue': [round(random.uniform(20.0, 500.0), 2) for _ in range(1000000)]
}

# 根据City列生成City_ID
data['City_ID'] = [cities[city] for city in data['City']]

# 创建DataFrame
df = pd.DataFrame(data)

# 显示前几行
print(df.head())

# 可选:保存为CSV文件
df.to_csv('D:/INSAN/Document Project/dummy_data.csv', index=False)

city_flagging表结构

city_idcity_nameis_availableis_runningis_success
1New York000
2Los Angeles000
4Houston000
5Phoenix000
3Chicago100

现有数据迁移代码

import pyodbc
import datetime
import time
import concurrent.futures

# 已迁移城市ID集合
migrated_city_ids = set()

# 从数据库获取is_available=1的城市
def check_flagging_status():
    conn = pyodbc.connect('Driver={PostgreSQL ODBC Driver(UNICODE)};server=localhost;port=5433;database=postgres;uid=postgres;pwd=postgres')
    cursor = conn.cursor()

    query = "SELECT * FROM city_flagging WHERE is_available = 1 AND is_running = 0"
    cursor.execute(query)

    available_cities = cursor.fetchall()

    cursor.close()
    conn.close()

    return available_cities

# 更新city_flagging表的is_running和is_success状态
def update_flagging_status(city_id, is_running=None, is_success=None):
    conn = pyodbc.connect('Driver={PostgreSQL ODBC Driver(UNICODE)};server=localhost;port=5433;database=postgres;uid=postgres;pwd=postgres')
    cursor = conn.cursor()

    try:
        if is_running is not None:
            query = "UPDATE city_flagging SET is_running = ? WHERE city_id = ?"
            cursor.execute(query, (is_running, city_id))
        
        if is_success is not None:
            query = "UPDATE city_flagging SET is_success = ? WHERE city_id = ?"
            cursor.execute(query, (is_success, city_id))
        
        conn.commit()
    except Exception as e:
        print(f"更新城市{city_id}状态时出错: {e}")
    finally:
        cursor.close()
        conn.close()

# 单个城市数据迁移函数
def migrate_data(city_id, city_name):
    if city_id in migrated_city_ids:
        print(f"城市{city_id}已迁移过,跳过。")
        return

    # 更新is_running为1(迁移中)
    update_flagging_status(city_id, is_running=1)

    date = datetime.datetime.now().date()
    start_datetime = datetime.datetime.now()

    conn_old = pyodbc.connect('Driver={PostgreSQL ODBC Driver(UNICODE)};server=localhost;port=5433;database=postgres;uid=postgres;pwd=postgres')
    cursor_old = conn_old.cursor()

    conn_new = pyodbc.connect('Driver={PostgreSQL ODBC Driver(UNICODE)};server=localhost;port=5433;database=postgres;uid=postgres;pwd=postgres')
    cursor_new = conn_new.cursor()

    try:
        query = f"SELECT * FROM dummy_data WHERE city_id = '{city_id}'"
        cursor_old.execute(query)
        transaksi_data = cursor_old.fetchall()

        print(f"城市{city_id}的数据迁移中...")
        insert_query = "INSERT INTO dummy_data_1 (id, name, city, revenue, city_id) VALUES (?, ?, ?, ?, ?)"
        cursor_new.executemany(insert_query, transaksi_data)

        conn_new.commit()

        end_datetime = datetime.datetime.now()

        log_query = """
        INSERT INTO migration_log (date, city_id, city_name, start_datetime, end_datetime) 
        VALUES (?, ?, ?, ?, ?)
        """
        cursor_new.execute(log_query, (date, city_id, city_name, start_datetime, end_datetime))
        conn_new.commit()

        # 更新is_success为1(迁移成功)
        update_flagging_status(city_id, is_success=1)

        migrated_city_ids.add(city_id)
        print(f"城市{city_id}的数据迁移完成。")

    except Exception as e:
        print(f"迁移城市{city_id}数据时出错: {e}")
    
    finally:
        cursor_old.close()
        conn_old.close()
        cursor_new.close()
        conn_new.close()

# 启动迁移
def start_migration(executor):
    city_list = check_flagging_status()

    if not city_list:
        print("没有可迁移的城市。")
        return False

    for city in city_list:
        city_id = city[0]
        city_name = city[1]
        executor.submit(migrate_data, city_id, city_name)

    return True

# 周期性检查迁移状态
def continuous_check(interval_seconds=2, max_no_migration_attempts=3):
    no_migration_attempts = 0
    # 全局线程池,不随每次start_migration销毁
    with concurrent.futures.ThreadPoolExecutor(max_workers=5) as executor:
        while no_migration_attempts < max_no_migration_attempts:
            print("检查城市标记状态...")
            migration_started = start_migration(executor)

            if migration_started:
                no_migration_attempts = 0
            else:
                no_migration_attempts += 1

            if no_migration_attempts >= max_no_migration_attempts:
                print("多次检查无可用迁移城市,停止进程。")
                break

            time.sleep(interval_seconds)

# 启动周期性检查
continuous_check(1, 3)

问题原因及解决方案

核心问题

原代码中start_migration函数使用with concurrent.futures.ThreadPoolExecutor(...),with语句会阻塞直到所有提交的线程执行完毕才会退出,导致continuous_check的下一次循环(检查是否有新城市可迁移)必须等上一批迁移线程全部完成才能执行,因此C的迁移无法在A、B运行时启动。

解决方案

  • 全局线程池复用:将线程池的创建移到continuous_check函数内,用with包裹整个循环过程,这样线程池会持续存在,每次检查到新城市时直接提交任务到现有线程池,无需等待之前的线程完成。
  • 线程安全处理:migrated_city_ids是多线程共享的集合,需要加锁避免并发修改问题,添加threading.Lock来保护对该集合的操作。
  • 定时启动逻辑优化:如果需要按指定时间启动不同城市的迁移,可以在city_flagging表中新增schedule_time字段,在check_flagging_status函数中加入时间判断,只有当前时间大于等于schedule_time且满足其他条件的城市才会被选中迁移。

修改后的关键代码片段

线程安全的迁移城市集合

import threading

migrated_city_ids = set()
migrated_lock = threading.Lock()

# 在migrate_data函数中修改集合时加锁
with migrated_lock:
    if city_id in migrated_city_ids:
        print(f"城市{city_id}已迁移过,跳过。")
        return

# ...

with migrated_lock:
    migrated_city_ids.add(city_id)

新增schedule_time字段后的检查逻辑

假设city_flagging表新增schedule_time字段(类型为TIMESTAMP),修改check_flagging_status函数:

def check_flagging_status():
    conn = pyodbc.connect('Driver={PostgreSQL ODBC Driver(UNICODE)};server=localhost;port=5433;database=postgres;uid=postgres;pwd=postgres')
    cursor = conn.cursor()

    query = """
        SELECT * FROM city_flagging 
        WHERE is_available = 1 AND is_running = 0 
        AND schedule_time <= CURRENT_TIMESTAMP
    """
    cursor.execute(query)

    available_cities = cursor.fetchall()

    cursor.close()
    conn.close()

    return available_cities

这样就可以通过设置schedule_time来控制不同城市的启动时间,同时线程池持续运行,新的迁移任务会立即被调度执行,实现真正的并行迁移。

内容的提问来源于stack exchange,提问作者Insan Cahya

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.18 07:14:54