跳到主要内容
极客日志极客日志面向AI+效率的开发者社区
首页博客GitHub 精选镜像AI 生图工具UI配色美学隐私政策关于联系
搜索内容 / 工具 / 仓库 / 镜像...⌘K搜索
注册
博客列表
C

MQTT 通信协议 Mosquitto 发布订阅 C 语言实现示例

基于 Mosquitto 库的 C 语言 MQTT 客户端开发指南。涵盖同步与异步通信模式,展示发布订阅核心回调机制及线程安全处理。提供完整代码示例与编译配置,帮助开发者快速构建稳定可靠的物联网消息服务。

FrontendX发布于 2025/1/21更新于 2026/7/2534 浏览
MQTT 通信协议 Mosquitto 发布订阅 C 语言实现示例

MQTT 客户端开发:基于 Mosquitto 的 C 语言实践

在物联网开发中,MQTT 协议因其轻量级和发布/订阅模式而被广泛采用。Mosquitto 作为成熟的开源 MQTT 代理库,提供了丰富的 C API。本文将通过实际代码演示如何使用 libmosquitto 库实现同步与异步的发布订阅功能。

准备工作

编译安装 mosquitto 后,主要需要用到以下文件:

  • libmosquitto.so.1 (动态链接库)
  • mosquitto.h (头文件)

下面的示例均使用标准 C 语言编写,只需包含头文件并链接 -lmosquitto 即可。

同步模式

同步模式下,程序会阻塞等待网络事件。这种方式逻辑简单,适合对实时性要求不高或单线程场景。

订阅端 (sub.c)

订阅端主要负责连接服务器、订阅主题并监听消息回调。

#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include "mosquitto.h"

#define HOST "localhost"
#define PORT 1883
#define KEEP_ALIVE 60
#define MSG_MAX_SIZE 512

static int running = 1;

void my_connect_callback(struct mosquitto *mosq, void *obj, int rc) {
    printf("Call the function: on_connect\n");
    if(rc) {
        printf("on_connect error!\n");
        exit(1);
    } else {
        // 参数:句柄、id、订阅的主题、qos
        if(mosquitto_subscribe(mosq, NULL, "topic1", 2)) {
            printf("Set the topic error!\n");
            exit(1);
        }
    }
}

void my_disconnect_callback(struct mosquitto *mosq, void *obj, int rc) {
    printf("Call the function: my_disconnect_callback\n");
    running = 0;
}

void my_subscribe_callback(struct mosquitto *mosq, void *obj, int mid, int qos_count, const int *granted_qos) {
    printf("Call the function: on_subscribe\n");
}

void my_message_callback(struct mosquitto *mosq, void *obj, const struct mosquitto_message *msg) {
    printf("Call the function: on_message\n");
    printf("Receive a message of %s : %s\n", (char *)msg->topic, (char *)msg->payload);
    if(0 == strcmp(msg->payload, "quit")) {
        mosquitto_disconnect(mosq);
    }
}

int main() {
    int ret;
    struct mosquitto *mosq;

    // 初始化 mosquitto 库
    ret = mosquitto_lib_init();
    if(ret) {
        printf("Init lib error!\n");
        return -1;
    }

    // 创建一个订阅端实例
    // 参数:id(不需要则为 NULL)、clean_start、用户数据
    mosq = mosquitto_new("sub_test", true, NULL);
    if(mosq == NULL) {
        printf("New sub_test error!\n");
        mosquitto_lib_cleanup();
        return -1;
    }

    // 设置回调函数
    mosquitto_connect_callback_set(mosq, my_connect_callback);
    mosquitto_disconnect_callback_set(mosq, my_disconnect_callback);
    mosquitto_subscribe_callback_set(mosq, my_subscribe_callback);
    mosquitto_message_callback_set(mosq, my_message_callback);

    // 连接至服务器
    // 参数:句柄、ip(host)、端口、心跳
    ret = mosquitto_connect(mosq, HOST, PORT, KEEP_ALIVE);
    if(ret) {
        printf("Connect server error!\n");
        mosquitto_destroy(mosq);
        mosquitto_lib_cleanup();
        return -1;
    }

    // 开始通信:循环执行,直到运行标志 running 被改变
    printf("Start!\n");
    while(running) {
        mosquitto_loop(mosq, -1, 1);
    }

    // 结束后的清理工作
    mosquitto_destroy(mosq);
    mosquitto_lib_cleanup();
    printf("End!\n");
    return 0;
}

发布端 (pub.c)

发布端负责向指定主题发送消息。

#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <unistd.h>
#include "mosquitto.h"

#define HOST "localhost"
#define PORT 1883
#define KEEP_ALIVE 60
#define MSG_MAX_SIZE 512

static int running = 1;

void my_connect_callback(struct mosquitto *mosq, void *obj, int rc) {
    printf("Call the function: my_connect_callback\n");
}

void my_disconnect_callback(struct mosquitto *mosq, void *obj, int rc) {
    printf("Call the function: my_disconnect_callback\n");
    running = 0;
}

void my_publish_callback(struct mosquitto *mosq, void *obj, int mid) {
    printf("Call the function: my_publish_callback\n");
}

int main() {
    int ret;
    struct mosquitto *mosq;
    char buff[MSG_MAX_SIZE];

    // 初始化 libmosquitto 库
    ret = mosquitto_lib_init();
    if(ret) {
        printf("Init lib error!\n");
        return -1;
    }

    // 创建一个发布端实例
    mosq = mosquitto_new("pub_test", true, NULL);
    if(mosq == NULL) {
        printf("New pub_test error!\n");
        mosquitto_lib_cleanup();
        return -1;
    }

    // 设置回调函数
    mosquitto_connect_callback_set(mosq, my_connect_callback);
    mosquitto_disconnect_callback_set(mosq, my_disconnect_callback);
    mosquitto_publish_callback_set(mosq, my_publish_callback);

    // 连接至服务器
    ret = mosquitto_connect(mosq, HOST, PORT, KEEP_ALIVE);
    if(ret) {
        printf("Connect server error!\n");
        mosquitto_destroy(mosq);
        mosquitto_lib_cleanup();
        return -1;
    }

    printf("Start!\n");
    // mosquitto_loop_start 作用是开启一个线程,在线程里不停地调用 mosquitto_loop() 来处理网络信息
    int loop = mosquitto_loop_start(mosq);
    if(loop != MOSQ_ERR_SUCCESS) {
        printf("mosquitto loop error\n");
        return 1;
    }

    while(fgets(buff, MSG_MAX_SIZE, stdin) != NULL) {
        /* 发布消息 */
        mosquitto_publish(mosq, NULL, "topic1", strlen(buff)+1, buff, 0, 0);
        memset(buff, 0, sizeof(buff));
    }

    mosquitto_destroy(mosq);
    mosquitto_lib_cleanup();
    printf("End!\n");
    return 0;
}

编译配置

使用 Makefile 进行编译管理:

all:
	@echo "Start compiling..."
	gcc -o sub sub.c -lmosquitto
	gcc -o pub pub.c -lmosquitto
	@echo "end"

sub:
	gcc -o sub sub.c -lmosquitto

pub:
	gcc -o pub pub.c -lmosquitto

clean:
	-rm sub pub

异步模式

异步模式利用多线程处理网络 IO,主线程不会被阻塞,更适合高并发或需要与其他业务逻辑交互的场景。

核心区别在于连接和循环函数的调用:

  • 同步:mosquitto_connect + mosquitto_loop (阻塞)
  • 异步:mosquitto_connect_async + mosquitto_loop_start (创建内部线程)

订阅端 (sub.c)

#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <unistd.h>
#include "mosquitto.h"

#define HOST "localhost"
#define PORT 1883
#define KEEP_ALIVE 60
#define MSG_MAX_SIZE 512

static int running = 1;

void my_connect_callback(struct mosquitto *mosq, void *obj, int rc) {
    printf("Call the function: on_connect\n");
    if(rc) {
        printf("on_connect error!\n");
        exit(1);
    } else {
        if(mosquitto_subscribe(mosq, NULL, "topic2", 2)) {
            printf("Set the topic error!\n");
            exit(1);
        }
    }
}

void my_disconnect_callback(struct mosquitto *mosq, void *obj, int rc) {
    printf("Call the function: my_disconnect_callback\n");
    running = 0;
}

void my_subscribe_callback(struct mosquitto *mosq, void *obj, int mid, int qos_count, const int *granted_qos) {
    printf("Call the function: on_subscribe\n");
}

void my_message_callback(struct mosquitto *mosq, void *obj, const struct mosquitto_message *msg) {
    printf("Call the function: on_message\n");
    printf("Receive a message of %s : %s\n", (char *)msg->topic, (char *)msg->payload);
    if(0 == strcmp(msg->payload, "quit")) {
        mosquitto_disconnect(mosq);
    }
}

int main() {
    int ret;
    struct mosquitto *mosq;

    ret = mosquitto_lib_init();
    if(ret) {
        printf("Init lib error!\n");
        return -1;
    }

    mosq = mosquitto_new("sub_test", true, NULL);
    if(mosq == NULL) {
        printf("New sub_test error!\n");
        mosquitto_lib_cleanup();
        return -1;
    }

    mosquitto_connect_callback_set(mosq, my_connect_callback);
    mosquitto_disconnect_callback_set(mosq, my_disconnect_callback);
    mosquitto_subscribe_callback_set(mosq, my_subscribe_callback);
    mosquitto_message_callback_set(mosq, my_message_callback);

    // 异步连接
    ret = mosquitto_connect_async(mosq, HOST, PORT, KEEP_ALIVE);
    if(ret) {
        printf("Connect server error!\n");
        mosquitto_destroy(mosq);
        mosquitto_lib_cleanup();
        return -1;
    }

    ret = mosquitto_loop_start(mosq);
    if(ret) {
        printf("Start loop error!\n");
        mosquitto_destroy(mosq);
        mosquitto_lib_cleanup();
        return -1;
    }

    printf("Start!\n");
    while(running) {
        sleep(1);
    }

    // 停止循环并清理
    mosquitto_loop_stop(mosq, false);
    mosquitto_destroy(mosq);
    mosquitto_lib_cleanup();
    printf("End!\n");
    return 0;
}

发布端 (pub.c)

#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <unistd.h>
#include "mosquitto.h"

#define HOST "localhost"
#define PORT 1883
#define KEEP_ALIVE 60
#define MSG_MAX_SIZE 512

static int running = 1;

void my_connect_callback(struct mosquitto *mosq, void *obj, int rc) {
    printf("Call the function: my_connect_callback\n");
}

void my_disconnect_callback(struct mosquitto *mosq, void *obj, int rc) {
    printf("Call the function: my_disconnect_callback\n");
    running = 0;
}

void my_publish_callback(struct mosquitto *mosq, void *obj, int mid) {
    printf("Call the function: my_publish_callback\n");
}

int main() {
    int ret;
    struct mosquitto *mosq;
    char buff[MSG_MAX_SIZE];

    ret = mosquitto_lib_init();
    if(ret) {
        printf("Init lib error!\n");
        return -1;
    }

    mosq = mosquitto_new("pub_test", true, NULL);
    if(mosq == NULL) {
        printf("New pub_test error!\n");
        mosquitto_lib_cleanup();
        return -1;
    }

    mosquitto_connect_callback_set(mosq, my_connect_callback);
    mosquitto_disconnect_callback_set(mosq, my_disconnect_callback);
    mosquitto_publish_callback_set(mosq, my_publish_callback);

    ret = mosquitto_connect_async(mosq, HOST, PORT, KEEP_ALIVE);
    if(ret) {
        printf("Connect server error!\n");
        mosquitto_destroy(mosq);
        mosquitto_lib_cleanup();
        return -1;
    }

    int loop = mosquitto_loop_start(mosq);
    if(loop != MOSQ_ERR_SUCCESS) {
        printf("mosquitto loop error\n");
        return 1;
    }

    printf("Start!\n");
    while(fgets(buff, MSG_MAX_SIZE, stdin) != NULL) {
        mosquitto_publish(mosq, NULL, "topic2", strlen(buff)+1, buff, 0, 0);
        memset(buff, 0, sizeof(buff));
    }

    mosquitto_loop_stop(mosq, false);
    mosquitto_destroy(mosq);
    mosquitto_lib_cleanup();
    printf("End!\n");
    return 0;
}

Makefile 配置同上。

可订阅可发布模式

在实际场景中,客户端可能既需要发布也需要订阅。下面展示一个双向收发的示例。

客户端代码 (pub_sub.c)

注意此处使用了 <stdbool.h> 来支持布尔类型。

#include <stdio.h>
#include <stdlib.h>
#include <mosquitto.h>
#include <string.h>
#include <stdbool.h>

#define HOST "localhost"
#define PORT 1883
#define KEEP_ALIVE 60
#define MSG_MAX_SIZE 512

bool session = true;

void my_message_callback(struct mosquitto *mosq, void *userdata, const struct mosquitto_message *message) {
    if(message->payloadlen) {
        printf("%s %s", message->topic, (char *)message->payload);
    } else {
        printf("%s (null)\n", message->topic);
    }
    fflush(stdout);
}

void my_connect_callback(struct mosquitto *mosq, void *userdata, int result) {
    int i;
    if(!result) {
        mosquitto_subscribe(mosq, NULL, "topic2 ", 2);
    } else {
        fprintf(stderr, "Connect failed\n");
    }
}

void my_subscribe_callback(struct mosquitto *mosq, void *userdata, int mid, int qos_count, const int *granted_qos) {
    int i;
    printf("Subscribed (mid: %d): %d", mid, granted_qos[0]);
    for(i=1; i<qos_count; i++) {
        printf(", %d", granted_qos[i]);
    }
    printf("\n");
}

void my_log_callback(struct mosquitto *mosq, void *userdata, int level, const char *str) {
    printf("%s\n", str);
}

int main() {
    struct mosquitto *mosq = NULL;
    char buff[MSG_MAX_SIZE];

    mosquitto_lib_init();
    mosq = mosquitto_new(NULL, session, NULL);
    if(!mosq) {
        printf("create client failed..\n");
        mosquitto_lib_cleanup();
        return 1;
    }

    mosquitto_connect_callback_set(mosq, my_connect_callback);
    mosquitto_message_callback_set(mosq, my_message_callback);
    mosquitto_subscribe_callback_set(mosq, my_subscribe_callback);

    if(mosquitto_connect(mosq, HOST, PORT, KEEP_ALIVE)) {
        fprintf(stderr, "Unable to connect.\n");
        return 1;
    }

    int loop = mosquitto_loop_start(mosq);
    if(loop != MOSQ_ERR_SUCCESS) {
        printf("mosquitto loop error\n");
        return 1;
    }

    while(fgets(buff, MSG_MAX_SIZE, stdin) != NULL) {
        mosquitto_publish(mosq, NULL, "topic1 ", strlen(buff)+1, buff, 0, 0);
        memset(buff, 0, sizeof(buff));
    }

    mosquitto_destroy(mosq);
    mosquitto_lib_cleanup();
    return 0;
}

运行此客户端时,它既能接收 topic2 的消息,也能向 topic1 发送消息,实现了双向通信。

目录

  1. MQTT 客户端开发:基于 Mosquitto 的 C 语言实践
  2. 准备工作
  3. 同步模式
  4. 订阅端 (sub.c)
  5. 发布端 (pub.c)
  6. 编译配置
  7. 异步模式
  8. 订阅端 (sub.c)
  9. 发布端 (pub.c)
  10. 可订阅可发布模式
  11. 客户端代码 (pub_sub.c)
  • 免费图片AI生成工具免费生成了解详情
  • Magick API 一键接入全球大模型注册送1000万token查看
  • 免费图片视频在线生成30秒,将你的创意变成现实开始设计
  • X/Twitter免费视频下载器免登陆无限额度免费视频解析下载了解详情
  • 100+免费在线小游戏爽一把
极客日志微信公众号二维码

微信扫一扫,关注极客日志

微信公众号「极客日志V2」,在微信中扫描左侧二维码关注。展示文案:极客日志V2 zeeklog

更多推荐文章

查看全部
  • Redis 配置密码不生效的排查与解决方案
  • 机器人系统设计核心:从架构拆解到工程落地实践
  • MiniMax 海螺 AI 视频:图片与文本生成高质量视频
  • Ubuntu 22.04 基于 ROS2 Humble 的 PX4 无人机仿真环境搭建
  • Trae AI 编程工具使用指南及竞品对比分析
  • VSCode Java 环境配置:解决 JDK 版本不一致问题
  • 算法基础:前缀和技巧与区间求和优化
  • AI 编程工具收费模式变革:Token 计费时代的开发者生存指南
  • VS Code 中 Git 的使用:从零到一入门教程
  • 基于 Trae 构建本地 AI 对话机器人
  • Cesium 无人机智能航线规划:航点动作组与 AI 识别
  • 基于 Vue3 与 Spring Boot 的若依框架快速搭建指南
  • OpenClaw 网关与子节点配对指南:构建分布式 AI 助手网络
  • Java 可执行 JAR 包打包:三种 Maven 方案对比
  • Claude Code 集成至 GitHub Actions 工作流指南
  • WebGIS 视角下的省域区县天气可视化实战攻略
  • 多源 BFS 算法原理及经典题目解析
  • 前端 JS 资源加载失败的重试与多源容灾方案
  • Java 集成 DeepFace 人脸识别微服务架构实践
  • SmolVLA 高算力适配:TensorRT 加速可行性分析与 ONNX 导出实操

相关免费在线工具

  • Base64 字符串编码/解码

    将字符串编码和解码为其 Base64 格式表示形式即可。 在线工具,Base64 字符串编码/解码在线工具,online

  • Base64 文件转换器

    将字符串、文件或图像转换为其 Base64 表示形式。 在线工具,Base64 文件转换器在线工具,online

  • Markdown转HTML

    将 Markdown(GFM)转为 HTML 片段,浏览器内 marked 解析;与 HTML转Markdown 互为补充。 在线工具,Markdown转HTML在线工具,online

  • HTML转Markdown

    将 HTML 片段转为 GitHub Flavored Markdown,支持标题、列表、链接、代码块与表格等;浏览器内处理,可链接预填。 在线工具,HTML转Markdown在线工具,online

  • JSON 压缩

    通过删除不必要的空白来缩小和压缩JSON。 在线工具,JSON 压缩在线工具,online

  • JSON美化和格式化

    将JSON字符串修饰为友好的可读格式。 在线工具,JSON美化和格式化在线工具,online