基于mosquitto.dll的MQTT客户端无法接收消息求助
Mosquitto.dll MQTT客户端无法触发onConnect/onMessage回调问题
- 客户端可正常发布消息,
onDisconnect、onLog回调正常触发,但onConnect和onMessage无响应 - 服务器日志显示客户端已成功连接、完成订阅,且消息已转发至该客户端
- 修改订阅主题、QoS参数后问题仍未解决
服务器日志
1696166724: New connection from 127.0.0.1:59203 on port 1883. 1696166724: New client connected from 127.0.0.1:59203 as mqttx_22be153f (p5, c1, k60, u'user'). 1696166724: No will message specified. 1696166724: Sending CONNACK to mqttx_22be153f (0, 0) 1696166724: Received SUBSCRIBE from mqttx_22be153f 1696166724: testtopic/# (QoS 0) 1696166724: mqttx_22be153f 0 testtopic/# 1696166724: Sending SUBACK to mqttx_22be153f 1696166727: New connection from ::1:59204 on port 1883. 1696166727: New client connected from ::1:59204 as mqttx_d84dc2b6 (p2, c1, k100, u'user'). 1696166727: No will message specified. 1696166727: Sending CONNACK to mqttx_d84dc2b6 (0, 0) 1696166727: Received SUBSCRIBE from mqttx_d84dc2b6 1696166727: # (QoS 0) 1696166727: mqttx_d84dc2b6 0 # 1696166727: Sending SUBACK to mqttx_d84dc2b6 1696166730: Received PUBLISH from mqttx_d84dc2b6 (d0, q0, r1, m0, 'testtopic', ... (6 bytes)) 1696166730: Sending PUBLISH to mqttx_22be153f (d0, q0, r0, m0, 'testtopic', ... (6 bytes)) 1696166730: Sending PUBLISH to mqttx_d84dc2b6 (d0, q0, r0, m0, 'testtopic', ... (6 bytes)) 1696166733: Received PUBLISH from mqttx_22be153f (d0, q1, r0, m24914, 'testtopic', ... (21 bytes)) 1696166733: Sending PUBLISH to mqttx_22be153f (d0, q0, r0, m0, 'testtopic', ... (21 bytes)) 1696166733: Sending PUBLISH to mqttx_d84dc2b6 (d0, q0, r0, m0, 'testtopic', ... (21 bytes)) 1696166733: Sending PUBACK to mqttx_22be153f (m24914, rc0) 1696166737: Client mqttx_d84dc2b6 closed its connection. 1696166740: Received DISCONNECT from mqttx_22be153f 1696166740: Client mqttx_22be153f disconnected. 1696166747: mosquitto version 2.0.14 terminating
可正常工作的客户端代码
//--------------------------------------------------------------------------- #include <vcl.h> #pragma hdrstop extern "C" { #include <mosquitto.h> } #include "Unit1.h" //--------------------------------------------------------------------------- #pragma package(smart_init) #pragma resource "*.dfm" TForm1 *Form1; //--------------------------------------------------------------------------- void TForm1:: onConnect(struct mosquitto* mosq, void* userdata, int result) { if (result == 0) { CheckBox1->Checked=1; } else { CheckBox1->Checked=0; } Edit1->Text=IntToStr(result) + " int "; } void TForm1::onMessage(struct mosquitto* mosq, void* userdata, const struct mosquitto_message* message) { Memo1->Lines->Add("massage"); } __fastcall TForm1::TForm1(TComponent* Owner) : TForm(Owner) { } //--------------------------------------------------------------------------- void __fastcall TForm1::Button1Click(TObject *Sender) { mosquitto_username_pw_set(mosq,"user", "*****"); int rc = mosquitto_connect(mosq, "localhost", 1883, 100); if (rc != MOSQ_ERR_SUCCESS) { Memo1->Lines->Add("mqtt not work" +IntToStr(rc) +" ") ; } rc= mosquitto_validate_utf8("#",1); if (rc != MOSQ_ERR_SUCCESS) { Memo1->Lines->Add("mqtt val not ok ") ; } else {Memo1->Lines->Add("mqtt val ok ") ;} mosquitto_subscribe(mosq, NULL, "#", 0); if (rc != MOSQ_ERR_SUCCESS) { Memo1->Lines->Add("mqtt sub not work ") ; } mosquitto_loop_start(mosq); } //--------------------------------------------------------------------------- void onConnect_(struct mosquitto* mosq, void* userdata, int result) { Form1->onConnect(mosq,userdata, result); } void onMessage_(struct mosquitto* mosq, void* userdata, const struct mosquitto_message* message) { Form1->onMessage(mosq,userdata,message); } void onDisconnect_(struct mosquitto* mosq, void* userdata, int result) { Form1->onConnect(mosq,userdata, result); } void onLog_(struct mosquitto* mosq, void* level, int cnt, const char *data) { String ans = data; Form1->Memo1->Lines->Add(ans); } void __fastcall TForm1::FormCreate(TObject *Sender) { mosquitto_lib_init(); mosq = mosquitto_new("mqttx_d84dc2b6", true, NULL); mosquitto_connect_callback_set(mosq, onConnect_ ); mosquitto_disconnect_callback_set(mosq,onDisconnect_); mosquitto_message_callback_set(mosq, onMessage_ ); mosquitto_log_callback_set(mosq,onLog_); } //--------------------------------------------------------------------------- void __fastcall TForm1::FormClose(TObject *Sender, TCloseAction &Action) { mosquitto_disconnect(mosq); // Отключение от брокера mosquitto_destroy(mosq); mosquitto_lib_cleanup(); } //--------------------------------------------------------------------------- void __fastcall TForm1::Button2Click(TObject *Sender) { mosquitto_disconnect(mosq); } //--------------------------------------------------------------------------- void __fastcall TForm1::Button3Click(TObject *Sender) { static int m; mosquitto_publish(mosq ,NULL, "testtopic",6,"test1 123; ", 0, true ); } //---------------------------------------------------------------------------
客户端日志
CONNECT Client mqttx_d84dc2b6 sending SUBSCRIBE (Mid: 1, Topic: #, QoS: 0, Options: 0x00) Client mqttx_d84dc2b6 sending PUBLISH (d0, q0, r1, m3, 'testtopic', ... (6 bytes))
问题排查与修复
核心问题分析
- UI线程与MQTT线程冲突:VCL组件仅允许主线程操作,而mosquitto回调运行在
mosquitto_loop_start创建的后台线程中,直接操作组件会导致线程安全问题,引发回调无响应。 - 订阅结果判断逻辑错误:代码中用被
mosquitto_validate_utf8覆盖的rc变量判断订阅是否成功,无法真实反映订阅操作结果。 - 日志回调参数不匹配:
onLog_函数参数签名错误,mosquitto日志回调的第二个参数应为int level,而非void* level,参数不匹配会导致回调执行异常。
修复方案
1. 修正日志回调参数
void onLog_(struct mosquitto* mosq, void* userdata, int level, const char *data) { String ans = data; // 切换到主线程更新UI TThread::Queue(NULL, [ans](){ Form1->Memo1->Lines->Add(ans); }); }
2. 线程安全更新UI
所有回调中操作VCL组件的代码,需通过TThread::Queue切换到主线程执行:
void TForm1:: onConnect(struct mosquitto* mosq, void* userdata, int result) { TThread::Queue(NULL, [result](){ CheckBox1->Checked = (result == 0); Edit1->Text = IntToStr(result) + " int "; }); } void TForm1::onMessage(struct mosquitto* mosq, void* userdata, const struct mosquitto_message* message) { TThread::Queue(NULL, [](){ Memo1->Lines->Add("message received"); // 可选:解析并显示消息内容 // String msgContent = String((char*)message->payload, message->payloadlen); // Memo1->Lines->Add("Content: " + msgContent); }); }
3. 修正订阅结果判断
单独保存订阅操作的返回值,避免变量被覆盖:
void __fastcall TForm1::Button1Click(TObject *Sender) { mosquitto_username_pw_set(mosq,"user", "*****"); int rc = mosquitto_connect(mosq, "localhost", 1883, 100); if (rc != MOSQ_ERR_SUCCESS) { Memo1->Lines->Add("mqtt connect failed: " + IntToStr(rc)); } rc = mosquitto_validate_utf8("#",1); if (rc != MOSQ_ERR_SUCCESS) { Memo1->Lines->Add("topic validation failed"); } else { Memo1->Lines->Add("topic validation ok"); } // 单独存储订阅返回值 int sub_rc = mosquitto_subscribe(mosq, NULL, "#", 0); if (sub_rc != MOSQ_ERR_SUCCESS) { Memo1->Lines->Add("mqtt subscribe failed: " + IntToStr(sub_rc)); } mosquitto_loop_start(mosq); }
4. 确保客户端ID唯一性
每次创建客户端时生成唯一ID,避免因ID重复导致连接异常:
void __fastcall TForm1::FormCreate(TObject *Sender) { mosquitto_lib_init(); // 基于时间戳生成唯一客户端ID String clientId = "mqttx_" + IntToStr(GetTickCount()); mosq = mosquitto_new(clientId.c_str(), true, NULL); mosquitto_connect_callback_set(mosq, onConnect_ ); mosquitto_disconnect_callback_set(mosq,onDisconnect_); mosquitto_message_callback_set(mosq, onMessage_ ); mosquitto_log_callback_set(mosq,onLog_); }
内容的提问来源于stack exchange,提问作者андрей раский
相关产品推荐
相关产品推荐

