点赞
评论
收藏
分享
举报
IoT 场景-02:通过Nginx JavaScript 实现会话保持
发表于2020-12-19 15:51

浏览 6.2k

文章标签

上一节我们介绍过了通过Nginx 实现MQTT 会话的负载均衡,那么如何解决每个MQTT Broker在消息分发的时候,只能收到一半信息的情况呢?

此时,我们就需要借助Nginx Plus的高级功能:Nginx JavaScript,来定制基于ClientID 的会话保持机制,将客户端与某个MQTT Broker进行会话保持,保障其信息的完整性,除非该MQTT Broker出现故障。


"Round Robin load balancing is an effective mechanism for distributing client connections across a group of servers. However, there are several reasons why it is not ideal for MQTT connections."


0x01 Nginx JavaScript 会话保持代码

mqtt.js 提取 MQTT ClientId 

var client_messages = 1;
2
var client_id_str = "-";
3
4
function getClientId(s) {
5
    if ( !s.fromUpstream ) {
6
        if ( s.buffer.toString().length == 0  ) { // Initial calls may
7
            s.log("No buffer yet");               // contain no data, so
8
            return s.AGAIN;                       // ask that we get called again
9
        } else if ( client_messages == 1 ) { // CONNECT is first packet from the client
            // CONNECT packet is 1, using upper 4 bits (00010000 to 00011111)
2
            var packet_type_flags_byte = s.buffer.charCodeAt(0);
3
            s.log("MQTT packet type+flags = " + packet_type_flags_byte.toString());
4
            if ( packet_type_flags_byte >= 16 && packet_type_flags_byte < 32 ) {
5
                // Calculate remaining length with variable encoding scheme
6
                var multiplier = 1;
7
                var remaining_len_val = 0;
8
                var remaining_len_byte;
9
                for (var remaining_len_pos = 1; remaining_len_pos < 5; remaining_len_pos++ ) {
10
                    remaining_len_byte = s.buffer.charCodeAt(remaining_len_pos);
11
                    if ( remaining_len_byte == 0 ) break; // Stop decoding on 0
12
                    remaining_len_val += (remaining_len_byte & 127) * multiplier;
13
                    multiplier *= 128;
14
                }
15
16
                // Extract ClientId based on length defined by 2-byte encoding
17
                var payload_offset = remaining_len_pos + 12; // Skip fixed head
1
                var client_id_len_msb = s.buffer.charCodeAt(payload_offset).toString(16);
2
                var client_id_len_lsb = s.buffer.charCodeAt(payload_offset + 1).toString(16);
3
                if ( client_id_len_lsb.length < 2 ) client_id_len_lsb = "0" + client_id_len_lsb;
4
                var client_id_len_int = parseInt(client_id_len_msb + client_id_len_lsb, 16);
5
                client_id_str = s.buffer.substr(payload_offset + 2, client_id_len_int);
6
                s.log("ClientId value  = " + client_id_str);
7
            } else {
8
                s.log("Received unexpected MQTT packet type+flags: " + packet_type_flags_byte.toString());
9
            }
10
        }
11
        client_messages++;
12
    }
13
    return s.OK;
14
}
15
16
function setClientId(s) {
17
    return client_id_str;
18
}


在配置文件nginx.conf 中进行引用 

upstream hive_mq {

server 172.17.0.3:1883; #mq1

server 172.17.0.2:1883; #mq2

zone tcp_mem 64k;

hash $mqtt_client_id consistent; # Session persistence keyed against ClientId

}

 

server {

listen 1883;

preread_buffer_size 1k; # Big enough to read CONNECT packet header

js_prereadgetClientId; # Parse CONNECT packet for ClientId

proxy_passhive_mq;

proxy_connect_timeout 1s; 


0x02 测试会话保持功能


此时,使用发布端持续发布数据,订阅客户端收到的信息会保持连续。

当停止其中任何一个MQTT Broker容器时,订阅客户端会有一定几率重新连接MQTT Broker并重新建立会话保持。



此时读者或许会问道,如果MQTT Broker容器意外故障,是否可以自动将故障点进行移除,并主动进行客户端的连接的重连? 好问题,请期待我们的下一节,使用NGINX Plus 进行主动MQTT健康检查。



参考:

https://dzone.com/articles/mqtt-load-balancing-and-session-persistence-with-nginx-plus













已修改于2023-03-09 02:07
本作品系原创
创作不易,留下一份鼓励
yuefeng

暂无个人介绍

关注



写下您的评论
发表评论
全部评论(1)

按点赞数排序

按时间排序

如果 sub 用 client id 1,pub 时也用 client id 1,就会连接冲突,但上例中怎么用 client id 做负载均衡的?

赞同

0

回复举报

发表于2024-03-15 18:15



回复6218018upea
回复
关于作者
yuefeng
这家伙很懒还未留下介绍~
6
文章
0
问答
2
粉丝
相关文章
众所周知,Nginx最常见的传统场景是Web服务器,HTTP反向代理以及负载均衡,此外,它在物联网的技术领域,也可以发挥同样可观的作用。本节我们会重点讨论使用Nginx实现物联网消息组件MQTTBroker的高可用。MQTT流量负载均衡带有健康检查的高可用实践基于MQTTClientID的会话保持0x01MQTT原理首先我们来熟悉一下MQTT协议。MQTT(MessageQueuingTelemetry Transport,消息队列遥测传输协议),是一种基于发布/订阅(publish/subscribe)模式的"轻量级"通讯协议,该协议构建于TCP/IP协议上,由IBM在1999年发布。MQTT最大优点在于,可以以极少的代码和有限的带宽,为连接远程设备提供实时可靠的消息服务。作为一种低开销、低带宽占用的即时通讯协议,使其在物联网、小型设备、移动应用等方面有较广泛的应用。MQTT协议原理1MQTT协议实现方式实现MQTT协议需要客户端和服务器端通讯完成,在通讯过程中,MQTT协议中有三种身份:发布者(Publish)、代理(Broker)(服务器)、订阅者(S
点赞 1
浏览 8.6k
为了使我们的NginxPlus具有对MQTTBroker的正常工作状况具有主动健康检查的能力,首先我们要了解MQTT的协议的连接。0x01MQTT协议第一节我们介绍过,MQTT是一种基于TCP建立的的一套协议,MQTT报文中也有类似于TCP三次握手的状态码。简单来讲,我们只需要构建一个合法的CONNECT连接,等待MQTT返回一个CONNACK的回应报文,就能确认MQTTBroker正在处于一个正常工作的状态。0x02NginxPlus构建CONNECT1matchmqtt_conn{12#SendCONNECTpacketwithclientID"nginxhealthcheck"13send\x10\x20\x00\x06\x4d\x51\x49\x73\x64\x70\x03\x02\x00\x3c\x00\x12\x6e\x67\x69\x6e\x78\x20\x68\x65\x61\x6c\x74\x68\x20\x63\x68\x65\x63\x6b;14expect\x2
点赞 1
浏览 5.1k
由于NGINXPlus在MQTT客户端与Broker的交互过程中处于一个核心代理的位置,我们可以很容易的在其原生的代理功能基础之上加入安全访问控制功能。0x01为什么需要访问控制对于一些MQTT的客户端,不管是消息的发布者还是订阅者,都有可能会使用贪婪模式,频繁的与MQTTBroker进行通信,并发布或获取大量数据。此时,MQTTBroker的通信信道会被这种客户端应用占用,导致服务降级。遇到这种类型的客户端,由于MQTTBroker原生并没有很好地限制机制,从而需要借助在NginxPlus的代理上面,进行配置,识别并限制过度的数据访问。0x02如何进行访问控制使用NginxPlus进行访问控制很简单,和HTTP协议的访问控制类似。0x03访问日志处理使用NginxPlus还可以自定义日志格式,通过与告警平台和自动化平台对接,可以有效识别异常客户端的通信,可以进行自动化的IP封禁以及限连和限流操作。至此,我们4节NGINXPlus在IoT场景中的应用已经告一段段落,希望大家喜欢本次的分享。参考:https://dzone.com/articles/mq
点赞 1
浏览 6.6k