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

推荐订阅源

S
Security Affairs
钛媒体:引领未来商业与生活新知
钛媒体:引领未来商业与生活新知
大猫的无限游戏
大猫的无限游戏
OSCHINA 社区最新新闻
OSCHINA 社区最新新闻
爱范儿
爱范儿
阮一峰的网络日志
阮一峰的网络日志
GbyAI
GbyAI
D
Docker
美团技术团队
N
Netflix TechBlog - Medium
罗磊的独立博客
V
Visual Studio Blog
人人都是产品经理
人人都是产品经理
freeCodeCamp Programming Tutorials: Python, JavaScript, Git & More
Hugging Face - Blog
Hugging Face - Blog
酷 壳 – CoolShell
酷 壳 – CoolShell
Jina AI
Jina AI
CTFtime.org: upcoming CTF events
CTFtime.org: upcoming CTF events
M
MIT News - Artificial intelligence
腾讯CDC
MongoDB | Blog
MongoDB | Blog
Last Week in AI
Last Week in AI
博客园 - 三生石上(FineUI控件)
博客园 - 叶小钗
V
V2EX
L
LangChain Blog
博客园 - 【当耐特】
B
Blog RSS Feed
量子位
U
Unit 42
Engineering at Meta
Engineering at Meta
小众软件
小众软件
宝玉的分享
宝玉的分享
H
Help Net Security
Microsoft Azure Blog
Microsoft Azure Blog
云风的 BLOG
云风的 BLOG
博客园 - 聂微东
博客园 - 司徒正美
The Cloudflare Blog
The GitHub Blog
The GitHub Blog
T
Tailwind CSS Blog
Exploit-DB.com RSS Feed
Exploit-DB.com RSS Feed
The Last Watchdog
The Last Watchdog
cs.AI updates on arXiv.org
cs.AI updates on arXiv.org
S
SegmentFault 最新的问题
博客园_首页
Attack and Defense Labs
Attack and Defense Labs
TaoSecurity Blog
TaoSecurity Blog
Apple Machine Learning Research
Apple Machine Learning Research
S
Security @ Cisco Blogs

博客园 - 海乐学习

Java Maven 开发的常用命令 C语言开发的常用命令 三汇Linux配置Config说明 Java实现优雅的关闭程序并执行清理流程的写法 C语言开发中优雅的关闭程序并执行清理流程的写法 将C语言开发的程序做成 麒麟系统的服务(Systemd 标准服务,Linux 通用) 实现开机自启动 ZeroMQ的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项目中读取配置文件 创建apache-maven项目 远程桌面连接时出现身份验证错误 要求的函数不受支持 这可能是由于CredSSP加密数据库修正 win10系统查看电脑从锁屏状态回到使用状态 apache-maven的常用命令 C语言在 Linux 中的常用命令 apache-maven安装配置 麒麟CentOS下安装ZeroMQ开发包 Window上用VS Code + Remote-SSH组件的方式来实现开发编译Linux上的C++程序 win10弹出 无法使用内置管理员账户打开 Microsoft Edge。请使用其他账户登录 在麒麟系统上安装Qwen3-TTS文字转语音 在麒麟系统上安装MaryTTS文字转语音 FTP上传Linux/Unix文件系统权限的修改方法 麒麟系统Kylin Linux Advanced Server 中安装 python3.10 将exe做成windows服务 java实现ftp上传 java实现TTS文字转语音wav (Jacob + SAPI) node.js和Next.js 编译部署说明
ZeroMQ中ZMQ_DEALER单帧数据(支持异步收发、主动推送事件)C写服务端 java写客户端
海乐学习 · 2026-06-16 · via 博客园 - 海乐学习

ZeroMQ中使用ZMQ_DEALER单帧数据(支持异步收发、主动推送事件)C写服务端 java写客户端

C的服务端代码

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 */

ZeroMQServer.c

/*
 * ============================================================================
 * ZeroMQServer.c - ZeroMQ 通信模块实现
 * ============================================================================
 */

#include "ZeroMQServer.h"
#include "Logger.h"
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <time.h>

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

void Zmq_SendMessage(const char *szMessage);

/*====================== 全局变量定义 ======================*/
ZMQ_HANDLE      g_Zmq;              // ZeroMQ 句柄
pthread_mutex_t g_ZmqMutex;         // ZeroMQ 互斥锁
Zmq_CommandHandler g_CommandHandler = NULL;  // 业务命令处理器

/*====================== ZeroMQ 服务端线程 ======================*/
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)
    {
        // 接收客户端消息(非阻塞)
        int size = zmq_recv(g_Zmq.pSocket, buffer, sizeof(buffer) - 1, 0);
        if (size >= 0)
        {
            buffer[size] = '\0';
            Log_Write("[ZeroMQServer] 收到消息:%s", buffer);

            // 处理消息并生成回复
            pthread_mutex_lock(&g_ZmqMutex);
            //由ShPbxServer来处理Json请求消息
            Zmq_ProcessMessage(buffer, szReply, sizeof(szReply));
            //得到Json返回值后,自动发送应答
            zmq_send(g_Zmq.pSocket, szReply, strlen(szReply), 0);
            pthread_mutex_unlock(&g_ZmqMutex);

            Log_Write("[ZeroMQServer] 回复:%s", szReply);
        }
        else if (zmq_errno() == ETERM)
        {
            break;  // 上下文被销毁
        }
        // 超时则继续循环,检查 running 标志
    }

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

/*====================== 处理客户端消息 ======================*/
static int Zmq_ProcessMessage(const char *szMsg, char *szReply, int nReplySize)
{
    // 1. 心跳命令 - 直接处理
    if (strcmp(szMsg, "PING") == 0)
    {
        snprintf(szReply, nReplySize, "PONG");
        //Log_Write("[ZeroMQServer] 心跳 PING -> PONG");
        return 0;
    }

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

    // 3. 未注册处理器时的默认响应
    snprintf(szReply, nReplySize, "ERROR|Command handler not registered");
    return -1;
}

/*====================== 初始化 ZeroMQ ======================*/
int Zmq_Init(void)
{
    // 1. 创建 ZeroMQ 上下文
    g_Zmq.pContext = zmq_ctx_new();
    if (!g_Zmq.pContext)
    {
        Log_Write("[ZeroMQServer] 创建上下文失败!");
        return -1;
    }

    // 2. 创建 DEALER 套接字(异步双向通信)
    g_Zmq.pSocket = zmq_socket(g_Zmq.pContext, ZMQ_DEALER);
    if (!g_Zmq.pSocket)
    {
        Log_Write("[ZeroMQServer] 创建套接字失败!");
        zmq_ctx_destroy(g_Zmq.pContext);
        return -1;
    }

    // 设置套接字标识
    zmq_setsockopt(g_Zmq.pSocket, ZMQ_ROUTING_ID, ZMQ_SERVER_NAME, strlen(ZMQ_SERVER_NAME));

    // 3. 绑定到端口
    int rc = zmq_bind(g_Zmq.pSocket, "tcp://*:5555");
    if (rc != 0)
    {
        Log_Write("[ZeroMQServer] 绑定端口失败:%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);

    // 4. 启动服务端线程
    pthread_create(&g_Zmq.hThread, NULL, Zmq_ServerThread, NULL);

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

/*====================== 清理 ZeroMQ ======================*/
void Zmq_Cleanup(void)
{
    if (!g_Zmq.bRunning)
    {
        return;
    }

    Log_Write("[ZeroMQServer] 正在关闭服务...");

    // 1. 停止线程
    g_Zmq.bRunning = 0;

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

    // 2. 关闭套接字
    if (g_Zmq.pSocket)
    {
        zmq_close(g_Zmq.pSocket);
        g_Zmq.pSocket = NULL;
    }

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

    // 4. 销毁互斥锁
    pthread_mutex_destroy(&g_ZmqMutex);

    Log_Write("[ZeroMQServer] 已安全关闭");
}

/*====================== 设置业务命令处理器 ======================*/
void Zmq_SetCommandHandler(Zmq_CommandHandler handler)
{
    g_CommandHandler = handler;
    Log_Write("[ZeroMQServer] 业务命令处理器已设置");
}

/*====================== 向客户端发送消息通知 ======================*/
void Zmq_SendMessage(const char *szMessage)
{
    if (!g_Zmq.bRunning || !g_Zmq.pSocket)
    {
        return;
    }
    // ZeroMQServer.c (通信层) - 只负责发送
    //向客户端发送消息
    pthread_mutex_lock(&g_ZmqMutex);
    zmq_send(g_Zmq.pSocket, szMessage, strlen(szMessage), 0);
    pthread_mutex_unlock(&g_ZmqMutex);

    Log_Write("[ZeroMQServer] 发送:%s", szMessage);
}

/*====================== JSON 解析辅助函数外部声明 ======================*/
/**
 * 从收到的消息 JSON 字符串中提取指定 key 的值
 * 此函数在 ShPbxServer.c 中实现
 */
extern const char* ParseMsgJsonValue(const char *szJson, const char *szKey, char *szBuffer, int nBufSize);

调用方法

初始化ZeroMQ

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

退出时 关闭 ZeroMQ 服务

 // 1. 先关闭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 (通信层) - 只负责发送

接收消息

/*====================== 内部函数声明 ======================*/
static int Zmq_HandleClientMessage(const char *szMsg, char *szReply, int nReplySize);
/**
 * 处理 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 客户端封装类 - 负责与 ShPbxServer 服务端通信
 * 支持断线自动重连、心跳检测、连接状态监听
 */
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;  // 初始重连延迟(毫秒)
    private static final int MAX_RECONNECT_DELAY = 30000;  // 最大重连延迟 30 秒

    // 心跳配置
    private static final int HEARTBEAT_INTERVAL = 60000;  // 心跳间隔 60 秒
    private static final int HEARTBEAT_TIMEOUT = 180000;   // 超时时间 180 秒(3 次无响应)

    // 最后通信时间
    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;

    /**
     * 构造函数 - 初始化 ZeroMQ 上下文和连接
     * @param host 服务器地址
     * @param port 服务器端口
     */
    public ZeroMQClient(String host, int port) {
        this.host = host;
        this.port = port;
        this.context = new ZContext();
        connect();
    }

    /**
     * 设置消息监听器
     * @param listener 消息监听器
     */
    public void setMessageListener(MessageListener listener) {
        this.listener = listener;
    }

    /**
     * 设置连接状态监听器
     * @param listener 连接状态监听器
     */
    public void setConnectionListener(ConnectionListener listener) {
        this.connectionListener = listener;
    }

    /**
     * 建立连接
     */
    private void connect() {
        try {
            if (dealerSocket != null) {
                dealerSocket.close();
            }
            dealerSocket = context.createSocket(SocketType.DEALER);
            dealerSocket.connect(String.format("tcp://%s:%d", host, port));
            isConnected.set(true);
            lastActivityTime = System.currentTimeMillis();
            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();
            }
        }
    }

    /**
     * 发送消息到服务器
     * @param message 要发送的消息
     * @return 是否发送成功
     */
    public boolean send(String message) {
        if (!isConnected.get() || dealerSocket == null) {
            logger.warn("发送失败:未连接");
            return false;
        }
        try {
            dealerSocket.send(message.getBytes(ZMQ.CHARSET), 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() {
        dealerSocket.setReceiveTimeOut(1000);
        while (isRunning && !Thread.currentThread().isInterrupted()) {
            try {
                byte[] msgBytes = dealerSocket.recv();
                if (msgBytes != null) {
                    String message = new String(msgBytes, ZMQ.CHARSET);
                    lastActivityTime = System.currentTimeMillis();
                    
                    // 收到 PONG 响应,更新连接状态
                    if ("PONG".equals(message)) {
                        isConnected.set(true);
                        reconnectDelay = 1000;  // 重置重连延迟
                    }
                    
                    if (listener != null) {
                        listener.onMessageReceived(message);
                    }
                }
            } catch (ZMQException e) {
                if (e.getErrorCode() != ZMQ.Error.EAGAIN.getCode()) {
                    logger.error("ZeroMQ 接收消息异常", 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 elapsed = System.currentTimeMillis() - lastActivityTime;
                if (elapsed > HEARTBEAT_TIMEOUT) {
                    // 超时,视为断开
                    if (isConnected.get()) {
                        logger.warn("心跳超时({}ms),判定为断开连接", elapsed);
                        isConnected.set(false);
                        if (connectionListener != null) {
                            connectionListener.onDisconnected();
                        }
                    }
                    // 触发重连
                    retryCount++;
                    if (connectionListener != null) {
                        connectionListener.onReconnecting(retryCount);
                    }
                    reconnect();
                } else if (!isConnected.get()) {
                    // 未连接状态,尝试重连
                    retryCount++;
                    if (connectionListener != null) {
                        connectionListener.onReconnecting(retryCount);
                    }
                    reconnect();
                } else {
                    // 连接正常,发送心跳
                    send("PING");
                    logger.debug("发送心跳 PING");
                }
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
                break;
            } catch (Exception e) {
                logger.error("心跳检测异常", e);
            }
        }
    }

    /**
     * 重连逻辑(带退避策略)
     */
    private void reconnect() {
        logger.info("开始重连... (延迟 {}ms)", reconnectDelay);
        try {
            Thread.sleep(reconnectDelay);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            return;
        }
        
        // 指数退避:每次延迟翻倍,最大不超过 MAX_RECONNECT_DELAY
        reconnectDelay = Math.min(reconnectDelay * 2, MAX_RECONNECT_DELAY);
        
        connect();
    }

    /**
     * 检查当前连接状态
     * @return 是否已连接
     */
    public boolean isConnected() {
        return isConnected.get();
    }

    /**
     * 停止接收
     */
    public void stopReceive() {
        isRunning = false;
        if (heartbeatThread != null) {
            heartbeatThread.interrupt();
        }
    }

    /**
     * 关闭连接,释放资源(使用 try-with-resources 确保资源正确关闭)
     */
    public void close() {
        logger.info("[ZeroMQ 客户端] 正在关闭...");
        
        // 1. 停止接收线程
        stopReceive();
        
        // 2. 关闭 ZeroMQ 资源
        if (dealerSocket != null) {
            try {
                dealerSocket.close();
                logger.debug("[ZeroMQ 客户端] DEALER 套接字已关闭");
            } catch (Exception e) {
                logger.error("[ZeroMQ 客户端] 关闭 DEALER 套接字失败", e);
            }
            dealerSocket = null;
        }
        
        if (context != null) {
            try {
                context.close();
                logger.debug("[ZeroMQ 客户端] ZeroMQ 上下文已销毁");
            } catch (Exception e) {
                logger.error("[ZeroMQ 客户端] 关闭 ZeroMQ 上下文失败", e);
            }
            context = null;
        }
        
        // 3. 更新连接状态
        isConnected.set(false);
        
        logger.info("[ZeroMQ 客户端] 已安全关闭");
    }
}

调用方法

初始化

private static ZeroMQClient zmqClient;
 String host ="192.168.1.122";
 int port = 5555;

        // 初始化 ZeroMQ 客户端(命令通道)
        zmqClient = new ZeroMQClient(host, port);

        // 设置消息监听器
        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);
            }
        });

        // 启动接收线程
        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 失败(客户端未连接或底层异常)");
}