Python multiprocessing实现Cassandra并行请求调用wait()出现pickle错误
问题触发原因
你遇到的报错本质是Python多进程间参数传递的序列化限制,具体原因有两点:
- Python的
multiprocessing库在给子进程传递任务参数时,会默认使用pickle序列化所有传入的对象,而你传入pool.map_async的参数是session.execute_async返回的ResponseFuture对象,该对象内部绑定了Cassandra连接会话、底层线程锁(_thread.RLock)、socket连接等不可序列化的资源,因此触发了序列化失败的报错。 - 额外补充:即使你绕过了序列化限制,主进程创建的Cassandra Session对象本身也不是进程安全的,跨进程共享连接资源会导致底层连接池、网络请求混乱,出现不可预期的错误。而且你已经使用了Cassandra驱动自带的异步查询接口,本身就可以实现查询并行,无需额外引入多进程增加复杂度。
解决方案
提供三种可落地的方案,你可以根据自己的业务场景选择:
方案1:直接使用异步API,去掉多余的多进程逻辑
这是改动最小的方案,Cassandra的异步查询本身已经支持并行执行,不需要多进程介入:
from cassandra.cluster import Cluster import json from functools import reduce cluster = Cluster(['xyz']) session = cluster.connect() query = session.prepare('SELECT stuff') session.default_timeout = 600000 session.default_fetch_size = 100 # 提交所有异步查询 queries = [ session.execute_async(query, ['2021-10-19'] + [i]) for i in range(32) ] # 直接等待所有异步查询返回结果,再执行compute和聚合 res = [compute(query_future.result()) for query_future in queries] final_response = reduce(aggregate, res) resp = json.dumps(final_response, sort_keys=True, indent=4).encode("utf-8") print("RESPONSE", resp)
方案2:改用多线程实现并行
如果你确实需要并行执行compute逻辑(比如compute是CPU密集型任务),可以改用多线程,线程间共享内存无需序列化对象,且Cassandra Session本身是线程安全的:
from concurrent.futures import ThreadPoolExecutor # ... 省略前面创建session、提交异步查询的代码 ... with ThreadPoolExecutor(max_workers=32) as pool: res = pool.map(lambda f: compute(f.result()), queries) # ... 省略后面聚合的代码 ...
方案3:坚持用多进程的情况下,让子进程自己创建Cassandra连接
如果你一定要用多进程,就不要在主进程创建查询和连接,只把查询参数传给子进程,子进程自己初始化连接执行查询:
def worker(i, date_str): # 子进程内部自己创建连接 cluster = Cluster(['xyz']) session = cluster.connect() query = session.prepare('SELECT stuff') session.default_timeout = 600000 session.default_fetch_size = 100 result = session.execute(query, [date_str] + [i]) processed = compute(result) session.shutdown() cluster.shutdown() return processed # 主进程逻辑 if __name__ == "__main__": import multiprocessing as mp from functools import reduce import json pool = mp.Pool(32) inter_obj = pool.starmap_async(worker, [(i, '2021-10-19') for i in range(32)]) inter_obj.wait() res = inter_obj.get() pool.close() pool.join() final_response = reduce(aggregate, res) resp = json.dumps(final_response, sort_keys=True, indent=4).encode("utf-8") print("RESPONSE", resp)
注意:该方案每个子进程都会创建独立的Cassandra连接,请注意控制并发进程数,避免连接数过多压垮Cassandra集群。
内容的提问来源于stack exchange,提问作者Shivangi Singh
相关产品推荐
相关产品推荐

