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

