使用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
相关产品推荐
相关产品推荐

