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

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

问题分析

  1. Client端资源泄漏:每次请求都销毁全局单例Context,多次term()会导致底层资源未正确释放,后续连接创建失败
  2. 无超时处理:Client的poll()未设置超时,若服务器无响应会无限阻塞
  3. Server异常崩溃风险:请求处理逻辑未捕获异常,单个请求出错会导致整个服务循环中断
  4. 僵尸连接积累: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 05:15:59