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

推荐订阅源

PCI Perspectives
PCI Perspectives
Threat Intelligence Blog | Flashpoint
Threat Intelligence Blog | Flashpoint
C
CXSECURITY Database RSS Feed - CXSecurity.com
P
Palo Alto Networks Blog
S
Schneier on Security
Scott Helme
Scott Helme
T
Threat Research - Cisco Blogs
K
Kaspersky official blog
Microsoft Azure Blog
Microsoft Azure Blog
T
The Exploit Database - CXSecurity.com
C
Cybersecurity and Infrastructure Security Agency CISA
T
Tenable Blog
G
GRAHAM CLULEY
OSCHINA 社区最新新闻
OSCHINA 社区最新新闻
Cyber Security Advisories - MS-ISAC
Cyber Security Advisories - MS-ISAC
I
Intezer
D
Docker
月光博客
月光博客
L
Lohrmann on Cybersecurity
Latest news
Latest news
B
Blog
罗磊的独立博客
M
MIT News - Artificial intelligence
S
Securelist
Know Your Adversary
Know Your Adversary
Help Net Security
Help Net Security
Recorded Future
Recorded Future
S
SegmentFault 最新的问题
N
Netflix TechBlog - Medium
T
Threatpost
H
Hacker News: Front Page
K
KPMG report finds enterprise disconnect between AI and its ROI | CIO
人人都是产品经理
人人都是产品经理
F
Fortinet All Blogs
博客园 - Franky
P
Proofpoint News Feed
大猫的无限游戏
大猫的无限游戏
Blog — PlanetScale
Blog — PlanetScale
有赞技术团队
有赞技术团队
博客园 - 【当耐特】
A
About on SuperTechFans
奇客Solidot–传递最新科技情报
奇客Solidot–传递最新科技情报
T
Tor Project blog
Google Online Security Blog
Google Online Security Blog
Application and Cybersecurity Blog
Application and Cybersecurity Blog
cs.AI updates on arXiv.org
cs.AI updates on arXiv.org
Engineering at Meta
Engineering at Meta
Webroot Blog
Webroot Blog
Security Archives - TechRepublic
Security Archives - TechRepublic
Microsoft Security Blog
Microsoft Security Blog

博客园 - 海乐学习

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项目中读取配置文件 创建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的DEALER双帧结构(路由空帧 + 业务帧)(支持异步收发、主动推送事件)C写服务端 java写客户端
海乐学习 · 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 失败(客户端未连接或底层异常)");
}