服务器端Python代码多核并行化实现方案与示例咨询
Python 构造器并行化实现方案
一、CPU核心数等硬件信息获取方法
不需要安装第三方依赖,直接用Python标准库就能跨平台获取核心信息,适配Windows开发机和Linux服务器:
import os import multiprocessing # 获取系统逻辑核心数(包含超线程虚拟核心) logical_core_num = os.cpu_count() # 获取可用物理核心数 if os.name == "nt": # Windows系统 physical_core_num = multiprocessing.cpu_count() // 2 else: # Linux/Mac系统,容器部署时也能正确识别分配给当前进程的核心配额 physical_core_num = len(os.sched_getaffinity(0)) print(f"逻辑核心数: {logical_core_num}, 可用物理核心数: {physical_core_num}")
开发调试阶段可以手动指定核心分配参数,不需要硬编码读取硬件信息,方便在少核的本地设备上调通流程。
二、group_a任务调度与核心资源分配实现
整体用层级多进程+CPU亲和性绑定的方案实现,完全匹配执行流程要求:每轮等group_a全量执行完成后再顺序跑group_b,group_a每个任务独占分配到的核心组,内部可以做二次并行,不会出现资源抢占。
基础调度逻辑
如果group_a单个任务不需要内部二次并行,直接用ProcessPoolExecutor提交所有group_a任务,调用wait()阻塞等待所有任务返回结果后,再顺序执行group_b的函数即可。
二次并行+核心绑定实现(匹配12核3任务每任务4核的场景)
实现逻辑:
- 启动任务前按总核心数和group_a任务数均分核心,生成每个任务对应的核心号列表
- 第一层启动和group_a任务数相等的进程,每个进程启动时先通过系统调用绑定到分配给自己的核心组
- 每个group_a任务内部启动和分配核心数相等的子进程池,跑内部并行子任务,因为已经做了核心绑定,子进程只会在分配的核心上运行,不会抢占其他任务的资源
- 所有group_a任务返回结果后,主进程顺序执行group_b的所有任务,完成后进入下一轮迭代
可直接参考的实现代码:
import os from concurrent.futures import ProcessPoolExecutor, wait, as_completed # -------------------------- # 业务函数模拟,替换成实际的文件对应函数即可 # -------------------------- def sub_task_handler(param): """group_a任务内部的子任务逻辑""" return param * 2 def run_group_a1(task_params, allocated_cores): # 绑定当前进程及所有子进程到分配的核心组 os.sched_setaffinity(0, allocated_cores) # 内部开和核心数相等的进程池做二次并行 with ProcessPoolExecutor(max_workers=len(allocated_cores)) as inner_pool: sub_futures = [inner_pool.submit(sub_task_handler, p) for p in task_params] sub_results = [f.result() for f in as_completed(sub_futures)] return {"task_id": "a1", "data": sub_results} def run_group_a2(task_params, allocated_cores): os.sched_setaffinity(0, allocated_cores) with ProcessPoolExecutor(max_workers=len(allocated_cores)) as inner_pool: sub_futures = [inner_pool.submit(sub_task_handler, p) for p in task_params] sub_results = [f.result() for f in as_completed(sub_futures)] return {"task_id": "a2", "data": sub_results} def run_group_a3(task_params, allocated_cores): os.sched_setaffinity(0, allocated_cores) with ProcessPoolExecutor(max_workers=len(allocated_cores)) as inner_pool: sub_futures = [inner_pool.submit(sub_task_handler, p) for p in task_params] sub_results = [f.result() for f in as_completed(sub_futures)] return {"task_id": "a3", "data": sub_results} def run_group_b1(params): return {"task_id": "b1", "data": "顺序执行任务1结果"} def run_group_b2(params): return {"task_id": "b2", "data": "顺序执行任务2结果"} # -------------------------- # 主迭代流程 # -------------------------- if __name__ == "__main__": # 核心资源分配 total_available_cores = len(os.sched_getaffinity(0)) if os.name != "nt" else os.cpu_count() group_a_task_count = 3 cores_per_a_task = total_available_cores // group_a_task_count # 生成每个group_a任务对应的核心组,12核场景下为[0,1,2,3]、[4,5,6,7]、[8,9,10,11] core_alloc_plan = [ list(range(i*cores_per_a_task, (i+1)*cores_per_a_task)) for i in range(group_a_task_count) ] # 按迭代轮次执行 total_iter_rounds = 10 # 替换成实际需要的迭代次数 for round_idx in range(total_iter_rounds): print(f"当前执行第{round_idx}轮迭代") # 加载当前轮次各任务的独立参数 a1_params = [1,2,3,4] a2_params = [5,6,7,8] a3_params = [9,10,11,12] # 并行执行所有group_a任务 a_task_list = [run_group_a1, run_group_a2, run_group_a3] a_param_list = [a1_params, a2_params, a3_params] a_results = [] with ProcessPoolExecutor(max_workers=group_a_task_count) as outer_pool: outer_futures = [ outer_pool.submit(task, param, cores) for task, param, cores in zip(a_task_list, a_param_list, core_alloc_plan) ] # 阻塞等待所有group_a任务执行完成 wait(outer_futures) for f in outer_futures: a_results.append(f.result()) # 此处写group_a返回值的统一处理逻辑 print("group_a执行完成,返回结果:", a_results) # 严格顺序执行group_b任务 b1_params = "param_for_b1" b2_params = "param_for_b2" b1_res = run_group_b1(b1_params) b2_res = run_group_b2(b2_params) # 此处写group_b返回值的统一处理逻辑 print("group_b执行完成,返回结果:", b1_res, b2_res)
实现注意事项
- 所有多进程启动逻辑必须放在
if __name__ == "__main__":块内,Windows环境下如果不做这个限制,会递归触发子进程启动直接报错。 - 如果group_a内部的子任务是IO密集型(比如读写大文件、调用外部接口),可以把内部的
ProcessPoolExecutor替换为ThreadPoolExecutor,不需要做核心绑定,线程切换开销更低,性能更好。 - 本地Windows开发调试时如果设备核心数不足,可以手动指定
cores_per_a_task参数(比如设为1),不需要完全匹配服务器核心配置,先把流程跑通再上服务器调优。 - 用
os.sched_setaffinity做核心绑定是系统级的资源限制,比单纯在代码里控制进程数更可靠,能避免不同任务的进程抢占同一核心导致的性能波动。
内容的提问来源于stack exchange,提问作者MIKKI MOUSE
相关产品推荐
相关产品推荐

