























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 失败(客户端未连接或底层异常)"); }
此内容由惯性聚合(RSS阅读器)自动聚合整理,仅供阅读参考。 原文来自 — 版权归原作者所有。