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

使用multiprocessing.Queue时遇TypeError:无法pickle _thread.RLock对象

解决TypeError: can't pickle _thread.RLock objects问题

错误根源是你试图将ThreadPoolExecutor.submit()返回的Future对象存入多进程Queue,而Future对象内部包含_thread.RLock这类无法被pickle序列化的线程锁结构,多进程队列传递对象必须经过序列化,因此触发报错。

解决方法

核心思路是不要跨进程传递Future对象,改为传递请求结果或调整执行逻辑,以下是两种可行方案:

方案1:在Worker内处理请求结果,仅传递最终结果到队列

如果需要收集请求响应,直接在worker线程里执行请求并拿到结果后再入队,而非传递Future:

from multiprocessing import Queue, Process, cpu_count
import time
from ratelimiter import RateLimiter
from concurrent import futures

q = Queue()

def predict(req):
    # 模拟请求逻辑
    time.sleep(0.001)
    return f"Response for {req}"

def worker(q, req, process_id, time_end):
    pool = futures.ThreadPoolExecutor(max_workers=1000)
    rate_limiter = RateLimiter(max_calls=1000, period=1)
    
    def handle_result(future):
        # 处理线程执行结果并放入队列
        result = future.result()
        q.put(result)
    
    while time.time() < time_end:
        with rate_limiter:
            future = pool.submit(predict, req)
            future.add_done_callback(handle_result)

if __name__ == "__main__":
    num_processes = cpu_count()
    req = "test_request"
    time_end = time.time() + 10  # 运行10秒
    processes = []
    
    for i in range(num_processes):
        p = Process(target=worker, args=(q, req, i, time_end))
        p.start()
        processes.append(p)
    
    # 示例:从队列取出结果
    try:
        while True:
            result = q.get(timeout=15)
            print(result)
    except:
        pass
    
    for p in processes:
        p.join()

方案2:移除多进程队列,直接在Worker内完成请求(无需跨进程传递)

如果不需要统一收集结果,完全可以去掉队列,让每个worker独立执行请求逻辑:

from multiprocessing import Process, cpu_count
import time
from ratelimiter import RateLimiter
from concurrent import futures

def predict(req):
    # 模拟请求逻辑
    time.sleep(0.001)
    print(f"Process {req} completed")

def worker(req, process_id, time_end):
    pool = futures.ThreadPoolExecutor(max_workers=1000)
    rate_limiter = RateLimiter(max_calls=1000, period=1)
    
    while time.time() < time_end:
        with rate_limiter:
            pool.submit(predict, f"{req}_{process_id}")

if __name__ == "__main__":
    num_processes = cpu_count()
    req = "test_request"
    time_end = time.time() + 10
    processes = []
    
    for i in range(num_processes):
        p = Process(target=worker, args=(req, i, time_end))
        p.start()
        processes.append(p)
    
    for p in processes:
        p.join()

额外注意事项

  • 每个worker里的RateLimiter是独立的,总请求速率为num_processes * max_calls/period,如果需要全局统一速率,要把RateLimiter改为进程间共享的实现(比如用Redis做分布式限流)。
  • ThreadPoolExecutor的max_workers=1000可能过高,IO密集型场景建议根据实际调整(比如200-500),避免资源耗尽。
  • 必须在if __name__ == "__main__":下启动进程,这是Windows系统多进程的强制要求,也能避免Linux下的重复初始化问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 09:57:41