MQTT 자동 재연결 모범 사례

freeFree Technical Resource

This content is free to read, suitable for basic learning and search traffic.

MQTT 자동 재연결 모범 사례
理解和应用 MQTT 协议是IoT领域中的一项重要技能。为了确保デバイス和サーバー·サーバー之间的稳定连接,我们需要深入了解并有效应用 MQTT クライアント側は的自動再接続。特性。下面,让我们像在探索一个神秘的冒险岛一样,深入探索 MQTT 協定と再接続机制。 MQTT 是一种基于 TCP 协议的发布/订阅模型协议。它像一艘经过风雨洗礼的海船,在IoT、センサーは网络和其他低带宽、不稳定网络环境中航行。但是,就像海上的风浪,网络环境中也充满了各种挑战:网络故障、弱い信号だ化、データ丢包等等,这些都可能使得 MQTT クライアント側は与サーバー·サーバー之间的连接中断。在IoT的大海中,常见的触发断线再接続的风浪包括网络环境恶劣或断网、サーバー·サーバーアップグレード、デバイス或クライアント側は再起動、以及其他网络因素等。 在这种情况下,我们如何确保我们的 "船"(MQTT クライアント側は)始终能与 "港口"(Server)保持稳定的连接呢?答案就是我们需要给我们的 "船" 装备一套自动导航系统——MQTT クライアント側は的自動再接続。逻辑。 设计一个优秀的 MQTT クライアント側は再接続逻辑,就像建造一艘坚固的海船。もし设计得不合理,那么我们的 "船" 可能会失去导航,静默不再受信来自 "港口" 的メッセージ,甚至可能会因为频繁地尝试再接続而无意识地攻击我们的 "港口",这就如同在海上无头乱窜,不仅消耗了自身的能量,也给 "港口" 带来了不必要的压力。然而,もし我们的再接続逻辑设计得合理,那么无论何时失去连接,我们的 "船" 都能稳定地自动导航,重新找到 "港口",并保持与其的连接。 设计 MQTT クライアント側は再接続逻辑时,我们需要考虑几个关键因素:
  1. 航行保活时间:在 MQTT 中,我们称其为 Keep Alive。这是一个タイマータイマー。,它会定期チェック我们的 "船" 是否与 "港口" 保持连接。我们需要根据实际的网络环境和应用需求,来設定一个合适的 Keep Alive。
  2. 再接続策略和退避:当我们的 "船" 失去了与 "港口" 的连接,我们不应立刻尝试重新连接,而应该設定一个合理的待機中时间,以免过度消耗リソース。这就像是当我们的船在海上迷路时,我们需要暂时停下,观察风向、测量海流,然后再制定新的航行路线。我们可以使用指数退避算法或者阶梯式的延时策略来実装这个功能。
  3. 连接ステータス管理:我们的 "船" 需要一个航海日志,来记录与 "港口" 的连接ステータス、连接断开的原因、已经订阅的信息等重要信息。在连接断开时,我们的 "船" 应该查阅航海日志,分析连接断开的原因,然后尝试重新连接 "港口"。
  4. 異常処理:在航行过程中,我们的 "船" 可能会遇到各种各样的問題,例如 "港口" 不可用、认证失敗、网络異常等。我们的 "船" 需要有一个应急计划,来应对这些問題。例如,当 "港口" 不可用时,我们的 "船" 可能需要寻找其他的 "港口";当认证失敗时,我们的 "船" 可能需要チェック自身的认证信息是否正确;当网络異常时,我们的 "船" 可能需要暂停航行,待機中网络恢复正常。
  5. 最大尝试次数限制:对于一些低消費電力はデバイス,我们可能需要考虑限制尝试再接続的次数,以避免过度消耗デバイス的电力。就像在海上迷航的 "船",当它已经尝试了很多次都无法找到 "港口" 时,可能就需要暂时停下,待機中更好的航行条件。
在设计了这个自动导航系统(MQTT クライアント側は的自動再接続。逻辑)之后,我们的 "船" 就能更好地在IoT的海洋中航行,无论面临何种挑战,都能始终保持与 "港口" 的稳定连接,从而确保我们的应用能够顺利进行。 来看一个实际的案例。我们以 Paho MQTT C 库为例,它为我们提供了一套丰富的航海工具——回调函数,让我们可以根据实际情况设定自动导航系统的工作方式。Paho 提供了全局回调、API 回调和异步方法回调,让我们可以在各种情况下都能保持与 "港口" 的连接。 这就是我们如何在IoT的海洋中航行的故事。希望経由这个故事,能够帮助你更好地理解 MQTT 協定と再接続机制,也希望你的 "船" 能在IoT的海洋中顺利航行。  
/*******************************************************************************
 * Copyright (c) 2012, 2022 IBM Corp., Ian Craggs
 *
 * All rights reserved。此程序和随附的资料
 * 根据Eclipse公共ライセンスはv2.0
 * 和Eclipse发行ライセンスはv1.0的条款提供。 
 *
 * Eclipse公共ライセンスは可在以下网址查阅 
 *   https://www.eclipse.org/legal/epl-2.0/
 * Eclipse发行ライセンスは可在以下网址查阅 
 *   http://www.eclipse.org/org/documents/edl-v10.php。
 *
 *******************************************************************************/

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

#if !defined(_WIN32)
#include <unistd.h>
#else
#include <windows.h>
#endif

#if defined(_WRS_KERNEL)
#include <OsWrapper.h>
#endif

// 定义需要使用的MQTT接続パラメータ,如brokerアドレス和クライアント側はID等
#define ADDRESS     "tcp://broker.emqx.io:1883"
#define CLIENTID    "PahoClientSub"
#define TOPIC       "nanomq/test"
#define PAYLOAD     "Hello World!"
#define QOS         1
#define TIMEOUT     10000L

// 定义在主线程中的逻辑Flag
int disc_finished = 0;
int subscribed = 0;
int finished = 0;

//首先声明 API 回调函数
void onConnect(void* context, MQTTAsync_successData* response);
void onConnectFailure(void* context, MQTTAsync_failureData* response);
void onSubscribe(void* context, MQTTAsync_successData* response);
void onSubscribeFailure(void* context, MQTTAsync_failureData* response);

// 下面2个是 Async 使用的回调函数
// 异步接続成功です。的回调函数,在接続成功です。的时候进行SubscribeOperation。
void conn_established(void *context, char *cause)
{
	printf("クライアント側は已重新连接!\n");
	MQTTAsync client = (MQTTAsync)context;
	MQTTAsync_responseOptions opts = MQTTAsync_responseOptions_initializer;
	int rc;

	printf("接続成功です。\n");

	printf("订阅主题 %s\n使用クライアント側は %s 并用QoS%d\n\n"
           "按Q<Enter>退出\n\n", TOPIC, CLIENTID, QOS);
	opts.onSuccess = onSubscribe;
	opts.onFailure = onSubscribeFailure;
	opts.context = client;
	if ((rc = MQTTAsync_subscribe(client, TOPIC, QOS, &opts)) != MQTTASYNC_SUCCESS)
	{
		printf("开始订阅失敗,戻る码 %d\n", rc);
		finished = 1;
	}
}

// 异步连受信到 Disconnectメッセージ时的回调,由于主に。断开的情况下不会收到 Disconnectメッセージ,所以此方法很少被触发
void disconnect_lost(void* context, MQTTProperties* properties,
		enum MQTTReasonCodes reasonCode)
{
	printf("クライアント側は已切断する!\n");
}

// 下面是クライアント側は全局回调函数,分别是连接断开和メッセージ到达
void conn_lost(void *context, char *cause)
{
	MQTTAsync client = (MQTTAsync)context;
	MQTTAsync_connectOptions conn_opts = MQTTAsync_connectOptions_initializer;
	int rc;

	printf("\n连接已断开\n");
	if (cause)
		printf("     原因: %s\n", cause);

	printf("正在再接続\n");
	conn_opts.keepAliveInterval = 20;
	conn_opts.cleansession = 1;
	conn_opts.maxRetryInterval = 16;
	conn_opts.minRetryInterval = 2;
	conn_opts.automaticReconnect = 1;
	
	//conn_opts.onSuccess = onConnect;
	conn_opts.onFailure = onConnectFailure;
	MQTTAsync_setConnected(client, client, conn_established);
	if ((rc = MQTTAsync_connect(client, &conn_opts)) != MQTTASYNC_SUCCESS)
	{
		printf("开始接続の失敗,戻る码 %d\n", rc);
		finished = 1;
	}
}

// 收到メッセージ时的全局回调函数,此处简单的打印メッセージ
int msgarrvd(void *context, char *topicName, int topicLen, MQTTAsync_message *message)
{
    printf("メッセージ已到达\n");
    printf("     主题: %s\n", topicName);
    printf("   消```C
息: ");

    /* 打印メッセージ内容 */
    char* payloadptr = message->payload;
    for(int i = 0; i < message->payloadlen; i++)
    {
        putchar(*payloadptr++);
    }
    putchar('\n');

    /* 释放メッセージ内存 */
    MQTTAsync_freeMessage(&message);
    MQTTAsync_free(topicName);

    return 1;
}

/* 异步切断する的回调函数 */
void onDisconnect(void* context, MQTTAsync_successData* response)
{
    printf("成功切断する\n");
    disc_finished = 1;
}

/* 异步接続成功です。的回调函数,在接続成功です。的时候进行订阅操作。 */
void onConnect(void* context, MQTTAsync_successData* response)
{
    MQTTAsync client = (MQTTAsync)context;
    MQTTAsync_responseOptions opts = MQTTAsync_responseOptions_initializer;
    int rc;

    printf("成功连接\n");

    printf("订阅主题 %s\n使用クライアント側は %s 并用QoS%d\n\n"
           "按Q<Enter>退出\n\n", TOPIC, CLIENTID, QOS);

    /* 开始订阅 */
    opts.onSuccess = onSubscribe;
    opts.onFailure = onSubscribeFailure;
    opts.context = client;
    if ((rc = MQTTAsync_subscribe(client, TOPIC, QOS, &opts)) != MQTTASYNC_SUCCESS)
    {
        printf("开始订阅失敗,戻る码 %d\n", rc);
        finished = 1;
    }
}

/* 异步接続の失敗的回调函数 */
void onConnectFailure(void* context, MQTTAsync_failureData* response)
{
    printf("接続の失敗\n");
    if (response && response->message)
    {
        printf("失敗信息: %s\n", response->message);
    }
    finished = 1;
}

/* 异步订阅成功的回调函数 */
void onSubscribe(void* context, MQTTAsync_successData* response)
{
    printf("成功订阅\n");
    subscribed = 1;
}

/* 异步订阅失敗的回调函数 */
void onSubscribeFailure(void* context, MQTTAsync_failureData* response)
{
    printf("订阅失敗\n");
    if (response && response->message)
    {
        printf("失敗信息: %s\n", response->message);
    }
    finished = 1;
}

/* 异步キャンセル订阅的回调函数 */
void onUnsubscribe(void* context, MQTTAsync_successData* response)
{
    printf("成功キャンセル订阅\n");
    finished = 1;
}

int main(int argc, char* argv[])
{
    MQTTAsync client;
    MQTTAsync_connectOptions conn_opts = MQTTAsync_connectOptions_initializer;
    int rc;
    MQTTAsync_message pubmsg = MQTTAsync_message_initializer;
    MQTTAsync_token token;

    /* 作成MQTTクライアント側は */
    MQTTAsync_create(&client, ADDRESS, CLIENTID, MQTTCLIENT_PERSISTENCE_NONE, NULL);

    /* 設定全局回调函数 */
    MQTTAsync_setCallbacks(client, client, conn_lost, msgarrvd, NULL);

    /* 設定连接选项 */
    conn_opts.keepAliveInterval = 20;
    conn_opts.cleansession = 1;
    conn_opts.automaticReconnect = 1;
    //conn_opts.onSuccess = onConnect;
    conn_opts.onFailure = onConnectFailure;
    MQTTAsync_setConnected(client, client, conn_established);

    /* 开始连接 */
    if ((rc = MQTTAsync_connect(client, M Q TT AS Y NC _ SU CC ESS)
    {
        print f (" 연 결 시작 실패 , 반환 코드 % d \ n ", r c);
        EX IT _ FA IL UR E 를 반환 합니다 .
    }

    While (아 직 도 !완료 됨)
    {
        # if defined (_ W IN 32)
            잠 들 기 (1 000)
        # else
            잠 들 다 (1)
        # 엔 디 프
    }

    if (sub s cribe)
    {
        if ((r c = M Q T TA syn c _ uns ub s cribe (cli ent , TOP IC , N ULL)) ! = M Q TT AS Y NC _ SU CC ESS)
        {
            print f (" 구 독 취소 실패 , 반환 코드 % d \ n ", r c);
            EX IT _ FA IL UR E 를 반환 합니다 .
        }
    }
    
    / * 연결 끊 기 */
    M Q T TA syn c _ dis conne ct Op tions disc _ op ts = M Q T TA syn c _ dis conne ct Op tions _ in itia li zer ;
    disc _ op ts . on S uc cess = on Dis conne ct ;
    if ((rc = MQTTAsync_disconnect(client, M Q TT AS Y NC _ SU CC ESS)
    {
        print f (" % d \ n " 반환 코드 에서 연결 해 제 시작 실패);
        EX IT _ FA IL UR E 를 반환 합니다 .
    }
    While (아 직 도 ! disc _ finished)
    {
        # if defined (_ W IN 32)
            잠 들 기 (1 000)
        # else
            잠 들 다 (1)
        # 엔 디 프
    }

    / * 클 라이언 트 제거 */
    MQTTAsync_destroy(

    EX IT _ SU CC ESS 를 반환 합니다 .
}
 
Related Tags
Put this resource to use in a real project?

Go to the Tool Center for message parsing, CRC verification and device debugging, or submit your requirements for selection and integration advice.

Engineer Membership

Turn this article into actionable debugging resources

After activation, you can use advanced message parsing, resource pack downloads, code examples, engineering cases and priority technical support, suitable for real project delivery.

Unlimited Advanced Tools
Resource & Code Packs
Complete Engineering Case Library
Priority Technical Support

Leave a Reply

Your email address will not be published. Required fields are marked *.