ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

EMQX+Python:从部署到接入的MQTT消息通信实战指南

EMQX+Python:从部署到接入的MQTT消息通信实战指南 这段时间被一批环境监测设备的消息接入折腾得够呛设备数量不算多但分布散、网络不稳定数据要实时上报后端还要能远程下发控制指令。刚开始图省事直接用HTTP轮询结果设备一多连接经常超时数据还会乱序。后来换成MQTT消息方案服务端用EMQX做Broker问题才算真正解决。EMQX是目前用得比较多的开源MQTT消息服务器对海量设备连接支持得非常好自带管理控制台规则引擎和数据集成也方便。这篇文章从EMQX是什么开始接着讲安装部署、基础功能最后用Python写代码做连接测试把整个流程完整跑通。想在自己的项目里引入MQTT中间件的朋友可以照着走一遍。1. 为什么消息通信选EMQX先把它看明白1.1 MQTT Broker到底在解决什么问题MQTT是一个发布/订阅模型的消息协议很多新手第一次接触时会把它和普通的消息队列搞混。用最直白的话说Broker就是一个消息中转站设备端往某个“主题”里扔消息订阅了这个主题的客户端就会收到。设备与设备之间不需要知道对方的IP只需要知道Broker的地址和主题名称就行。这个机制解决了我项目里最头疼的问题设备端网络经常断开、IP不固定。如果用HTTP直连后端设备端必须知道后端服务的地址而且每个设备都要管理一套连接状态。MQTT把连接状态收敛到了Broker这一层设备只要连上Broker发消息和收消息都通过主题来路由。相当于大家不在一个办公室但都在同一个微信群里不管谁什么时候上线群消息都还在。MQTT本身是为低带宽、高延迟、不稳定的网络设计的所以它的报文头很小控制报文非常精简。典型场景包括传感器数据采集、远程设备控制、移动端消息推送。对物联网项目来说这是很合适的通信方式。1.2 EMQX在同类Broker中强在哪市面上MQTT Broker不少轻量级的有Mosquitto中大型的还有VerneMQ、HiveMQ等。我最初图省事在本地用Mosquitto做了个测试单机单节点的情况下确实够用。但当我开始考虑后续要接几十万设备、需要把消息流转到数据库和大数据平台时Mosquitto的生态和集成能力就显得吃力了。EMQX是开源项目基于Erlang/OTP开发这个技术栈的强项就是高并发、低延迟和分布式容错。它能支撑百万级设备连接在同类开源Broker里属于第一梯队。EMQX的另一个优势是内置规则引擎和“数据集成”能力消息进来以后可以直接通过规则转发到MySQL、PostgreSQL、Kafka、Redis、HTTP服务等省去了自研消费者服务的成本。我做过一个简单对比适合大家参考对比项EMQXMosquittoVerneMQ部署难度中等有Docker镜像低轻量中等大规模连接能力强百万级弱适合几千级较强控制台管理界面内置完整无独立图形界面有基础界面规则引擎与数据集成内置功能丰富需要外部扩展需要外部扩展集群能力强自动集群弱支持适用阶段开发测试到生产均可小型项目/教学中大型但生态相对小众如果你的项目就是几个设备做演示Mosquitto当然够用。但考虑到后续扩展性我建议直接从EMQX入手省得以后迁移。1.3 什么场景适合引入EMQX我梳理了一下下面这几种场景用EMQX比较合适物联网设备接入传感器、智能硬件、工业网关设备量大且网络不稳定的场景。车联网车载终端实时上传位置、状态平台下行控制指令。移动消息推送App服务端向客户端推送通知。实时数据管道设备数据先进入MQTT主题再通过规则引擎分发到Kafka或数据库。带状态反馈的远程控制不只上报数据还要下发指令并确认设备执行结果。但也得说句实在话如果只是单机版测试、仅需在几个客户端之间传个消息没必要上EMQX。我之前在个人电脑上装EMQX只用了不到十分钟但它提供的不仅仅是Broker功能而是一个完整的消息处理底座。这个底座对后续项目的价值远大于消息转发本身。2. 版本、端口与部署方式动手前先盘三件事2.1 版本选择开源版还是企业版在EMQX官网能看到两个大版本一个是开源版一个是企业版。开源版的核心能力已经非常强MQTT Broker、Dashboard、认证鉴权、规则引擎、数据集成、Cluster集群都包含在内。企业版主要增加了更多数据集成连接器、与特定云服务的深度整合、多活容灾方案等。还有一点需要留意EMQX 4.x和5.x的配置方式、Dashboard界面差异不小。4.x用的是老版emqx.conf风格5.x改成了HOCON格式更结构化了。如果你在网上搜到的教程大部分是4.x照搬到5.x可能会碰到参数名对不上的问题。我这次用的是5.x所以下面的操作都以5.x为准。对大多数人来说直接上新版就好别在老版本上摇摆。如果是学习和项目预研开源版完全够用。我第一次部署直接拉了最新开源稳定版后续所有测试包括Python接入运行都很顺畅。2.2 端口规划清单EMQX装起来不难但端口规划必须在动手前想清楚。我列一张表说明默认端口的作用端口协议/用途说明1883MQTT over TCP最常见的MQTT连接端口8883MQTT over SSL/TLS加密连接生产环境建议开启8083MQTT over WebSocket浏览器或Web端设备使用8084MQTT over WSSWebSocket的TLS加密版本18083Dashboard管理界面浏览器访问管理控制台5369RPC端口EMQX节点间通信使用集群部署时需放行如果你的EMQX装在云服务器上一定要在安全组或防火墙里放行需要用到的端口。我第一次在云服务器上部署1883端口没放行本地客户端怎么都连不上查了很久才发现是安全组拦截。这个坑很基础但特别容易踩。2.3 部署方式取舍常用的部署方式有这么几种我根据自己的实际经验排个优先级Docker部署环境隔离启动快适合开发和测试甚至单机生产也可以用。Linux二进制包部署适合生产环境通过systemd管理服务稳定可控。Windows zip包部署适合在Windows机器上做本地实验不推荐生产使用。Kubernetes部署适合大规模集群通过EMQX Operator管理学习成本高一些。我平时习惯先用Docker跑通功能然后再考虑具体部署形态。Docker方式之所以推荐是因为不用关心依赖环境一条命令就能拉起完整服务后续想删掉也干净。3. Docker快速部署与Windows包安装的完整过程3.1 Docker方式一条命令拉起服务确保本机已经装好Docker在命令行里执行docker run -d --name emqx \ -p 1883:1883 -p 8083:8083 \ -p 8084:8084 -p 8883:8883 \ -p 18083:18083 \ emqx/emqx:5.8.4如果本地还没有这个镜像Docker会自动下载。启动成功后用浏览器打开http://localhost:18083就能看到EMQX管理控制台登录页面默认用户名是admin密码是public首次登录后建议马上修改。这个命令里有几个细节值得说明。-d表示后台运行--name emqx给容器起名后面端口映射把宿主机的1883、8083、8084、8883、18083分别映射到容器内。容器内的EMQX默认只监听这些端口如果你后面要用其他端口需要在这里再加映射。如果想持久化配置和数据建议加上数据卷挂载docker run -d --name emqx \ -p 1883:1883 -p 8083:8083 \ -p 8084:8084 -p 8883:8883 \ -p 18083:18083 \ -v emqx_data:/opt/emqx/data \ -v emqx_etc:/opt/emqx/etc \ -v emqx_log:/opt/emqx/log \ emqx/emqx:5.8.4这样容器被删除后配置和数据还在宿主机上。后面升级镜像或迁移容器不会被清空。3.2 Windows二进制包方式下载、配置、启动服务有些朋友的测试环境是Windows不想装Docker。EMQX官方也提供了Windows安装包从官网下载zip压缩包后解压到一个路径中不含中文和空格的目录比如D:\emqx。进入解压后的目录打开PowerShell或cmdcd bin .\emqx start如果一切正常命令行会提示服务启动成功。可以通过下面的命令确认服务状态.\emqx ctl status也可以直接在浏览器访问http://localhost:18083。EMQX 5.x在Windows上以zip包方式运行不需要额外安装Erlang运行时因为发行包已经内置了。Windows方式有一个地方和Docker不一样它默认读取的是解压目录下的etc/emqx.conf配置文件。如果你想修改监听端口或其他参数编辑这个文件后用.\emqx stop停止服务再.\emqx start启动服务。需要注意Windows下有些安全软件会拦截端口监听如果你发现启动成功但外网访问不了先检查杀毒软件和防火墙。3.3 验证部署是否成功部署完了不能只看Dashboard能打开就说成功我习惯做三层验证第一检查端口是否在监听。Windows命令行下执行netstat -ano | findstr :1883Linux下用netstat -tlnp | grep 1883能看到LISTENING状态说明端口服务起来了。第二看EMQX日志。Docker方式用docker logs emqxWindows方式看log目录下的emqx.log。日志里出现EMQX 5.8.4 is started successfully之类的字样才算启动完整。第三用MQTT客户端连一次。可以先不写Python直接用一个叫MQTTX的图形化客户端工具快速验证1883端口通不通。如果MQTTX能连上说明Broker核心功能已经就绪可以进入下一步了。4. Dashboard上手先把这些基础功能过一遍4.1 登录与全局概览打开Dashboard后最直观的就是“概览”页面。这里能看到当前集群节点状态、在线连接数、主题数、订阅数以及实时消息收发速率。对运维而言这个页面是判断Broker健康状态的第一入口。我第一次搭好之后主要是通过这里观察消息流量变化确认是否有异常突发。比如设备上报频率突然升高消息速率曲线会有明显波动这时候就要考虑是不是设备端逻辑出了问题。默认账号admin/public是很多人容易忽略的安全隐患。登录后第一时间去“系统设置”里改密码或者单独创建一个管理员账号把默认账号禁用。这个动作花不了两分钟但能避免匿名访问管理后台的低级事故。4.2 客户端管理在Dashboard左侧菜单找到“客户端”能看到当前所有连接上来的客户端列表。每个客户端有Client ID、IP地址、协议版本、连接状态等信息。这个页面在调试阶段非常实用。我写Python代码测试时经常开两个订阅端然后在Dashboard里确认它们都成功连上来了。如果客户端连接后马上断开也能在这里看到“已断开”状态方便定位问题。Dashboard支持直接踢掉指定客户端。如果发现某个异常连接一直在占用资源点一下就可以强制断开。生产环境慎用这个操作最好先确认业务影响。4.3 主题与订阅管理“主题”页面列出了Broker当前活跃的主题能看到每个主题的消息数量。旁边的“订阅”页面可以查看谁在订阅哪个主题用的QoS等级是多少。这个功能排场特别大。有一次我调试时发现订阅端收不到消息排查了半天最后打开订阅页面一看原来订阅端订阅的是sensor/#而发布端发到了sensor_data/#主题前缀根本不一致。通过订阅页面一眼就能看出主题不匹配的问题。4.4 认证鉴权配置EMQX默认是允许匿名连接的也就是说任何客户端只要知道Broker地址和端口不需要账号密码就能连上来收发消息。这个配置在测试阶段很方便但一旦放到生产环境等于把消息通道裸奔在外面。建议去“访问控制”里的“认证”功能添加一个“内置数据库”认证器。步骤大致是点击“认证”页面里的“创建”选择“内置数据库”填写认证器名称然后添加用户输入用户名、密码。创建完成后再在“认证”列表里把新认证器启用。启用认证之后Python连接时就需要传入用户名和密码了后面测试代码里我会加上这一块。这么做的好处是防止陌生客户端随意连接尤其是设备暴露在公网的情况下这是必须做的防护。4.5 保留消息、遗嘱与共享订阅概念这三个概念不属于某个具体页面但在EMQX里是很核心的功能我在这里先点明白后面的Python实验里再实际验证。保留消息发布消息时设置retainTrueBroker会把该主题的最新一条消息存下来。之后新的订阅者一订阅这个主题就会立刻收到这条保留消息而不是空等。适合用来发布设备状态、版本号之类的“最新值”信息。遗嘱消息客户端在连接时提前设定一个“遗嘱”比如“本设备离线了”。当客户端因为异常断线比如断电、网络断开而不是正常断开时Broker会代替这个客户端把遗嘱消息发出去。这个机制可以用于在线状态监控。共享订阅多个订阅端订阅同一个共享主题后消息会分摊到不同订阅端而不是每个订阅端都收到全量消息。格式是$share/组名/主题。这个功能可以用来做消费者负载均衡。5. Python接进来才算落地paho-mqtt实测从订阅到发布5.1 环境准备Python连接EMQX最常用的库是paho-mqtt它是Eclipse Paho项目的一部分API设计简洁社区文档也多。安装命令pip install paho-mqtt安装完成后可以查看版本确认成功pip show paho-mqtt我建议在虚拟环境里操作避免污染系统Python环境。项目里建一个.venv目录python -m venv .venv source .venv/bin/activate # Windows下执行 .venv\Scripts\activate pip install paho-mqtt5.2 最小连接代码先跑通连接先写一个最简脚本只做连接和保持在线import paho.mqtt.client as mqtt def on_connect(client, userdata, flags, rc): if rc 0: print(连接成功) else: print(连接失败返回码, rc) client mqtt.Client() client.on_connect on_connect # 如果开启了认证取消下面两行注释 # client.username_pw_set(你的用户名, 你的密码) client.connect(127.0.0.1, 1883, 60) client.loop_forever()这段代码的关键点on_connect回调函数会在连接成功或失败时被调用rc是返回码0表示成功。client.connect第三个参数60是keepalive心跳间隔单位是秒。最后一行的loop_forever是阻塞式事件循环让程序一直处理网络报文保持连接不退出。运行这个脚本控制台输出“连接成功”说明Python已经能和EMQX通信了。此时回到Dashboard的“客户端”页面应该能看到一个新增的在线客户端。5.3 实现订阅端订阅端的作用是监听某个主题的消息。新建一个subscriber.pyimport paho.mqtt.client as mqtt def on_connect(client, userdata, flags, rc): print(连接成功开始订阅) client.subscribe(sensor/temperature, qos1) def on_message(client, userdata, msg): print(f主题: {msg.topic}, 消息: {msg.payload.decode()}) client mqtt.Client() client.on_connect on_connect client.on_message on_message client.connect(127.0.0.1, 1883, 60) client.loop_forever()这里subscribe的qos1表示订阅端的QoS等级后面我会专门解释QoS的含义。现在只需要知道on_message回调会在每次收到消息时触发里面参数msg包含主题、消息内容、QoS等信息。先运行订阅端脚本让它保持监听状态。这个阶段会一直阻塞所以建议另开一个终端窗口来跑发布端。5.4 实现发布端再新建一个publisher.pyimport paho.mqtt.client as mqtt import time client mqtt.Client() client.connect(127.0.0.1, 1883, 60) client.loop_start() for i in range(10): payload f温度 {20 i} 度 client.publish(sensor/temperature, payload, qos1) print(已发布, payload) time.sleep(1) client.loop_stop() client.disconnect()执行发布端后再看订阅端的终端窗口应该能连续收到10条消息。运行流程是先起订阅端再起发布端订阅端实时打印收到的消息。我在实际操作中体会到loop_start会在后台开启一个线程循环处理网络事件主线程可以继续做别的事情比如循环发布消息。如果代码里忘了调用loop_start或loop_foreverpublish 之后消息可能发不出去这是新手最容易犯的错误。5.5 进阶测试遗嘱消息与保留消息遗嘱消息验证起来很有意思。先写一个客户端脚本连接时设置遗嘱import paho.mqtt.client as mqtt def on_connect(client, userdata, flags, rc): print(已连接) client.subscribe(device/status) def on_message(client, userdata, msg): print(f收到状态: {msg.payload.decode()}) client mqtt.Client() client.will_set(device/status, 设备异常离线, qos1, retainTrue) client.on_connect on_connect client.on_message on_message client.connect(127.0.0.1, 1883, 60) client.loop_forever()运行这个脚本后打开Dashboard在客户端列表里找到这个连接然后直接踢掉它。也可以更粗暴一点直接关闭运行脚本的终端窗口。此时另一个订阅了device/status主题的客户端就会收到“设备异常离线”这条遗嘱消息。保留消息的验证更简单。发布端发布一条带保留标志的消息import paho.mqtt.client as mqtt client mqtt.Client() client.connect(127.0.0.1, 1883, 60) client.loop_start() client.publish(device/latest, 最新状态online, qos1, retainTrue) print(已发布保留消息) client.loop_stop() client.disconnect()然后重新启动一个订阅端订阅device/latest你会发现它一订阅就立刻收到了这条消息即使发布端已经停止运行了。这就是保留消息的作用它让Broker记录每个主题的最后一条消息给后续订阅者提供“初始快照”。5.6 用MQTTX做辅助验证除了自己写Python代码我还习惯用MQTTX这个图形化客户端做交叉验证。MQTTX支持TCP、WebSocket连接可以同时开多个连接窗口手动发布和订阅非常直观。连接时填上Broker地址127.0.0.1端口1883客户端ID随便写一个不重复的字符串。如果EMQX开启了认证需要在连接设置里填用户名和密码。MQTTX的价值在于当Python代码收不到消息时用它手动发一条能快速判断问题出在Broker还是Python代码。我在调试阶段经常这么干先让MQTTX订阅某个主题再用Python发布如果可以收到说明Broker路由正常问题在Python客户端订阅逻辑。反过来再用Python订阅MQTTX发布就能定位到具体是哪一端写错了。6. 运行一段时间后发现的问题与调优配置6.1 连接闪断与心跳参数测试环境一切正常一到真实网络环境就经常掉线这是MQTT项目最常见的烦恼。原因多半是网络设备中间层把长时间空闲的TCP连接断开了。比如NAT网关或运营商防火墙经常会清理空闲连接。解决办法是合理设置心跳。客户端连接时第三个参数就是这个值我建议根据网络稳定性设置网络稳定设置60秒到120秒。移动网络或弱网环境建议30秒到45秒既能及时感知断线又不会产生太多空包。EMQX侧也会主动检测心跳超时相关参数在配置文件里但大多数情况下不建议频繁改服务端心跳参数优先调客户端这边的keepalive。6.2 消息丢失与QoS选择MQTT的QoS有三个等级选择不当是消息丢失或重复的根本原因QoS 0最多一次消息可能丢失适合对实时性要求高但允许丢失的场景。QoS 1至少一次保证送达但可能重复需要业务端做幂等处理。QoS 2恰好一次防止重复和丢失但确认流程复杂开销最大。我在生产环境一般用QoS 1因为大部分设备数据丢失是不可接受的而重复消息通过增加消息ID或时间戳可以轻松去重。千万不要为了性能全部用QoS 0尤其是控制指令类消息丢了可就不是数据误差的问题了。另外一个容易踩坑的是“持久会话”。如果客户端连接时clean_sessionFalseBroker会保存该客户端的订阅关系和离线消息重连后可以继续收到离线期间的消息。对于需要可靠接收的设备端建议开启持久会话代价是Broker内存占用会增加。6.3 性能调优与系统限制EMQX默认配置针对大多数场景已经够用但如果你要连接几千甚至上万设备需要关注文件描述符限制。Docker方式部署时默认容器文件描述符数量有限可以在docker run时加上docker run -d --name emqx \ --ulimit nofile1048576:1048576 \ -p 1883:1883 -p 18083:18083 \ emqx/emqx:5.8.4Linux系统层面也要检查ulimit -n如果过低需要修改系统限制。这个参数直接决定了EMQX能建立的TCP连接数上限连接数上不去时优先检查这里。另外EMQX 5.x的配置都集中在emqx.conf比如最大连接数、监听器并发数等。没有明确需求前不建议随手改改错一个参数可能导致服务启动失败。改之前先备份原文件改完用emqx ctl check之类的命令检查配置。6.4 常见坑和排查路径我把自己踩过和帮朋友排查过的问题集中整理一下连接被拒绝报Connection refused先看端口启动没再看客户端连的地址是不是127.0.0.1如果是远程服务器检查防火墙和安全组是否放行1883端口。能连接但登录失败检查认证是否开启用户名密码是否正确。EMQX默认不打印密码错误细节需要看日志确认。订阅不到消息用Dashboard的订阅页面确认主题是否一致注意MQTT主题是大小写敏感的。QoS设置为2但性能下降明显确认业务是否真的需要QoS 2大多数场景QoS 1足够。Docker容器重启后EMQX没起来docker run时加上--restartalways这样宿主机重启后容器会自动拉起。Dashboard登录密码忘了在EMQX的bin目录执行emqx ctl admins passwd admin 新密码重置。排查问题时我习惯的路径是先确认网络通不通再看端口和日志最后用Dashboard和MQTTX交叉验证客户端连接和主题订阅。顺着这个链路走一遍大部分问题都能快速定位。我自己的经验是部署EMQX本身不难难的是真正理解消息路由和连接生命周期。把上面这些基础功能跑通一遍之后再做业务层的消息设计会顺手很多。如果你也准备在自己的项目里接入EMQX建议先别急着堆功能花一个下午把这些基本功过一遍后面省下的时间绝对物超所值。
返回列表