惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

Stack Overflow Blog
Stack Overflow Blog
J
Java Code Geeks
Last Week in AI
Last Week in AI
人人都是产品经理
人人都是产品经理
博客园 - 【当耐特】
钛媒体:引领未来商业与生活新知
钛媒体:引领未来商业与生活新知
C
Check Point Blog
月光博客
月光博客
腾讯CDC
Engineering at Meta
Engineering at Meta
博客园 - Franky
Vercel News
Vercel News
D
Docker
OSCHINA 社区最新新闻
OSCHINA 社区最新新闻
F
Fortinet All Blogs
Microsoft Security Blog
Microsoft Security Blog
奇客Solidot–传递最新科技情报
奇客Solidot–传递最新科技情报
雷峰网
雷峰网
Google DeepMind News
Google DeepMind News
Martin Fowler
Martin Fowler
GbyAI
GbyAI
B
Blog
Hugging Face - Blog
Hugging Face - Blog
T
Tailwind CSS Blog

博客园 - 海乐学习

VS2026创建.net8的控制台程序并能在银河麒麟V10的Liunx上运行 在麒麟系统中卸载RabbitMQ Java Maven 开发的常用命令 C语言开发的常用命令 三汇Linux配置Config说明 Java实现优雅的关闭程序并执行清理流程的写法 C语言开发中优雅的关闭程序并执行清理流程的写法 将C语言开发的程序做成 麒麟系统的服务(Systemd 标准服务,Linux 通用) 实现开机自启动 ZeroMQ中ZMQ_DEALER单帧数据(支持异步收发、主动推送事件)C写服务端 java写客户端 win10系统中 关闭专注助手 执行 xx.sh 脚本文件时出现: /bin/bash^M:解释器错误: 没有那个文件或目录 运行编译打包好的 xxx.jar 中没有主清单属性 出现这个错误 将Java编译的 .jar文件做成 麒麟系统的服务(Systemd 标准服务,Linux 通用) 实现开机自启动 将Java编译的 .jar文件做成windows服务 实现开机自启动 方法二 RabbitMQ在麒麟系统中离线安装说明 VMware启动虚似机后出现 无法获取快照信息: 锁定文件失败 模块“Snapshot”启动失败。未能启动虚拟机。 C语言在Linux中开发完整Demo包含读配置文件写日志和定时器Timer C语言在Linux中开发读取配置文件app.conf 麒麟V10 Server系统中搭建C语言开发环境 麒麟ServerV10 修改IP4地址 麒麟ServerV10 配置IP4 当系统中有两个版本的Maven时,用IDEA创建Maven有时会出错 麒麟ServerV10安装 espeak-ng 和 ffmpeg 方法 C语言在Linux中开发没有界面纯后台运行的Demo程序(含日志和Timer) C语言在Linux中开发使用定时器Timer在界面上显示时间 C语言在Linux中开发带界面的程序(含每小时日志) C语言在Linux中开发第一个项目Hello Word 在apache-maven项目中使用log4写日志 在apache-maven项目中解决中文乱码问题 在apache-maven项目中读取配置文件
ZeroMQ的DEALER双帧结构(路由空帧 + 业务帧)(支持异步收发、...
海乐学习 · 2026-06-17 · via 博客园 - 海乐学习

整体说明
C 服务端 DEALER <-> Java 客户端 DEALER 对等通信 处理双帧结构(路由空帧 + 业务帧)
C 主动推送事件、回复客户端指令统一走多帧接口,加锁防止多线程打乱帧分片
Java 客户端发送增加 sendMore 分段,接收循环读取全部帧,丢弃路由帧只拿业务内容
Java 客户端手动 connect,避免监听器未注册就触发连接回调
心跳 PING/PONG 完整互通,断线重连逻辑修复
Java 客户端每60秒自动发送 PING,C 服务端 收到Java客户端消息:PING,回复 PONG
Java 客户端连续三次(180秒超时),未收到PING,判定断线,启动重连机制
C 服务端代码调用Zmq_SendMessage推送通道事件,Java客户端 监听器完整收到事件字符串
Ctrl+C 关闭 C 服务端,Java 客户端自动检测断开,退避重连;重启 C 服务端后自动恢复连接

C 服务端代码

ZeroMQServer.c

/*
 * ============================================================================
 * ZeroMQServer.c - ZeroMQ 通信模块实现
 * 适配:DEALER 对等套接字 与 Java DEALER 客户端互通
 * 核心修复:增加多帧读写封装,解决对等套接字路由帧丢消息问题
 * ============================================================================
 */
#include "ZeroMQServer.h"
#include "Logger.h"
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <time.h>
#include <zmq.h>

/*====================== 内部函数声明 ======================*/
static void *Zmq_ServerThread(void *arg);
static int Zmq_ProcessMessage(const char *szMsg, char *szReply, int nReplySize);
//【修复新增】对等DEALER多帧读写封装
static int Zmq_ReadFullMsg(char *outBuf, int bufLen);
static int Zmq_WriteFullMsg(const char *data);

void Zmq_SendMessage(const char *szMessage);

/*====================== 全局变量定义 ======================*/
ZMQ_HANDLE      g_Zmq;              // ZeroMQ 句柄
pthread_mutex_t g_ZmqMutex;         // ZeroMQ 互斥锁:所有收发操作加锁,防止多线程拆分帧错乱
Zmq_CommandHandler g_CommandHandler = NULL;  // 业务命令处理器(对接ShPbxServer)

/*====================== 【修复新增】读取完整对等DEALER消息 ======================*/
// 对等DEALER消息固定两帧:1.空路由ID帧  2.业务数据帧
// 丢弃第一帧路由占位,输出纯业务字符串到outBuf
static int Zmq_ReadFullMsg(char *outBuf, int bufLen)
{
    int rc;
    zmq_msg_t idMsg;
    zmq_msg_init(&idMsg);

    // 第一步:读取并丢弃路由占位帧
    rc = zmq_msg_recv(&idMsg, g_Zmq.pSocket, ZMQ_DONTWAIT);
    if (rc == -1)
    {
        zmq_msg_close(&idMsg);
        return -1;
    }
    zmq_msg_close(&idMsg);

    // 第二步:读取真正业务数据帧
    rc = zmq_recv(g_Zmq.pSocket, outBuf, bufLen - 1, ZMQ_DONTWAIT);
    if (rc < 0)
        return -1;

    outBuf[rc] = '\0'; // 字符串结尾补0
    return rc;
}

/*====================== 【修复新增】发送完整对等DEALER消息 ======================*/
// 对等DEALER发送规范:先发空分段帧(ZMQ_SNDMORE),再发业务数据
static int Zmq_WriteFullMsg(const char *data)
{
    if (!g_Zmq.bRunning || !g_Zmq.pSocket)
        return -1;
    int dataLen = strlen(data);

    // 第一段:空路由占位帧,标记后续还有数据
    int rc = zmq_send(g_Zmq.pSocket, "", 0, ZMQ_SNDMORE);
    if (rc == -1)
    {
        Log_Write("[ZeroMQServer] 发送分段帧失败 err:%s", zmq_strerror(zmq_errno()));
        return -1;
    }
    // 第二段:实际业务消息,无SNDMORE代表结束
    rc = zmq_send(g_Zmq.pSocket, data, dataLen, 0);
    if (rc == -1)
    {
        Log_Write("[ZeroMQServer] 发送业务帧失败 err:%s", zmq_strerror(zmq_errno()));
        return -1;
    }
    return 0;
}

/*====================== ZeroMQ 服务端线程(【修复】替换原有单帧recv) ======================*/
static void *Zmq_ServerThread(void *arg)
{
    char buffer[ZMQ_BUFFER_SIZE];
    char szReply[ZMQ_BUFFER_SIZE];
    int timeout = ZMQ_TIMEOUT;

    // 设置收发超时,非阻塞轮询
    zmq_setsockopt(g_Zmq.pSocket, ZMQ_RCVTIMEO, &timeout, sizeof(timeout));
    zmq_setsockopt(g_Zmq.pSocket, ZMQ_SNDTIMEO, &timeout, sizeof(timeout));

    Log_Write("[ZeroMQServer] 服务端线程已启动");

    while (g_Zmq.bRunning)
    {
        pthread_mutex_lock(&g_ZmqMutex);
        //【修复】使用多帧读取函数替代原zmq_recv单帧读取
        int size = Zmq_ReadFullMsg(buffer, sizeof(buffer));
        if (size >= 0)
        {
            Log_Write("[ZeroMQServer] 收到客户端消息:%s", buffer);
            // 交给业务层处理PING/JSON指令
            Zmq_ProcessMessage(buffer, szReply, sizeof(szReply));

            //【修复】多帧发送回复给客户端
            int sendRet = Zmq_WriteFullMsg(szReply);
            if (sendRet == 0)
            {
                Log_Write("[ZeroMQServer] 回复客户端:%s", szReply);
            }
            else
            {
                Log_Write("[ZeroMQServer] 回复消息发送失败");
            }
        }
        pthread_mutex_unlock(&g_ZmqMutex);

        // 上下文销毁,退出循环
        if (zmq_errno() == ETERM)
        {
            break;
        }
        // 超时无消息,继续循环
    }

    Log_Write("[ZeroMQServer] 服务端线程已退出");
    return NULL;
}

/*====================== 处理客户端消息(【原有逻辑保留】) ======================*/
static int Zmq_ProcessMessage(const char *szMsg, char *szReply, int nReplySize)
{
    // 1. 心跳PING指令,直接返回PONG
    if (strcmp(szMsg, "PING") == 0)
    {
        snprintf(szReply, nReplySize, "PONG");
        return 0;
    }

    // 2. 业务JSON命令,转发给ShPbxServer上层处理器
    if (g_CommandHandler != NULL)
    {
        return g_CommandHandler(szMsg, szReply, nReplySize);
    }

    // 3. 未注册业务处理器默认报错
    snprintf(szReply, nReplySize, "ERROR|Command handler not registered");
    return -1;
}

/*====================== 初始化 ZeroMQ(【修复】删除无效ROUTING_ID配置) ======================*/
int Zmq_Init(void)
{
    // 1. 创建ZeroMQ上下文
    g_Zmq.pContext = zmq_ctx_new();
    if (!g_Zmq.pContext)
    {
        Log_Write("[ZeroMQServer] 创建zmq上下文失败!");
        return -1;
    }

    // 2. 创建DEALER套接字(对等双向通信,支持主动推送事件)
    g_Zmq.pSocket = zmq_socket(g_Zmq.pContext, ZMQ_DEALER);
    if (!g_Zmq.pSocket)
    {
        Log_Write("[ZeroMQServer] 创建DEALER套接字失败!");
        zmq_ctx_destroy(g_Zmq.pContext);
        return -1;
    }

    //【修复】DEALER对等模式不需要手动设置ROUTING_ID,删除该行
    // zmq_setsockopt(g_Zmq.pSocket, ZMQ_ROUTING_ID, ZMQ_SERVER_NAME, strlen(ZMQ_SERVER_NAME));

    // 3. 绑定监听端口5555
    int rc = zmq_bind(g_Zmq.pSocket, "tcp://*:5555");
    if (rc != 0)
    {
        Log_Write("[ZeroMQServer] 绑定5555端口失败:%s", zmq_strerror(zmq_errno()));
        zmq_close(g_Zmq.pSocket);
        zmq_ctx_destroy(g_Zmq.pContext);
        return -1;
    }

    // 初始化运行标记、互斥锁
    g_Zmq.bRunning = 1;
    pthread_mutex_init(&g_ZmqMutex, NULL);

    // 启动后台收发线程
    pthread_create(&g_Zmq.hThread, NULL, Zmq_ServerThread, NULL);

    Log_Write("[ZeroMQServer] 已启动,监听 tcp://*:5555");
    return 0;
}

/*====================== 清理ZeroMQ资源(【修复】先关socket唤醒阻塞recv,防止死锁) ======================*/
void Zmq_Cleanup(void)
{
    if (!g_Zmq.bRunning)
    {
        return;
    }
    Log_Write("[ZeroMQServer] 正在关闭ZeroMQ服务...");
    g_Zmq.bRunning = 0;

    //【修复】优先关闭套接字,唤醒阻塞在recv的线程,避免timedjoin卡死
    if (g_Zmq.pSocket)
    {
        zmq_close(g_Zmq.pSocket);
        g_Zmq.pSocket = NULL;
    }

    // 等待线程最多3秒退出
    struct timespec ts;
    ts.tv_sec = time(NULL) + 3;
    ts.tv_nsec = 0;
    pthread_timedjoin_np(g_Zmq.hThread, NULL, &ts);

    // 销毁上下文
    if (g_Zmq.pContext)
    {
        zmq_ctx_destroy(g_Zmq.pContext);
        g_Zmq.pContext = NULL;
    }

    // 销毁互斥锁
    pthread_mutex_destroy(&g_ZmqMutex);
    Log_Write("[ZeroMQServer] ZeroMQ 已安全关闭");
}

/*====================== 注册上层业务命令处理器(【原有逻辑保留】) ======================*/
void Zmq_SetCommandHandler(Zmq_CommandHandler handler)
{
    g_CommandHandler = handler;
    Log_Write("[ZeroMQServer] 业务命令处理器已注册");
}

/*====================== 主动向客户端推送通道事件(【修复】替换单帧send为多帧封装) ======================*/
// 上层ShPbxServer调用:通道振铃/通话/挂机事件主动推送
void Zmq_SendMessage(const char *szMessage)
{
    if (!g_Zmq.bRunning || !g_Zmq.pSocket)
    {
        return;
    }
    pthread_mutex_lock(&g_ZmqMutex);
    int rc = Zmq_WriteFullMsg(szMessage);
    pthread_mutex_unlock(&g_ZmqMutex);

    if (rc == 0)
    {
        Log_Write("[ZeroMQServer] 主动推送事件:%s", szMessage);
    }
    else
    {
        Log_Write("[ZeroMQServer] 主动推送事件发送失败");
    }
}

/*====================== JSON解析函数外部声明(对接ShPbxServer) ======================*/
extern const char* ParseMsgJsonValue(const char *szJson, const char *szKey, char *szBuffer, int nBufSize);

ZeroMQServer.h

/*
 * ============================================================================
 * ZeroMQServer.h - ZeroMQ 通信模块头文件
 * ============================================================================
 */

#ifndef ZERO_MQ_SERVER_H
#define ZERO_MQ_SERVER_H

#include <zmq.h>
#include <pthread.h>

/*====================== ZeroMQ 配置 ======================*/
// 如果未在 ShPbxServer.h 中定义,则在此定义
#ifndef ZMQ_PORT
#define ZMQ_PORT            5555            // ZeroMQ 监听端口
#endif
#ifndef ZMQ_SERVER_NAME
#define ZMQ_SERVER_NAME     "ShPbxServer"   // ZeroMQ 服务端标识
#endif
#ifndef ZMQ_BUFFER_SIZE
#define ZMQ_BUFFER_SIZE     1024            // 消息缓冲区大小
#endif
#ifndef ZMQ_TIMEOUT
#define ZMQ_TIMEOUT         1000            // ZeroMQ 超时时间 (ms)
#endif

/*====================== ZeroMQ 句柄结构 ======================*/
typedef struct {
    void        *pContext;          // ZeroMQ 上下文
    void        *pSocket;           // ZeroMQ 套接字
    int         bRunning;           // 运行标志
    pthread_t   hThread;            // ZeroMQ 线程句柄
} ZMQ_HANDLE;

/*====================== 全局变量声明 ======================*/
extern ZMQ_HANDLE      g_Zmq;               // ZeroMQ 句柄
extern pthread_mutex_t g_ZmqMutex;          // ZeroMQ 互斥锁

/*====================== 函数声明 ======================*/

/**
 * 业务命令处理回调函数类型
 * @param szMsg 客户端消息
 * @param szReply 回复缓冲区
 * @param nReplySize 回复缓冲区大小
 * @return 0=成功,-1=失败
 */
typedef int (*Zmq_CommandHandler)(const char *szMsg, char *szReply, int nReplySize);

/**
 * 初始化 ZeroMQ 服务端
 * @return 0=成功,-1=失败
 */
int Zmq_Init(void);

/**
 * 清理 ZeroMQ 资源
 */
void Zmq_Cleanup(void);

/**
 * 向ZeroMQ客户端发送消息
 */
void Zmq_SendMessage(const char *szMessage);

/**
 * 设置业务命令处理器
 * @param handler 业务命令处理函数指针
 */
void Zmq_SetCommandHandler(Zmq_CommandHandler handler);

#endif /* SH_ZMQ_H */

调用方法

/*====================== 内部函数声明 ======================*/
static int Zmq_HandleClientMessage(const char *szMsg, char *szReply, int nReplySize);

初始化

    // 初始化 ZeroMQ 服务端
    if (Zmq_Init() != 0)
    {
        Log_Write("[ShPbxServer][警告] ZeroMQ 初始化失败,ZeroMQ 服务不可用");
    }
    else
    {
        // 注册业务命令处理器
        Zmq_SetCommandHandler(Zmq_HandleClientMessage);
        Log_Write("[ShPbxServer][信息] ZeroMQ 服务端已启动,监听端口:%d", ZMQ_PORT);
    }

关闭

// 先关闭ZeroMQ(内部等待服务线程退出)
Zmq_Cleanup();
Log_Write("[ShPbxServer][信息] ZeroMQ 服务已关闭");

发送消息

// 构建 JSON 格式消息
snprintf(szJson, sizeof(szJson),
             "{\"type\":\"%s\",\"EventName\":\"%s\",\"StateName\":\"%s\","
             "\"Event\":%d,\"Reference\":%d,\"Param\":%d,\"User\":%d}",
             szEventType, szEventName, szStateName,
             wEvent, nChId, nParam, (int)dwUser);

//向ZeroMQ客户端发送消息 
Zmq_SendMessage(szJson);// ZeroMQServer.c (通信层) - 只负责发送

接收消息

/**
 * 处理 ZeroMQ 客户端消息
 * @param szMsg 客户端消息
 * @param szReply 回复缓冲区
 * @param nReplySize 回复缓冲区大小
 * @return 0=成功,-1=失败
 */
int Zmq_HandleClientMessage(const char *szMsg, char *szReply, int nReplySize)
{
    // 1. 判断是否为 JSON 格式
    if (strlen(szMsg) > 0 && szMsg[0] == '{')
    {
        // 解析 JSON 格式:{"Type":"Action","ActionName":"SsmGetMaxCh","ChannelId":"","Parameters":""}
        char szType[64] = {0};
        char szActionName[128] = {0};
        char szChannelId[256] = {0};
        char szParameters[512] = {0};

        ParseMsgJsonValue(szMsg, "Type", szType, sizeof(szType));
        ParseMsgJsonValue(szMsg, "ActionName", szActionName, sizeof(szActionName));
        ParseMsgJsonValue(szMsg, "ChannelId", szChannelId, sizeof(szChannelId));
        ParseMsgJsonValue(szMsg, "Parameters", szParameters, sizeof(szParameters));

        Log_Write("[ShPbxServer] 收到Api请求 - Type:%s, ActionName:%s, ChannelId:%s, Parameters:%s", szType, szActionName, szChannelId, szParameters);

        // 处理 Action 类型请求
        if (strcmp(szType, "Action") == 0)
        {
            //具体业务处理
         }
    }
    // 未知命令
    snprintf(szReply, nReplySize, "ERROR|Unknown command: %s", szMsg);
    Log_Write("[ShPbxServer] 未知命令:%s", szMsg);
    return -1;
}

Java 客户端

ZeroMQClient.java

package com.ShPbxClient;

import org.zeromq.SocketType;
import org.zeromq.ZContext;
import org.zeromq.ZMQ;
import org.zeromq.ZMQException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.util.concurrent.atomic.AtomicBoolean;

/**
 * ZeroMQ DEALER 客户端,对接C服务端DEALER对等套接字
 * 修复点:
 * 1. send增加sendMore分段,匹配服务端双帧协议
 * 2. 接收循环读取全部帧,丢弃路由空帧,提取业务消息
 * 3. 构造函数移除自动connect,新增手动connect,防止监听器未注册回调丢失
 * 4. 重连时销毁旧socket,清空缓存脏帧
 * 5. 完整心跳PING/PONG互通,断线指数退避重连
 */
public class ZeroMQClient {
    private static final Logger logger = LoggerFactory.getLogger(ZeroMQClient.class);

    // 服务端地址端口配置
    private final String host;
    private final int port;

    // ZeroMQ底层资源
    private ZContext context;
    private ZMQ.Socket dealerSocket;

    // 全局运行标记:控制收发、心跳线程启停
    private volatile boolean isRunning = false;
    // 原子连接状态标记,多线程安全读写
    private final AtomicBoolean isConnected = new AtomicBoolean(false);

    // 重连退避配置
    private int reconnectDelay = 1000;        // 初始重连延迟1秒
    private static final int MAX_RECONNECT_DELAY = 30000; // 最大延迟30秒

    // 心跳配置
    private static final int HEARTBEAT_INTERVAL = 60000;  // 60秒发一次PING
    private static final int HEARTBEAT_TIMEOUT = 180000;  // 180秒无响应判定断开

    // 最后一次收发消息时间戳,用于心跳超时判断
    private volatile long lastActivityTime = System.currentTimeMillis();
    // 心跳后台线程引用
    private Thread heartbeatThread;

    // 消息回调:接收服务端推送事件、指令回复
    public interface MessageListener {
        void onMessageReceived(String message);
    }
    private MessageListener listener;

    // 连接状态回调:上线/断开/重连通知业务层
    public interface ConnectionListener {
        void onConnected();
        void onDisconnected();
        void onReconnecting(int retryCount);
    }
    private ConnectionListener connectionListener;

    /**
     * 构造函数:仅初始化上下文,【修复】不再自动连接
     * @param host C服务端IP
     * @param port C服务端端口5555
     */
    public ZeroMQClient(String host, int port) {
        this.host = host;
        this.port = port;
        this.context = new ZContext();
        // 原代码 this.connect(); 已删除,改为外部手动调用
    }

    /**
     * 设置消息接收监听器(必须connect前调用)
     */
    public void setMessageListener(MessageListener listener) {
        this.listener = listener;
    }

    /**
     * 设置连接状态监听器(必须connect前调用)
     */
    public void setConnectionListener(ConnectionListener listener) {
        this.connectionListener = listener;
    }

    /**
     * 【修复新增】手动发起连接,外部设置完监听器再调用
     */
    public void connect() {
        try {
            // 存在旧套接字先关闭,释放资源
            if (dealerSocket != null) {
                dealerSocket.close();
                dealerSocket = null;
            }
            // 创建DEALER套接字连接服务端
            dealerSocket = context.createSocket(SocketType.DEALER);
            dealerSocket.connect(String.format("tcp://%s:%d", host, port));

            // 更新连接状态
            isConnected.set(true);
            lastActivityTime = System.currentTimeMillis();
            reconnectDelay = 1000; // 重置重连延迟

            logger.info("ZeroMQ 成功连接服务端 tcp://{}:{}", host, port);
            if (connectionListener != null) {
                connectionListener.onConnected();
            }
        } catch (Exception e) {
            isConnected.set(false);
            logger.error("ZeroMQ 连接服务端失败 tcp://{}:{}", host, port, e);
            if (connectionListener != null) {
                connectionListener.onDisconnected();
            }
        }
    }

    /**
     * 【修复】多帧分段发送消息,匹配C端DEALER双帧协议
     * @param message 待发送PING/JSON指令
     * @return true发送成功 false失败/未连接
     */
    public boolean send(String message) {
        if (!isConnected.get() || dealerSocket == null) {
            logger.warn("发送失败:当前未连接服务端");
            return false;
        }
        try {
            byte[] payload = message.getBytes(ZMQ.CHARSET);
            // 对等DEALER规范:先发空分段帧,标记后续还有数据
            dealerSocket.sendMore(new byte[0]);
            // 发送真实业务数据,无sendMore代表消息结束
            dealerSocket.send(payload, ZMQ.DONTWAIT);

            lastActivityTime = System.currentTimeMillis();
            return true;
        } catch (ZMQException e) {
            logger.error("发送消息异常", e);
            isConnected.set(false);
            if (connectionListener != null) {
                connectionListener.onDisconnected();
            }
            return false;
        }
    }

    /**
     * 启动接收线程 + 心跳检测线程
     * 连接成功后调用
     */
    public void startReceive() {
        isRunning = true;
        // 异步消息接收线程
        new Thread(this::receiveLoop, "ZeroMQ-Receiver").start();
        // 心跳后台线程
        heartbeatThread = new Thread(this::heartbeatLoop, "ZeroMQ-Heartbeat");
        heartbeatThread.start();
        logger.info("ZeroMQ 接收线程、心跳线程已启动");
    }

    /**
     * 【修复】消息接收循环,完整读取双帧,丢弃路由帧提取业务内容
     */
    private void receiveLoop() {
        // 接收超时1秒,非阻塞轮询
        dealerSocket.setReceiveTimeOut(1000);
        while (isRunning && !Thread.currentThread().isInterrupted()) {
            try {
                // 第一帧:空路由占位帧,直接丢弃
                byte[] frameRoute = dealerSocket.recv();
                if (frameRoute == null) continue;
                // 判断是否存在下一帧业务数据
                boolean hasNextFrame = dealerSocket.hasReceiveMore();
                if (!hasNextFrame) continue;

                // 第二帧:真实业务消息(PING回复PONG、服务端推送事件JSON)
                byte[] dataFrame = dealerSocket.recv();
                if (dataFrame == null) continue;
                String msg = new String(dataFrame, ZMQ.CHARSET);

                // 更新活跃时间戳,心跳超时重置
                lastActivityTime = System.currentTimeMillis();

                // 服务端返回PONG,恢复正常连接状态、重置重连延迟
                if ("PONG".equals(msg)) {
                    isConnected.set(true);
                    reconnectDelay = 1000;
                    logger.debug("收到心跳PONG");
                }
                // 回调业务层处理消息
                if (listener != null) {
                    listener.onMessageReceived(msg);
                }
            } catch (ZMQException e) {
                // EAGAIN=超时无消息,忽略;其他异常判定断开
                if (e.getErrorCode() != ZMQ.Error.EAGAIN.getCode()) {
                    logger.error("接收消息异常,连接断开", e);
                    isConnected.set(false);
                    if (connectionListener != null) {
                        connectionListener.onDisconnected();
                    }
                    break;
                }
            }
        }
    }

    /**
     * 心跳检测循环:定时发PING、超时自动重连
     */
    private void heartbeatLoop() {
        int retryCount = 0;
        while (isRunning && !Thread.currentThread().isInterrupted()) {
            try {
                Thread.sleep(HEARTBEAT_INTERVAL);
                long idleMs = System.currentTimeMillis() - lastActivityTime;

                // 1. 长时间无任何通信,判定心跳超时断开
                if (idleMs > HEARTBEAT_TIMEOUT) {
                    if (isConnected.get()) {
                        logger.warn("心跳超时{}ms,判定与服务端断开", idleMs);
                        isConnected.set(false);
                        if (connectionListener != null) connectionListener.onDisconnected();
                    }
                    retryCount++;
                    if (connectionListener != null) connectionListener.onReconnecting(retryCount);
                    reconnect();
                }
                // 2. 当前未连接,持续尝试重连
                else if (!isConnected.get()) {
                    retryCount++;
                    if (connectionListener != null) connectionListener.onReconnecting(retryCount);
                    reconnect();
                }
                // 3. 连接正常,发送心跳PING
                else {
                    send("PING");
                    logger.debug("发送心跳PING");
                }
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
                break;
            } catch (Exception e) {
                logger.error("心跳线程运行异常", e);
            }
        }
    }

    /**
     * 【修复】断线重连,先销毁旧socket清空脏帧,指数退避延迟
     */
    private void reconnect() {
        logger.info("准备重连,等待{}ms", reconnectDelay);
        try {
            // 销毁旧套接字,清除残留缓存消息
            if (dealerSocket != null) {
                dealerSocket.close();
                dealerSocket = null;
            }
            Thread.sleep(reconnectDelay);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            return;
        }
        // 延迟翻倍,上限30秒
        reconnectDelay = Math.min(reconnectDelay * 2, MAX_RECONNECT_DELAY);
        // 发起新连接
        connect();
    }

    /**
     * 获取当前连接状态
     */
    public boolean isConnected() {
        return isConnected.get();
    }

    /**
     * 停止收发、心跳线程
     */
    public void stopReceive() {
        isRunning = false;
        if (heartbeatThread != null) {
            heartbeatThread.interrupt();
        }
    }

    /**
     * 彻底关闭客户端,释放所有ZeroMQ资源
     */
    public void close() {
        logger.info("[ZeroMQ客户端] 开始释放资源关闭");
        // 1. 停止后台线程
        stopReceive();
        // 2. 关闭套接字
        if (dealerSocket != null) {
            try {
                dealerSocket.close();
                logger.debug("[ZeroMQ客户端] DEALER套接字已关闭");
            } catch (Exception e) {
                logger.error("[ZeroMQ客户端] 关闭套接字异常", e);
            }
            dealerSocket = null;
        }
        // 3. 销毁上下文
        if (context != null) {
            try {
                context.close();
                logger.debug("[ZeroMQ客户端] ZContext上下文已销毁");
            } catch (Exception e) {
                logger.error("[ZeroMQ客户端] 销毁上下文异常", e);
            }
            context = null;
        }
        // 标记断开
        isConnected.set(false);
        logger.info("[ZeroMQ客户端] 完全关闭完成");
    }
}

调用方法

初始化

private static ZeroMQClient zmqClient;
String host ="192.168.1.122";
 int port = 5555;
//1. 初始化 ZeroMQ 客户端(命令通道) zmqClient = new ZeroMQClient(host, port); //2. 设置消息监听器 zmqClient.setMessageListener(App::handleServerMessage); // 设置连接状态监听器 zmqClient.setConnectionListener(new ZeroMQClient.ConnectionListener() { @Override public void onConnected() { logger.info("[连接状态] 服务器已连接"); } @Override public void onDisconnected() { logger.warn("[连接状态] 服务器连接断开"); } @Override public void onReconnecting(int retryCount) { logger.info("[连接状态] 正在重连服务器... (第 {} 次)", retryCount); } }); // 3. 手动发起连接 zmqClient.connect(); // 4. 启动接收、心跳线程 zmqClient.startReceive();

关闭退出

  // 添加关闭钩子
  Runtime.getRuntime().addShutdownHook(new Thread(() -> {
          
            if (zmqClient != null) {
                zmqClient.close();
            }
           
  }));

收到消息

   /**
     * 处理服务器消息
     * @param message 收到的消息
     */
    public static void handleServerMessage(String message) {
        logger.info("收到:{}", message);
        //具体业务代码
    }

发送消息

String jsonRequest = JSON_MAPPER.writeValueAsString(json);
 if (zmqClient.send(jsonRequest)) {
       logger.info("[请求] 成功发送 SsmHangup: {}", jsonRequest);
} else {
       logger.error("[请求] 发送 SsmHangup 失败(客户端未连接或底层异常)");
}