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

服务器端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核的场景)

实现逻辑:

  1. 启动任务前按总核心数和group_a任务数均分核心,生成每个任务对应的核心号列表
  2. 第一层启动和group_a任务数相等的进程,每个进程启动时先通过系统调用绑定到分配给自己的核心组
  3. 每个group_a任务内部启动和分配核心数相等的子进程池,跑内部并行子任务,因为已经做了核心绑定,子进程只会在分配的核心上运行,不会抢占其他任务的资源
  4. 所有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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 17:45:37