找回密码
立即注册
搜索
热搜: Java Python Linux Go
发回帖 发新帖

4159

积分

0

好友

537

主题
发表于 1 小时前 | 查看: 4| 回复: 0

一、ZeroMQ 介绍

ZeroMQ(ØMQ / 0MQ / zmq)是一个高性能异步消息库。它看起来像一个可嵌入的网络库,但用起来却像一个并发框架。你知道它为什么叫 "Zero" 吗?因为它追求零 Broker(无中心服务器)、零延迟、零管理成本和尽可能低的学习门槛。

不是 RabbitMQ 或 Kafka 那样的独立消息中间件,而是一个直接嵌入到你程序里的库。它的 Socket API 风格和 BSD Socket 很像,但一个 zmq socket 背后可以自动管理多条连接、自动重连、消息排队与异步 I/O。

特性一览

特性 说明
协议 ZMTP(ZeroMQ Message Transport Protocol)
消息模型 面向消息帧(frame),而非字节流,支持多帧消息
传输层 tcp://(跨机器)、ipc://(同机进程间,Unix socket)、inproc://(同进程线程间)、pgm://epgm://(多播)、ws://(WebSocket)
性能 单机可达百万级 msg/s,微秒级延迟
异步 I/O 后台 I/O 线程自动处理收发队列与重连
线程模型 socket 非线程安全(必须由创建它的线程使用,但 context 是线程安全的)
安全 CurveZMQ(基于 Curve25519 的加密认证)、ZAP 认证协议

消息模式(Messaging Patterns)

模式 Socket 对 用途
请求-应答 REQREP 严格一问一答的 RPC
发布-订阅 PUBSUB 广播,订阅者按前缀过滤。注意:慢订阅者会被丢消息
管道/流水线 PUSHPULL 任务分发(负载均衡)与结果收集
异步请求-应答 DEALERROUTER 无锁步限制的异步 RPC、代理、负载均衡 broker
排他配对 PAIRPAIR 主要用于 inproc:// 的线程间信令
扩展订阅 XPUB / XSUB 可感知订阅消息的 PUB/SUB,专用于构建代理

一个关键认知:bind/connect 与 client/server 无关。 通常是“地址稳定的一端 bind,动态的一端 connect”。真正厉害的是,connect 可以先于 bind 发起——ZeroMQ 会自动重连,这完全区别于原生 Socket。

二、核心概念梳理

1. Context(上下文)

  • 通过 zmq_ctx_new() 创建,负责管理所有 socket 和后台 I/O 线程。
  • 一个进程通常只需要一个 context
  • context 是线程安全的,但 socket 不是。

2. Message(消息)

  • ZeroMQ 传输的基本单位是帧(frame),一条消息可以包含多帧(multipart message)。
  • 多帧消息是原子的:要么全部收到,要么全都收不到。
  • ROUTER socket 收到消息时会自动添加一帧“来源身份(identity)”,回复时也必须带上这一帧以完成路由。

3. 高水位线(HWM, High-Water Mark)

  • 每个 socket 都有收发队列上限(默认 1000 条)。
  • 队列满了怎么办?行为取决于 socket 类型:PUB丢弃新消息,而 PUSH/REQ/DEALER 则会阻塞(或是在设置了 ZMQ_DONTWAIT 时返回 EAGAIN)。
  • 可以通过 ZMQ_SNDHWMZMQ_RCVHWM 选项调整。

4. 架构分层

┌─────────────────────────────────────┐
│   C 应用程序                        │
├─────────────────────────────────────┤
│  CZMQ(高级 C API)                 │  ← 更好用
├─────────────────────────────────────┤
│  libzmq(核心底层 API)             │  ← 必须装
├─────────────────────────────────────┤
│  传输层(TCP / IPC / inproc / PGM) │
└─────────────────────────────────────┘

CZMQ 依赖 libzmq,但 libzmq 不依赖 CZMQ。

三、libzmq 与 CZMQ 的关系与区别

C 语言绑定页,官方给 C 开发者两个选择:

  1. CZMQ — 高级 C 绑定,官方推荐。额外提供了 poller、线程管理、安全辅助等。
  2. libzmq — 核心底层库,提供最原始的 C API(zmq.h)。

libzmq(核心引擎)

  • 项目地址:https://github.com/zeromq/libzmq
  • 文档:https://libzmq.readthedocs.io/
  • C++ 实现,但提供纯 C API(zmq.h)。
  • 实现了所有 socket 类型、ZMTP 协议和消息收发引擎。
  • 痛点:API 偏底层,略显啰嗦。zmq_msg_t 需要手动 init/close,没有事件循环,多线程用起来比较复杂。

CZMQ(高级封装)

核心组件概览:

组件 作用
zsock_t 封装 socket,自动管理 context、错误与关闭
zstr_* 收发字符串,一行搞定
zmsg_t / zframe_t 自动管理多帧消息的生命周期
zloop_t 事件循环(风格类似 libev/libuv)
zpoller_t 多 socket 轮询
zactor_t 将线程与 PAIR 管道打包为 actor 模型
zauth_t / zcert_t 认证与 Curve 加密封装
zhashx_t / zlistx_t 通用容器
zbeacon / zgossip UDP 服务发现 / gossip 协议

详细区别对照

维度 libzmq CZMQ
定位 核心引擎 高级封装
API 风格 底层、显式资源管理 高层、自动管理
消息处理 zmq_msg_t 手动 init/close zmsg_t / zstr_* 自动
事件循环 zloop_t
Actor 模型 zactor_t
认证/加密 手动配置 ZAP zauth_t 封装
依赖 仅 libstdc++ 等 依赖 libzmq
学习曲线 较陡 平缓
适合场景 极致控制、嵌入式裁剪、零拷贝 快速开发、工程化

经验之谈:99% 的 C 语言 ZeroMQ 项目都应该用 CZMQ。 它能让你少写一半样板代码,并免费获得事件循环、actor 模型等生产级特性。只有在资源极度受限,或需要用 zmq_msg_t 做零拷贝时,才考虑原生 libzmq。C++ 项目则可以选择 cppzmq (https://github.com/zeromq/cppzmq)。

四、常用 API

4.1 libzmq 常用 API

上下文管理

void *zmq_ctx_new (void);                         // 创建 context
int zmq_ctx_term (void *ctx);                     // 销毁(阻塞直到所有 socket 关闭)
int zmq_ctx_set (void *ctx, int option, int v);  // 如 ZMQ_IO_THREADS
int zmq_ctx_shutdown (void *ctx);                // 令阻塞调用返回 ETERM

Socket 生命周期

void *zmq_socket (void *ctx, int type);  // ZMQ_REQ/REP/PUB/SUB/PUSH/PULL/DEALER/ROUTER/PAIR...
int zmq_bind (void *s, const char *endpoint);     // "tcp://*:5555" "ipc:///tmp/x" "inproc://a"
int zmq_connect (void *s, const char *endpoint);
int zmq_close (void *s);
int zmq_setsockopt (void *s, int option, const void *val, size_t len);
int zmq_getsockopt (void *s, int option, void *val, size_t *len);

常用 socket 选项:

选项 说明
ZMQ_SUBSCRIBE / ZMQ_UNSUBSCRIBE SUB 订阅前缀("" = 订阅全部)
ZMQ_SNDHWM / ZMQ_RCVHWM 收发高水位线(默认 1000)
ZMQ_LINGER close 时未发消息的滞留时间(ms),建议设为 0 或有限值
ZMQ_SNDTIMEO / ZMQ_RCVTIMEO 收发超时(ms),-1 为永久阻塞
ZMQ_IDENTITY (ZMQ_ROUTING_ID) DEALER/REQ 的路由身份
ZMQ_TCP_KEEPALIVE TCP 保活
ZMQ_RECONNECT_IVL 自动重连间隔(ms)

简单收发(字节缓冲区)

int zmq_send (void *s, const void *buf, size_t len, int flags); // flags: 0 / ZMQ_DONTWAIT / ZMQ_SNDMORE
int zmq_recv (void *s, void *buf, size_t len, int flags);       // 超长会被截断!返回实际消息长度

消息对象收发(零拷贝 / 多帧)

int zmq_msg_init (zmq_msg_t *msg);
int zmq_msg_init_size (zmq_msg_t *msg, size_t size);
int zmq_msg_init_data (zmq_msg_t *msg, void *data, size_t size,
                          zmq_free_fn *ffn, void *hint);       // 零拷贝
void  *zmq_msg_data (zmq_msg_t *msg);
size_t zmq_msg_size (zmq_msg_t *msg);
int zmq_msg_send (zmq_msg_t *msg, void *s, int flags);
int zmq_msg_recv (zmq_msg_t *msg, void *s, int flags);
int zmq_msg_more (zmq_msg_t *msg);                            // 是否还有后续帧
int zmq_msg_close (zmq_msg_t *msg);

多路复用与代理

typedef struct { void *socket; int fd; short events; short revents; } zmq_pollitem_t;
int zmq_poll (zmq_pollitem_t *items, int nitems, long timeout_ms);  // events: ZMQ_POLLIN/ZMQ_POLLOUT

int zmq_proxy (void *frontend, void *backend, void *capture);      // 内置代理(如 XSUB<->XPUB)

错误处理

int zmq_errno (void);          // EAGAIN / ETERM / EINTR ...
const char *zmq_strerror (int errnum);
void zmq_version (int *major, int *minor, int *patch);

4.2 CZMQ 常用 API

zsock —— socket 封装

zsock_t *zsock_new_pub (const char *endpoint);   // "@tcp://*:5555"  @ 表示 bind
zsock_t *zsock_new_sub (const char *endpoint, const char *subscribe); // > 表示 connect(默认)
zsock_t *zsock_new_req / _rep / _push / _pull / _dealer / _router / _pair (const char *endpoint);
void zsock_destroy (zsock_t **self_p);
int zsock_bind (zsock_t *self, const char *format, ...);
int zsock_connect (zsock_t *self, const char *format, ...);
void     zsock_set_sndhwm / zsock_set_rcvtimeo / zsock_set_linger (zsock_t *, int);

endpoint 前缀约定:@ 代表 bind,> 代表 connect。例如 "@tcp://*:5555"">tcp://192.168.1.10:5555"

zstr —— 字符串收发

int zstr_send (void *dest, const char *string);
int zstr_sendx (void *dest, const char *s1, ..., NULL);  // 多帧字符串
char *zstr_recv (void *source);                            // 需 zstr_free 释放
int zstr_recvx (void *source, char **s1, ..., NULL);
void zstr_free (char **string_p);

zmsg / zframe —— 多帧消息

zmsg_t   *zmsg_new (void);
int zmsg_addstr (zmsg_t *self, const char *string);
int zmsg_addmem (zmsg_t *self, const void *data, size_t size);
int zmsg_send (zmsg_t **self_p, void *dest);    // 发送后自动销毁
zmsg_t   *zmsg_recv (void *source);
char     *zmsg_popstr (zmsg_t *self);
zframe_t *zmsg_pop (zmsg_t *self);
void zmsg_destroy (zmsg_t **self_p);

zframe_t *zframe_new (const void *data, size_t size);
byte     *zframe_data (zframe_t *self);
size_t zframe_size (zframe_t *self);

zpoller —— 多 socket 轮询

zpoller_t *zpoller_new (void *reader, ...);           // NULL 结尾
void      *zpoller_wait (zpoller_t *self, int timeout_ms); // 返回就绪 socket
bool zpoller_expired (zpoller_t *self);          // 超时?
bool zpoller_terminated (zpoller_t *self);       // 收到中断(Ctrl-C)?
void zpoller_destroy (zpoller_t **self_p);

zloop —— 事件循环

zloop_t *zloop_new (void);
int zloop_reader (zloop_t *self, zsock_t *sock,
                       zloop_reader_fn handler, void *arg); // socket 可读回调
int zloop_timer (zloop_t *self, size_t delay_ms, size_t times,
                       zloop_timer_fn handler, void *arg);   // 定时器
int zloop_start (zloop_t *self);                       // 阻塞运行
void zloop_destroy (zloop_t **self_p);
// 回调返回 0 继续循环,返回 -1 退出循环

zactor —— actor 并发模型

typedef void (zactor_fn) (zsock_t *pipe, void *args);  // actor 线程函数
zactor_t *zactor_new (zactor_fn task, void *args);      // 创建线程 + PAIR 管道
void zactor_destroy (zactor_t **self_p);           // 发 "$TERM" 并 join
// 与 actor 通信:zstr_send(actor, "CMD") / zstr_recv(actor)
// actor 内部:必须先 zsock_signal(pipe, 0) 表示就绪

五、ARM Linux 移植(交叉编译)

以常见的 arm-linux-gnueabihf 工具链为例(32 位 ARM,如 i.MX6/AM335x/树莓派)。64 位平台(如 RK3588、树莓派 4 64 位),将三元组换成 aarch64-linux-gnu 即可。

5.0 准备工具链

# Ubuntu 上安装通用交叉工具链(或使用芯片厂商 SDK 里的工具链)
sudo apt install gcc-arm-linux-gnueabihf g++-arm-linux-gnueabihf
# 64 位:
sudo apt install gcc-aarch64-linux-gnu g++-aarch64-linux-gnu

# 统一变量(下文使用)
export HOST=arm-linux-gnueabihf          # 或 aarch64-linux-gnu
export PREFIX=$HOME/zmq-arm              # 交叉编译产物安装目录
export CC=${HOST}-gcc
export CXX=${HOST}-g++

5.1 交叉编译 libzmq

git clone --depth 1 https://github.com/zeromq/libzmq.git
# 或者使用 Gitee 镜像
git clone --depth 1 https://gitee.com/mirrors/libzmq.git

cd libzmq
./autogen.sh          # 需要 autoconf automake libtool pkg-config

./configure \
    --host=$HOST \
    --prefix=$PREFIX \
    --enable-static \
    --disable-shared=no \
    --without-docs \
    --disable-Werror \
    --disable-curve         # 若不需要加密可关闭,省去 libsodium 依赖
# 需要加密则改为: --with-libsodium 并先交叉编译 libsodium

make -j$(nproc)
make install
# 产物: $PREFIX/lib/libzmq.so / libzmq.a 和 $PREFIX/include/zmq.h

也可用 CMake:

cmake -B build \
  -DCMAKE_SYSTEM_NAME=Linux -DCMAKE_SYSTEM_PROCESSOR=arm \
  -DCMAKE_C_COMPILER=${HOST}-gcc -DCMAKE_CXX_COMPILER=${HOST}-g++ \
  -DCMAKE_INSTALL_PREFIX=$PREFIX -DWITH_DOCS=OFF -DENABLE_CURVE=OFF
cmake --build build -j && cmake --install build

libzmq交叉编译构建过程终端输出,显示多个目标的编译进度和链接成功信息

5.2 交叉编译 CZMQ

git clone --depth 1 https://github.com/zeromq/czmq.git
cd czmq
./autogen.sh

# 让 configure 找到刚才交叉编译的 libzmq
export PKG_CONFIG_PATH=$PREFIX/lib/pkgconfig
# export PKG_CONFIG_SYSROOT_DIR=/            # 视工具链 sysroot 而定

./configure \
    --host=$HOST \
    --prefix=$PREFIX \
    --with-libzmq=$PREFIX \
    --without-docs

make -j$(nproc)
make install
# 产物: $PREFIX/lib/libczmq.so 和 $PREFIX/include/czmq.h

或者使用 CMake

# 0. 前置条件:libzmq 已交叉编译安装到 $PREFIX
#    $PREFIX/lib/libzmq.so 和 $PREFIX/include/zmq.h 必须存在

# 1. 拉取源码
git clone --depth 1 https://github.com/zeromq/czmq.git
cd czmq

# 2. 清理旧构建(保险)
rm -rf build-arm

# 3. CMake 配置
cmake -S . -B build-arm \
    -DCMAKE_SYSTEM_NAME=Linux \
    -DCMAKE_SYSTEM_PROCESSOR=arm \
    -DCMAKE_C_COMPILER=${HOST}-gcc \
    -DCMAKE_CXX_COMPILER=${HOST}-g++ \
    -DCMAKE_INSTALL_PREFIX=$(pwd)/czmq-arm-install \
    -DCMAKE_PREFIX_PATH=$PREFIX \
    -DCMAKE_BUILD_TYPE=Release \
    -DBUILD_TESTING=OFF \
    -DBUILD_DOC=OFF \
    -DBUILD_STATIC=ON \
    -DBUILD_SHARED=ON \
    -DWITH_LIBCURL=OFF \
    -DWITH_LIBMICROHTTPD=OFF \
    -DWITH_LZ4=OFF \
    -DWITH_SYSTEMD=OFF \
    -DWITH_NSS=OFF

# 注意 BUILD_DOC 这个变量 CZMQ 的 CMakeLists.txt 里根本没有定义。CZMQ 的文档构建只在 autotools 体系(Makefile.am + asciidoc)里存在,CMake 构建系统不提供文档开关。

# 4. 并行编译
cmake --build build-arm -j$(nproc)

# 5. 安装到 ./czmq-arm-install
cmake --install build-arm

CZMQ交叉编译CMake构建过程终端截图,显示多个目标编译进度和微信公众号水印

5.3 编译你的应用

${HOST}-gcc main.c \
    -I$PREFIX/include \
    -L$PREFIX/lib \
    -lczmq -lzmq -lpthread -lstdc++ -lm \
    -o app_arm

# 静态链接(部署最省事,但体积较大):
${HOST}-gcc main.c -I$PREFIX/include \
    $PREFIX/lib/libczmq.a $PREFIX/lib/libzmq.a \
    -lpthread -lstdc++ -lm -static-libgcc -o app_arm_static

5.4 部署到目标板

# 动态链接时需同时拷贝 .so 库
scp app_arm root@<board-ip>:/usr/local/bin/
scp $PREFIX/lib/libzmq.so.5 $PREFIX/lib/libczmq.so.4 root@<board-ip>:/usr/lib/

# 板上运行
./app_arm
# 若提示找不到库: export LD_LIBRARY_PATH=/usr/lib:$LD_LIBRARY_PATH 或 ldconfig

5.5 嵌入式构建系统集成

构建系统 方式
Buildroot make menuconfig → Target packages → Libraries → Networking → 勾选 zeromqczmq
Yocto 在镜像配方添加 IMAGE_INSTALL:append = " zeromq czmq"(meta-oe 层提供)
Debian/Ubuntu ARM 板 (树莓派等) 直接 sudo apt install libzmq3-dev libczmq-dev

5.6 嵌入式裁剪与注意事项

  • --disable-curve --without-docs 可显著减小体积,再配合 -Osstrip 使用。
  • libzmq 是 C++ 实现,必须链接 -lstdc++(或用 ${HOST}-g++ 链接)。
  • 资源受限设备注意调低 ZMQ_SNDHWM/RCVHWM,防止队列吃光内存。
  • ZeroMQ 需要 pthread;确认工具链 libc 支持(glibc/musl 均可,uclibc 需较新版本)。
  • 弱网或现场环境,建议设置 ZMQ_TCP_KEEPALIVE=1、合理的 ZMQ_RECONNECT_IVL,并充分利用 ZeroMQ 的自动重连特性——你再也不用自己写重连逻辑了。

六、完整示例代码

6.1 REQ/REP —— Hello World(libzmq 原生,官方示例)

server.c

// gcc server.c -lzmq -o server
#include <zmq.h>
#include <stdio.h>
#include <unistd.h>
#include <string.h>
#include <assert.h>

int main (void)
{
    void *context = zmq_ctx_new ();
    void *responder = zmq_socket (context, ZMQ_REP);
    int rc = zmq_bind (responder, "tcp://*:5555");
    assert (rc == 0);

    while (1) {
        char buffer [10];
        zmq_recv (responder, buffer, 10, 0);
        printf ("Received Hello\n");
        sleep (1);                     // 模拟处理耗时
        zmq_send (responder, "World", 5, 0);
    }
    zmq_close (responder);
    zmq_ctx_destroy (context);
    return 0;
}

黑色终端界面显示服务端程序输出,列出多行"Received Hello"和公众号水印

client.c

// gcc client.c -lzmq -o client
#include <zmq.h>
#include <string.h>
#include <stdio.h>
#include <unistd.h>

int main (void)
{
    printf ("Connecting to hello world server…\n");
    void *context = zmq_ctx_new ();
    void *requester = zmq_socket (context, ZMQ_REQ);
    zmq_connect (requester, "tcp://localhost:5555");

    for (int request_nbr = 0; request_nbr != 10; request_nbr++) {
        char buffer [10];
        printf ("Sending Hello %d…\n", request_nbr);
        zmq_send (requester, "Hello", 5, 0);
        zmq_recv (requester, buffer, 10, 0);
        printf ("Received World %d\n", request_nbr);
    }
    zmq_close (requester);
    zmq_ctx_destroy (context);
    return 0;
}

黑色终端界面显示客户端与服务端通信日志,循环发送"Hello"和接收"World"消息

6.2 PUB/SUB —— 数据广播(CZMQ)

publisher.c

// gcc publisher.c -lczmq -lzmq -o publisher
#include <czmq.h>

int main (void)
{
    zsock_t *pub = zsock_new_pub ("@tcp://*:5556");   // @ = bind
    assert (pub);

    int seq = 0;
    while (!zsys_interrupted) {                        // Ctrl-C 优雅退出
        // 第一帧作为 topic,SUB 端按前缀过滤
        zstr_sendf (pub, "sensor.temp %d %.2f", seq++, 25.0 + (rand () % 100) / 10.0);
        zclock_sleep (1000);                           // 1s
    }
    zsock_destroy (&pub);
    return 0;
}

subscriber.c

// gcc subscriber.c -lczmq -lzmq -o subscriber
#include <czmq.h>

int main (void)
{
    // 只订阅 "sensor.temp" 前缀;传 "" 则订阅全部
    zsock_t *sub = zsock_new_sub (">tcp://localhost:5556", "sensor.temp");
    assert (sub);

    while (!zsys_interrupted) {
        char *msg = zstr_recv (sub);
        if (!msg) break;                               // 被中断
        printf ("Got: %s\n", msg);
        zstr_free (&msg);
    }
    zsock_destroy (&sub);
    return 0;
}

黑色终端界面显示订阅者接收到的传感器温度数据日志

6.3 PUSH/PULL —— 任务分发流水线(CZMQ,Producer/Consumer)

producer.c

// gcc producer.c -lczmq -lzmq -o producer
#include <czmq.h>

int main (void)
{
    zsock_t *push = zsock_new_push ("@tcp://*:5557");
    assert (push);

    for (int task = 0; task < 100 && !zsys_interrupted; task++) {
        zstr_sendf (push, "task-%03d", task);          // 自动负载均衡到各 worker
        printf ("Dispatched task-%03d\n", task);
        zclock_sleep (10);
    }
    zstr_send (push, "END");
    zsock_destroy (&push);
    return 0;
}

worker.c

// gcc worker.c -lczmq -lzmq -o worker   (可启动多个实例)
#include <czmq.h>

int main (void)
{
    zsock_t *pull = zsock_new_pull (">tcp://localhost:5557");
    assert (pull);

    while (!zsys_interrupted) {
        char *task = zstr_recv (pull);
        if (!task) break;
        if (streq (task, "END")) { zstr_free (&task); break; }
        printf ("[worker %d] processing %s\n", getpid (), task);
        zclock_sleep (100);                            // 模拟耗时
        zstr_free (&task);
    }
    zsock_destroy (&pull);
    return 0;
}

6.4 zloop 事件循环 —— 多 socket + 定时器(CZMQ)

// gcc loop_demo.c -lczmq -lzmq -o loop_demo
#include <czmq.h>

static int on_sub_msg (zloop_t *loop, zsock_t *reader, void *arg)
{
    char *msg = zstr_recv (reader);
    printf ("SUB got: %s\n", msg);
    zstr_free (&msg);
    return 0;                      // 返回 -1 退出循环
}

static int on_timer (zloop_t *loop, int timer_id, void *arg)
{
    printf ("heartbeat @ %ld\n", (long) zclock_mono ());
    return 0;
}

int main (void)
{
    zsock_t *sub = zsock_new_sub (">tcp://localhost:5556", "");
    zloop_t *loop = zloop_new ();

    zloop_reader (loop, sub, on_sub_msg, NULL);  // socket 可读事件
    zloop_timer (loop, 1000, 0, on_timer, NULL);  // 每 1000ms,0 = 无限次

    zloop_start (loop);                           // 阻塞,Ctrl-C 退出

    zloop_destroy (&loop);
    zsock_destroy (&sub);
    return 0;
}

6.5 zactor —— 多线程 actor 模型(CZMQ)

// gcc actor_demo.c -lczmq -lzmq -o actor_demo
#include <czmq.h>

// actor 线程:通过 PAIR 管道与主线程通信,同时可拥有自己的 socket
static void worker_actor (zsock_t *pipe, void *args)
{
    zsock_signal (pipe, 0);                       // 必须:通知主线程"已就绪"

    zsock_t *pull = zsock_new_pull (">tcp://localhost:5557");
    zpoller_t *poller = zpoller_new (pipe, pull, NULL);
    bool terminated = false;

    while (!terminated) {
        void *which = zpoller_wait (poller, -1);
        if (which == pipe) {                      // 主线程发来的命令
            char *cmd = zstr_recv (pipe);
            if (streq (cmd, "$TERM"))             // zactor_destroy 会发送 $TERM
                terminated = true;
            else if (streq (cmd, "STATUS"))
                zstr_send (pipe, "alive");
            zstr_free (&cmd);
        }
        else if (which == pull) {                 // 业务数据
            char *task = zstr_recv (pull);
            printf ("[actor] %s\n", task);
            zstr_free (&task);
        }
        else if (zpoller_terminated (poller))
            break;
    }
    zpoller_destroy (&poller);
    zsock_destroy (&pull);
}

int main (void)
{
    zactor_t *actor = zactor_new (worker_actor, NULL);  // 创建线程

    zstr_send (actor, "STATUS");                        // 向 actor 发命令
    char *reply = zstr_recv (actor);
    printf ("actor says: %s\n", reply);
    zstr_free (&reply);

    zclock_sleep (5000);                                // 让 actor 干 5 秒活

    zactor_destroy (&actor);                            // 发 $TERM 并 join 线程
    return 0;
}

6.6 官方最小示例 —— inproc PUSH/PULL(CZMQ,摘自 zeromq.org)

#include <czmq.h>
int main (void)
{
    zsock_t *push = zsock_new_push ("inproc://example");
    zsock_t *pull = zsock_new_pull ("inproc://example");

    zstr_send (push, "Hello, World");

    char *string = zstr_recv (pull);
    puts (string);
    zstr_free (&string);

    zsock_destroy (&pull);
    zsock_destroy (&push);
    return 0;
}

6.7 Makefile(本机 + ARM 交叉编译二合一)

# 用法:
#   make            → 本机编译
#   make CROSS=arm-linux-gnueabihf- PREFIX=$(HOME)/zmq-arm  → ARM 交叉编译
CROSS   ?=
PREFIX  ?= /usr
CC      := $(CROSS)gcc
CFLAGS  := -O2 -Wall -I$(PREFIX)/include
LDFLAGS := -L$(PREFIX)/lib -lczmq -lzmq -lpthread -lstdc++ -lm

BINS := server client publisher subscriber producer worker loop_demo actor_demo

all: $(BINS)

%: %.c
    $(CC) $(CFLAGS) $< $(LDFLAGS) -o $@

clean:
    rm -f $(BINS)

七、选型建议

需求场景 推荐
C 项目、快速开发、需要事件循环/多线程 CZMQ
资源极度受限嵌入式、需要零拷贝控制 libzmq 原生
C++ 项目 cppzmq(header-only)或 libzmq
单机多线程、追求极限延迟、不跨进程 无锁队列(如 Concurrency Kit)更快
同机跨进程、简单场景 POSIX MQ 也可,但无自动重连/跨机能力
未来可能跨进程/跨机器扩展 ZeroMQ 性价比最高(同一套 API 换个 endpoint 即从 inproc → ipc → tcp)

与 POSIX MQ / 无锁队列的定位对比:

方案 跨进程 跨机器 线程安全 性能 复杂度
POSIX MQ 部分
无锁队列 (CK) 极高
ZeroMQ (CZMQ/libzmq) actor 模型

八、参考链接

官方

libzmq

CZMQ

其他绑定与生态

九、附:单线程 / 多线程 / 多进程可运行示例(examples/ 目录)

examples/
├── Makefile                         # make 编译全部;支持 CROSS= ARM 交叉编译
├── single_thread/
│   ├── st_inproc_pushpull.c         # 单线程 inproc 自发自收(最小示例)
│   └── st_zpoller_server.c          # 单线程 zpoller 多路复用 REP+PULL+心跳
├── multi_thread/
│   ├── mt_libzmq_inproc.c           # libzmq + pthread:inproc 任务池(共享 context,socket 线程私有)
│   └── mt_czmq_zactor.c             # CZMQ zactor 线程池 + zpoller 收集结果(推荐写法)
└── multi_process/
    ├── mp_fork_pubsub.c             # fork + ipc:// 发布订阅(fork 后各进程再建 socket)
    ├── mp_ventilator.c              # 并行流水线:任务分发器(PUSH:5557)
    ├── mp_worker.c                  # 并行流水线:工作进程(可开多个,自动负载均衡)
    └── mp_sink.c                    # 并行流水线:结果收集器(PULL:5558,统计并行加速)

编译与运行:

cd examples && make

# 单线程
./single_thread/st_inproc_pushpull
./single_thread/st_zpoller_server          # Ctrl-C 退出

# 多线程
./multi_thread/mt_libzmq_inproc
./multi_thread/mt_czmq_zactor

# 多进程(fork 单文件版)
./multi_process/mp_fork_pubsub

# 多进程(流水线版,三个终端按序启动)
./multi_process/mp_sink
./multi_process/mp_worker                  # 可开多个
./multi_process/mp_ventilator              # 回车开始分发

三种并发模型的 ZeroMQ 要点总结:

模型 传输 关键规则
单线程 inproc:// 或任意 zpoller/zloop 多路复用,一个线程服务多个 socket
多线程 inproc:// 全进程共享 1 个 context;socket 线程私有;只传消息不共享内存
多进程 ipc://(同机)/ tcp://(跨机) context 不能跨 fork —— 必须 fork 后各进程自建 socket;换 endpoint 即从同机扩展到跨机

作为 云栈社区 的开发者,如果你正在为 分布式系统 寻找轻量级通信方案,ZeroMQ 是一个值得深入研究的利器。它灵活的 IPC 抽象和极低的上手门槛,使其在嵌入式与服务器领域都广受欢迎。希望这篇从 API 到交叉编译的详实指南,能帮你扫清上手与移植的障碍。




上一篇:嵌入式音频对讲全链路深度拆解:硬件链路、AEC与WebRTC弱网实战量产
下一篇:C++ 模板元编程实际用途:编译期类型安全与零开销优化
您需要登录后才可以回帖 登录 | 立即注册

手机版|小黑屋|网站地图|云栈社区 ( 苏ICP备2022046150号-2 )

GMT+8, 2026-8-4 05:39 , Processed in 1.128698 second(s), 42 queries , Gzip On.

Powered by Discuz! X3.5

© 2025-2026 云栈社区.

快速回复 返回顶部 返回列表