为何触发ZeroMQ Errno 4 ZMQException?技术排查求助
问题分析:ZeroMQ Errno 4异常排查
错误日志
org.zeromq.ZMQException: Errno 4 at org.zeromq.ZMQ$Socket.mayRaise(ZMQ.java:3732) ~[jeromq-0.5.3.jar:na] at org.zeromq.ZMQ$Socket.recv(ZMQ.java:3530) ~[jeromq-0.5.3.jar:na] at com.forexassistant.service.zeromq.CurrencyStrengthZeroMQ.sendCurrencyStrengthRequest(CurrencyStrengthZeroMQ.java:30) ~[classes/:na] at com.forexassistant.service.algorithmlogic.AlgorithmLogic.getCurrencyStrength(AlgorithmLogic.java:209) ~[classes/:na] at sun.reflect.GeneratedMethodAccessor111.invoke(Unknown Source) ~[na:na] at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) ~[na:1.8.0_241] at java.lang.reflect.Method.invoke(Method.java:498) ~[na:1.8.0_241] at org.springframework.scheduling.support.ScheduledMethodRunnable.run(ScheduledMethodRunnable.java:84) ~[spring-context-5.3.22.jar:5.3.22] at org.springframework.scheduling.support.DelegatingErrorHandlingRunnable.run(DelegatingErrorHandlingRunnable.java:54) ~[spring-context-5.3.22.jar:5.3.22] at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511) [na:1.8.0_241] at java.util.concurrent.FutureTask.runAndReset(FutureTask.java:308) [na:1.8.0_241] at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.access$301(ScheduledThreadPoolExecutor.java:180) [na:1.8.0_241] at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:294) [na:1.8.0_241] at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) [na:1.8.0_241] at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) [na:1.8.0_241] at java.lang.Thread.run(Thread.java:748) [na:1.8.0_241] Suppressed: org.zeromq.ZMQException: Errno 4 at zmq.Ctx.terminate(Ctx.java:304) ~[jeromq-0.5.3.jar:na] at org.zeromq.ZMQ$Context.term(ZMQ.java:671) ~[jeromq-0.5.3.jar:na] at org.zeromq.ZContext.destroy(ZContext.java:136) ~[jeromq-0.5.3.jar:na] at org.zeromq.ZContext.close(ZContext.java:463) ~[jeromq-0.5.3.jar:na] at com.forexassistant.service.zeromq.CurrencyStrengthZeroMQ.sendCurrencyStrengthRequest(CurrencyStrengthZeroMQ.java:37) ~[classes/:na] ... 13 common frames omitted
问题根源
结合客户端、服务端代码及报错信息,异常由以下多方面问题共同导致:
客户端重复关闭ZContext
客户端使用try-with-resources自动管理ZContext生命周期,但代码中又手动调用context.close(),导致Context被多次终止,触发Errno 4异常。客户端阻塞recv无超时,引发线程泄漏
socket2.recv(0)是无限阻塞调用,若服务端未发送响应(比如服务端send失败),当前线程会一直挂起。Spring每5秒触发一次任务,会快速耗尽线程池资源,后续任务被中断时,recv操作触发中断错误(Errno 4对应EINTR)。服务端重复发送空消息
MessageHandler返回空的ZmqMsg,OnTimer中又执行pushSocket.send(reply, true)发送空消息,导致send失败(打印“###ERROR### Sending message”),进而客户端无法收到响应,陷入阻塞循环。频繁创建销毁Context/Socket引发资源压力
客户端每5秒创建新的ZContext和Socket,频繁的连接/断开会导致服务端Socket的连接队列压力增大,结合服务端设置的HWM=1(高水位线),极易触发发送失败。
修复方案
客户端修复
- 移除手动调用
context.close(),依赖try-with-resources自动管理生命周期 - 给recv操作添加超时时间,避免无限阻塞
- 复用Context和Socket,减少频繁创建销毁的资源消耗
修复后示例代码:
import org.zeromq.SocketType; import org.zeromq.ZContext; import org.zeromq.ZMQ; import org.springframework.stereotype.Component; import javax.annotation.PostConstruct; import javax.annotation.PreDestroy; @Component public class CurrencyStrengthZeroMQ { private ZContext context; private ZMQ.Socket pushSocket; private ZMQ.Socket pullSocket; private static final int RECV_TIMEOUT = 5000; // 5秒超时 @PostConstruct public void init() { context = new ZContext(); // 初始化PUSH Socket pushSocket = context.createSocket(SocketType.PUSH); pushSocket.connect("tcp://localhost:32868"); // 初始化PULL Socket并设置超时 pullSocket = context.createSocket(SocketType.PULL); pullSocket.connect("tcp://localhost:32869"); pullSocket.setRcvTimeout(RECV_TIMEOUT); } public String sendCurrencyStrengthRequest() { try { String msg = "GET_CURRENCY_STRENGTHS"; // 阻塞发送,确保消息送达 pushSocket.send(msg.getBytes(ZMQ.CHARSET), 0); byte[] reply = pullSocket.recv(0); if (reply != null) { return new String(reply, ZMQ.CHARSET); } return null; // 超时或无响应返回null } catch (ZMQException e) { e.printStackTrace(); return null; } } @PreDestroy public void destroy() { pushSocket.close(); pullSocket.close(); context.close(); } }
服务端修复
- 移除
OnTimer中多余的空消息发送逻辑,保留InformPullClient的响应发送即可 - 确保消息处理流程中不产生无效的空消息发送操作
修复后OnTimer代码:
void OnTimer() { if(CheckServerStatus() == true) { bool hasAllOtherRequest = pullSocket.recv(request, true); if(hasAllOtherRequest){ if (request.size() > 0) { MessageHandler(request); // 内部已完成响应发送,无需额外处理 } } // 定期更新价格逻辑保留 if (GetTickCount() >= lastUpdateMillis + MILLISECOND_TIMER_PRICES) OnTick(); } }
内容的提问来源于stack exchange,提问作者jose1278
相关产品推荐
相关产品推荐

