如何使用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_id | city_name | is_available | is_running | is_success |
|---|---|---|---|---|
| 1 | New York | 0 | 0 | 0 |
| 2 | Los Angeles | 0 | 0 | 0 |
| 4 | Houston | 0 | 0 | 0 |
| 5 | Phoenix | 0 | 0 | 0 |
| 3 | Chicago | 1 | 0 | 0 |
现有数据迁移代码
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
相关产品推荐
相关产品推荐

