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

使用throttler与joblib处理重任务时遇"无法解包非可迭代函数对象"错误

限流HTTP请求结合多进程处理CPU密集任务的错误解决

问题场景

需要以2次/秒的速率发送限流HTTP请求,请求返回后用多进程执行CPU密集型计算。使用throttler控制请求速率,joblib实现多进程,但改用parallel(delayed(heavy_calculate)(i))时触发错误:

TypeError: cannot unpack non-iterable function object

错误原因

Parallel对象的调用需要接收可迭代的delayed任务列表,而非单个delayed包装的任务。直接传递单个delayed任务时,Parallel会尝试将其作为可迭代对象遍历,而delayed对象本身不支持迭代,因此触发类型错误。

解决方案

1. 修正Parallel调用方式

将单个delayed任务放入列表中传递给parallel,再从返回的结果列表中取出对应值:

# 替换原来的错误调用
results = parallel([delayed(heavy_calculate)(i)])[0]

2. 完整修正后的代码

import asyncio
from math import sqrt
import time
from joblib import Parallel, delayed
from throttler import throttle


@throttle(rate_limit=1, period=0.5)
async def request(i):
    """模拟限流HTTP请求"""
    time.sleep(0.3)
    print(f"get response: {i} {time.time()}")
    return i


def heavy_calculate(i):
    """模拟CPU密集型计算"""
    n = i
    for _ in range(30_000_000):
        n = sqrt(n)
        n = n ** 2
    print(f"heavy_calculate: {n} {time.time()}")
    return n


async def many_tasks(count: int):
    print(f"=== START === {time.time()}")
    # 初始化多进程池
    with Parallel(n_jobs=-1) as parallel:
        results_list = []
        coros = [request(i) for i in range(count)]
        # 按请求完成顺序处理结果
        for coro in asyncio.as_completed(coros):
            i = await coro
            # 提交单个计算任务到多进程池,取结果列表第一个元素
            calc_result = parallel([delayed(heavy_calculate)(i)])[0]
            results_list.append(calc_result)
    print(f"=== END === {time.time()}")
    print(results_list)

# 启动任务
asyncio.run(many_tasks(10))

额外优化建议

如果希望避免在协程中阻塞事件循环(Parallel多进程调用本身不占用协程线程CPU,但等待结果会阻塞当前协程),可以将多进程计算任务提交到asyncio的线程池执行器中,让事件循环继续处理其他请求:

async def many_tasks(count: int):
    print(f"=== START === {time.time()}")
    # 初始化多进程池
    with Parallel(n_jobs=-1) as parallel:
        results_list = []
        coros = [request(i) for i in range(count)]
        # 获取事件循环的线程池执行器
        loop = asyncio.get_running_loop()
        for coro in asyncio.as_completed(coros):
            i = await coro
            # 在线程池中执行多进程调用,不阻塞事件循环
            calc_result = await loop.run_in_executor(
                None, 
                parallel, 
                [delayed(heavy_calculate)(i)]
            )
            results_list.append(calc_result[0])
    print(f"=== END === {time.time()}")
    print(results_list)

内容的提问来源于stack exchange,提问作者Xaree Lee

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 14:53:17