ZeroMQ运行一段时间后无响应问题排查与解决
解决ZeroMQ异步通信长期运行后无响应问题
我用ZeroMQ Python库实现API与进程的异步通信,API作为客户端向进程发起请求。前几次请求正常,但运行一段时间后服务无响应,IPC无返回结果,curl调用API时请求超时。推测问题出在连接管理或阻塞,以下是相关代码及修复方案:
原代码
server.py
from zmq import POLLIN, ROUTER from zmq.asyncio import Context, Poller from zmq.auth.asyncio import AsyncioAuthenticator from json import loads, dumps async def start(self): """Starts the zmq server asyncronously and handles incoming requests""" context = Context() auth = AsyncioAuthenticator(context) auth.start() auth.configure_plain(domain="*", passwords={"user": IPC_TOKEN}) auth.allow("127.0.0.1") socket = context.socket(ROUTER) socket.plain_server = True socket.bind("tcp://*:5555") poller = Poller() poller.register(socket, POLLIN) while True: socks = dict(await poller.poll()) if socket in socks and socks[socket] == POLLIN: message = await socket.recv_multipart() identity, request = message decoded = loads(request.decode()) res = await getattr(self, decoded["route"])(decoded["data"]) if res: await socket.send_multipart([identity, dumps(res).encode()]) else: await socket.send_multipart([identity, b'{"status":"ok"}'])
client.py
from zmq import DEALER, POLLIN from zmq.asyncio import Context, Poller import json import uuid async def make_request(route: str, data: dict) -> dict: context = Context.instance() socket = context.socket(DEALER) socket.identity = uuid.uuid4().hex.encode('utf-8') socket.plain_username = b"user" socket.plain_password = IPC_TOKEN.encode("UTF-8") socket.connect("tcp://localhost:5555") request = json.dumps({"route": route, "data": data}).encode('utf-8') socket.send(request) poller = Poller() poller.register(socket, POLLIN) while True: events = dict(await poller.poll()) if socket in events and events[socket] == POLLIN: multipart = json.loads((await socket.recv_multipart())[0].decode()) socket.close() context.term() return multipart
问题分析
- Client端资源泄漏:每次请求都销毁全局单例
Context,多次term()会导致底层资源未正确释放,后续连接创建失败 - 无超时处理:Client的
poll()未设置超时,若服务器无响应会无限阻塞 - Server异常崩溃风险:请求处理逻辑未捕获异常,单个请求出错会导致整个服务循环中断
- 僵尸连接积累:ROUTER套接字未配置心跳,客户端异常断开后,服务器会保留无效连接身份,耗尽资源
修复后代码
server.py
from zmq import POLLIN, ROUTER, TCP_KEEPALIVE, TCP_KEEPALIVE_IDLE, TCP_KEEPALIVE_INTVL from zmq.asyncio import Context, Poller from zmq.auth.asyncio import AsyncioAuthenticator from json import loads, dumps async def start(self): """Starts the zmq server asyncronously and handles incoming requests""" context = Context() auth = AsyncioAuthenticator(context) auth.start() auth.configure_plain(domain="*", passwords={"user": IPC_TOKEN}) auth.allow("127.0.0.1") socket = context.socket(ROUTER) socket.plain_server = True # 配置TCP心跳,自动清理僵尸连接 socket.setsockopt(TCP_KEEPALIVE, 1) socket.setsockopt(TCP_KEEPALIVE_IDLE, 300) socket.setsockopt(TCP_KEEPALIVE_INTVL, 60) socket.bind("tcp://*:5555") poller = Poller() poller.register(socket, POLLIN) while True: socks = dict(await poller.poll()) if socket in socks and socks[socket] == POLLIN: try: message = await socket.recv_multipart() identity, request = message decoded = loads(request.decode()) res = await getattr(self, decoded["route"])(decoded["data"]) if res: await socket.send_multipart([identity, dumps(res).encode()]) else: await socket.send_multipart([identity, b'{"status":"ok"}']) except Exception as e: # 捕获异常避免服务崩溃,可替换为日志记录 print(f"请求处理出错: {str(e)}") await socket.send_multipart([identity, b'{"status":"error","msg":"服务器内部错误"}'])
client.py
from zmq import DEALER, POLLIN, TCP_KEEPALIVE, TCP_KEEPALIVE_IDLE, TCP_KEEPALIVE_INTVL from zmq.asyncio import Context, Poller import json import uuid # 全局复用Context,避免频繁创建销毁 _global_context = None def get_global_context(): global _global_context if _global_context is None or _global_context.closed: _global_context = Context.instance() return _global_context async def make_request(route: str, data: dict) -> dict: context = get_global_context() socket = context.socket(DEALER) socket.identity = uuid.uuid4().hex.encode('utf-8') socket.plain_username = b"user" socket.plain_password = IPC_TOKEN.encode("UTF-8") # 配置TCP心跳维持连接健康 socket.setsockopt(TCP_KEEPALIVE, 1) socket.setsockopt(TCP_KEEPALIVE_IDLE, 300) socket.setsockopt(TCP_KEEPALIVE_INTVL, 60) socket.connect("tcp://localhost:5555") request = json.dumps({"route": route, "data": data}).encode('utf-8') await socket.send(request) # 改为异步send契合异步模型 poller = Poller() poller.register(socket, POLLIN) try: # 设置5秒超时,避免无限阻塞 events = dict(await poller.poll(5000)) if socket in events and events[socket] == POLLIN: multipart = await socket.recv_multipart() return json.loads(multipart[0].decode()) else: raise TimeoutError("请求超时") finally: # 确保socket无论是否出错都能关闭 socket.close() # 应用退出时调用,销毁全局Context def cleanup_context(): global _global_context if _global_context: _global_context.term() _global_context = None
关键修复点
Server端
- 给ROUTER套接字添加TCP心跳配置,自动清理僵尸连接,避免资源积累
- 用
try-except包裹请求处理逻辑,防止单个请求异常导致服务崩溃
Client端
- 全局复用
Context,避免频繁创建销毁单例Context导致的资源泄漏 - 给
poller.poll()设置超时时间,防止服务器无响应时客户端无限阻塞 - 将同步
socket.send()改为异步await socket.send(),契合异步IO执行逻辑 - 用
finally块确保套接字无论是否出错都能关闭,避免连接资源泄漏 - 添加TCP心跳配置,维持连接健康状态
内容的提问来源于stack exchange,提问作者Kile Alkuri
相关产品推荐
相关产品推荐

