如何在Python中并行发起多个gRPC Unary API调用?
嘿,这个问题我刚好折腾过!其实Python里实现gRPC Unary调用的并行有好几种靠谱方案,你之前觉得asyncio不适用可能是没找对正确的姿势——gRPC其实有专门的异步客户端支持asyncio的,下面给你详细说说几个可行的路子:
方案1:gRPC异步客户端 + asyncio(推荐,IO密集场景最优)
这个方案是性能最高的,因为异步IO的开销极小,特别适合批量发起IO密集的gRPC调用。
首先你需要生成异步gRPC客户端代码,用grpcio-tools的时候要加上--grpc_asyncio_python_out参数,示例命令:
python -m grpc_tools.protoc -I./proto --python_out=./gen --grpc_asyncio_python_out=./gen your_service.proto
然后写Python代码实现并发调用:
import asyncio import grpc from gen import your_service_pb2 from gen import your_service_pb2_grpc async def single_grpc_call(stub, request): try: response = await stub.YourUnaryMethod(request) return (request, response, None) except grpc.aio.AioRpcError as e: return (request, None, e) async def main(): # 准备你的批量请求列表 requests = [your_service_pb2.YourRequest(param=i) for i in range(10)] # 创建异步gRPC通道和stub,用async with自动管理生命周期 async with grpc.aio.insecure_channel('localhost:50051') as channel: stub = your_service_pb2_grpc.YourServiceStub(channel) # 把所有调用任务打包,用asyncio.gather并发执行 tasks = [single_grpc_call(stub, req) for req in requests] results = await asyncio.gather(*tasks) # 遍历处理结果 for req, resp, err in results: if err: print(f"请求 {req.param} 失败: {err.details()}") else: print(f"请求 {req.param} 成功: {resp.result}") if __name__ == "__main__": asyncio.run(main())
方案2:多线程 + 同步gRPC客户端
如果不想折腾异步代码,用线程池是个更简单的选择——gRPC的同步客户端是线程安全的,完全可以放心在多线程环境下使用。
示例代码:
import grpc from concurrent.futures import ThreadPoolExecutor from gen import your_service_pb2 from gen import your_service_pb2_grpc def single_grpc_call(stub, request): try: response = stub.YourUnaryMethod(request) return (request, response, None) except grpc.RpcError as e: return (request, None, e) def main(): requests = [your_service_pb2.YourRequest(param=i) for i in range(10)] # 创建同步gRPC通道和stub with grpc.insecure_channel('localhost:50051') as channel: stub = your_service_pb2_grpc.YourServiceStub(channel) # 用线程池并发执行调用,max_workers根据你的并发需求调整 with ThreadPoolExecutor(max_workers=10) as executor: futures = [executor.submit(single_grpc_call, stub, req) for req in requests] # 逐个获取结果 for future in futures: req, resp, err = future.result() if err: print(f"请求 {req.param} 失败: {err.details()}") else: print(f"请求 {req.param} 成功: {resp.result}") if __name__ == "__main__": main()
这个方案的优势是代码逻辑同步,容易理解和调试,适合中小规模的并发场景。
方案3:多进程(仅适合CPU密集型场景)
如果你的gRPC调用还附带大量CPU密集的预处理/后处理任务,多进程可以利用多核CPU的优势。不过要注意:每个进程需要单独创建gRPC通道,不能跨进程共享通道对象。
示例代码:
import grpc from concurrent.futures import ProcessPoolExecutor from gen import your_service_pb2 from gen import your_service_pb2_grpc def single_grpc_call(request): # 每个进程单独创建通道和stub with grpc.insecure_channel('localhost:50051') as channel: stub = your_service_pb2_grpc.YourServiceStub(channel) try: response = stub.YourUnaryMethod(request) return (request, response, None) except grpc.RpcError as e: return (request, None, e) def main(): requests = [your_service_pb2.YourRequest(param=i) for i in range(10)] # 进程数建议和CPU核心数匹配 with ProcessPoolExecutor(max_workers=4) as executor: results = list(executor.map(single_grpc_call, requests)) for req, resp, err in results: if err: print(f"请求 {req.param} 失败: {err.details()}") else: print(f"请求 {req.param} 成功: {resp.result}") if __name__ == "__main__": main()
这个方案的开销比前两个大,所以只建议在确实有CPU密集任务的时候使用。
一些额外注意事项
- 异步方案里一定要用
async with管理通道,避免资源泄漏; - 线程池的
max_workers不要设得过大,避免占用过多服务器资源,一般参考gRPC服务器的并发能力调整; - 不管用哪种方案,都要做好异常处理,避免单个调用失败影响整个批量任务。
内容的提问来源于stack exchange,提问作者rajan sthapit
相关产品推荐
相关产品推荐

