在多进程池中运行异构数据加载作业(限3并发)的技术实现咨询
控制并行数的夜间批处理数据加载方案
嘿,这个需求我太懂了——既要并行提速,又得防着数据库因为负载过高炸锅对吧?刚好Python的标准库就能搞定,给你一步步捋清楚怎么做:
核心思路
你的数据加载任务是典型的IO密集型(大部分时间在等数据库返回数据),用线程池来控制并行数是最优解:线程切换开销小,既能同时跑多个任务提升效率,又能把并行数牢牢限制在3个,避免把数据库冲垮。
具体实现步骤
1. 先把所有加载任务整理好
你已经导入了15个数据源的模块,假设每个模块里都有一个load()函数是实际执行数据加载的入口(如果你的函数名不一样,直接替换成你自己的就行),先把这些函数都放进一个列表里:
import load_BW as BW import load_CS as CS import load_DM as DM import load_ECG as ECG import load_FA as FA import load_FW as FW import load_LB as LB import load_MA as MA import load_MI as MI import load_OE as OE import load_OM as OM import load_PC as PC import load_PP as PP import load_TS as TS # 补上第15个模块,比如叫load_XY import load_XY as XY # 把所有加载函数攒成一个列表,方便后续批量处理 load_tasks = [ BW.load, CS.load, DM.load, ECG.load, FA.load, FW.load, LB.load, MA.load, MI.load, OE.load, OM.load, PC.load, PP.load, TS.load, XY.load ]
2. 用线程池控制并行数
用Python标准库的concurrent.futures.ThreadPoolExecutor,直接把最大并行数设为3就行。同时加上异常处理,免得某个任务失败把整个批处理都搞崩:
from concurrent.futures import ThreadPoolExecutor, as_completed def run_nightly_load(): # 初始化线程池,最多同时跑3个任务 with ThreadPoolExecutor(max_workers=3) as executor: # 把所有任务提交到线程池,顺便记录每个任务的名字方便排查问题 future_map = {executor.submit(task): task.__name__ for task in load_tasks} # 逐个处理完成的任务,不管成功失败都打日志 for future in as_completed(future_map): task_name = future_map[future] try: # 如果你的加载函数有返回值(比如加载了多少条数据),可以在这里获取 result = future.result() print(f"✅ 任务 {task_name} 搞定啦!结果:{result}") except Exception as e: print(f"❌ 任务 {task_name} 翻车了:{str(e)}") # 要是想给失败的任务加重试,就在这里重新提交任务就行,比如: # executor.submit(task) if __name__ == "__main__": run_nightly_load()
3. 可选:给任务加重试机制(应对数据库断开)
既然数据库负载高会断开连接,那给加载任务加个重试逻辑就很有必要了。推荐用tenacity库(先装一下:pip install tenacity),能轻松实现指数退避重试:
from tenacity import retry, stop_after_attempt, wait_exponential, retry_if_exception_type # 替换成你用的数据库对应的连接异常类型,比如MySQL的是mysql.connector.errors.OperationalError import psycopg2 # 写个带重试的包装函数 @retry( stop=stop_after_attempt(3), # 最多重试3次 wait=wait_exponential(multiplier=1, min=2, max=10), # 等待时间2s→4s→8s,最多等10s retry=retry_if_exception_type((psycopg2.OperationalError, ConnectionResetError)) # 只在连接异常时重试 ) def safe_run_load(task_func): return task_func() # 然后修改任务提交的地方,用包装后的函数 future_map = {executor.submit(safe_run_load, task): task.__name__ for task in load_tasks}
为啥选线程池不是进程池?
- 数据加载是IO密集型任务,线程切换的开销比进程小太多,效率更高。
- 进程池会占更多系统资源,15个小任务完全没必要用。
- 要是你的加载任务里有大量CPU密集的操作(比如复杂数据清洗),再换成
ProcessPoolExecutor就行,数据库场景下线程池绝对够用。
最后:把脚本设成每晚自动跑
- Linux/macOS:用
cron定时任务,比如在终端输crontab -e,加一行:
意思是每天凌晨2点自动执行脚本,日志输出到指定文件。0 2 * * * /usr/bin/python3 /你脚本的绝对路径/load_script.py >> /日志路径/load_log.log 2>&1 - Windows:用「任务计划程序」,新建任务,指定Python解释器的路径和你的脚本路径,设置触发时间为每晚就行。
内容的提问来源于stack exchange,提问作者summersmd
相关产品推荐
相关产品推荐

