做物联网设备接入的时候,最头疼的就是设备端和服务端怎么稳定通信。设备数量一多、上下线一频繁,自己维护的WebSocket服务很快就不行了。后来我把消息层整体切到了EMQX——一款基于Erlang/OTP开发的开源MQTT消息服务器,这些问题才算真正解决下来。这篇文章我就从零开始,把EMQX的简介、安装部署、基础功能和Python代码测试全过程整理出来,给正在调研或准备上手的你做个参考。
不管你是刚开始接触MQTT的嵌入式开发者,还是准备自建消息服务的后端工程师,这篇文章都够用。我不会只贴官方文档,而是把我实际部署和测试时踩过的坑、验证过的方案都写出来,你可以直接照着操作。
1. EMQX解决了什么问题:核心定位与选型逻辑
1.1 MQTT协议和EMQX到底是什么
先说MQTT。它是一套专门为物联网场景设计的轻量级发布订阅消息协议,运行在TCP之上,设计目标就三个:省带宽、省电量、适应不稳定网络。设备端和服务端解耦,谁都不需要知道对方IP和在线状态,只要约定好主题,就能互相通信。这套机制非常适合传感器采集、远程控制、移动推送这类场景。
EMQX就是MQTT协议的服务端实现,大家俗称消息Broker。它干的事情很简单也很核心:接收客户端发来的消息,按主题规则路由给所有订阅了该主题的客户端。设备A发一条消息到主题sensor/temp,服务端和App只要订阅了这个主题,就能实时收到。你不需要关心设备IP、不需要知道设备用的什么语言、不需要维护长连接池,只管收发消息就行。
EMQX底层基于Erlang/OTP语言,这个语言最牛的地方是天然支持高并发、高可用和热更新。所以EMQX一台单机就能扛百万级连接,这在Java、Go写的消息服务器里并不容易做到。官方描述是“百万级并发连接,毫秒级消息转发”,我实际部署在4核8G的云服务器上,跑了几千个模拟客户端,消息延迟基本在个位数毫秒级别,确实没让人失望。
1.2 为什么选EMQX:横向对比后我的结论
市面上MQTT服务器不少,Mosquitto老牌经典,RabbitMQ也支持MQTT插件,网上还有EMQX和VerneMQ等选择。我做一个简单的对比,这些都是我实际用过或调研过的:
| 方案 | 优势 | 劣势 | 适用场景 |
|---|---|---|---|
| Mosquitto | 轻量、部署简单、内存占用极小 | 功能较基础,高并发和集群能力弱 | 原型验证、设备量少的项目 |
| RabbitMQ MQTT插件 | 复用已有消息队列设施 | 原生AMQP协议服务,MQTT只是插件,性能和功能受限 | 团队已有RabbitMQ且设备量不大 |
| EMQX | 原生MQTT协议,高并发、规则引擎、数据集成、集群强大 | 相对较重,部署学习成本略高 | 物联网项目、车联网、智能家居、工业采集 |
| 自研协议服务 | 完全可控 | 投入巨大,容易踩坑,稳定性难保证 | 需求极度特殊 |
我最终选EMQX的原因就两条。第一,它原生实现了MQTT 3.1、3.1.1和5.0协议,协议合规性最好,兼容性没问题。第二,它自带规则引擎和数据集成功能,消息进来之后可以直接转发到MySQL、Kafka、InfluxDB等外部系统,省了再写一套消息消费服务的成本。对标Mosquitto,虽然轻量但要对接数据库还得自己写中间件,设备量上来后性能也吃紧。
1.3 哪些场景最适合用EMQX
结合我个人的项目经验,下面这几类场景用EMQX是特别合适的:
- 智能家居:门锁、传感器、开关状态上报,App远程下发指令,要求消息实时、稳定。
- 车联网:车辆GPS轨迹上报、远程诊断、OTA指令下发,一个平台接入几万辆车很常见。
- 工业数据采集:PLC、传感器通过网关接入,采集数据需要转发到MES、数据库或监控平台。
- 移动推送:App服务端下发通知,移动端网络不稳定,MQTT的持久会话正好应对。
- 教学和实验:想深入理解MQTT协议,自己搭建一套完整环境,EMQX的Dashboard能直观看到消息流转。
如果你的场景只是几台设备传点数据,用Mosquitto就行。但只要是正经做产品、设备量预期会增长、后续还要做数据处理和分析,我建议直接上EMQX,省得后面迁移。
2. 安装部署:两种主流方式的完整记录
2.1 前置准备:下载与端口规划
安装之前,先把版本搞清楚。EMQX 4.x是老版本,5.x是重写后的新版本,两者配置格式和Dashboard界面差异很大。建议新项目直接用5.x,官方对5.x的维护周期也长。我这边写这篇文章时用的是EMQX 5.8.4,实际你用的时候去官网或者GitHub Releases看最新稳定版就行,命令里的版本号替换一下即可。
EMQX启动后会监听好几个端口,安装前最好提前规划,避免和已有服务冲突:
| 端口 | 作用 |
|---|---|
| 1883 | MQTT TCP协议端口,设备接入主端口 |
| 8883 | MQTT SSL/TLS加密端口 |
| 8083 | MQTT over WebSocket端口,浏览器端常用 |
| 8084 | MQTT over WSS加密WebSocket端口 |
| 18083 | Dashboard管理控制台端口,网页访问 |
我用Docker居多,因为升级回滚都方便。如果你要在Windows上玩,Docker Desktop也是个不错的选择。下面两种方式,你按自己的环境选。
2.2 方式一:Docker快速部署
服务器上先装好Docker,然后执行:
docker run -d --name emqx \ -p 1883:1883 \ -p 8883:8883 \ -p 8083:8083 \ -p 8084:8084 \ -p 18083:18083 \ emqx/emqx:5.8.4跑起来之后,用docker ps确认容器状态是Up,然后浏览器访问http://服务器IP:18083,看到登录页就说明部署成功了。默认账号是admin,默认密码是public,第一次登录后系统会引导你修改,这一步不要跳过。
我实际部署时发现一个细节:如果你只想快速验证,端口只映射1883和18083就够了,其他端口后面用到再加。因为多映射一个端口就多一个暴露面,生产环境尤其要注意缩小监听范围。另外,容器数据默认是放在容器层里的,一旦容器删除,配置和数据都没了。稳妥的做法是用-v参数把/opt/emqx/data和/opt/emqx/etc挂载到宿主机,方便持久化和备份。
2.3 方式二:Linux裸机部署与systemd托管
有些团队安全要求严格,不让用Docker,那就走传统部署。以Ubuntu 22.04为例,先到EMQX官网下载对应的包,或者用wget直接拉:
wget https://www.emqx.com/zh/downloads/broker/5.8.4/emqx-5.8.4-ubuntu22.04-amd64.tar.gz tar -zxvf emqx-5.8.4-ubuntu22.04-amd64.tar.gz cd emqx ./bin/emqx start启动后执行./bin/emqx_ctl status,看到Node 'emqx@127.0.0.1' is started就说明起来了。停止用./bin/emqx stop,重启用./bin/emqx restart,这些命令在bin目录下都有。
裸机部署还有个坑:默认EMQX是以前台进程方式挂着,你关掉SSH窗口可能就把进程带死了。虽然./bin/emqx start本身是后台启动,但服务器重启后并不会自动拉起。所以我建议配置一个systemd服务,让EMQX开机自启、崩溃自动重启:
[Unit] Description=EMQX Broker After=network.target [Service] Type=forking ExecStart=/opt/emqx/bin/emqx start ExecStop=/opt/emqx/bin/emqx stop ExecReload=/opt/emqx/bin/emqx restart Restart=on-failure RestartSec=10 [Install] WantedBy=multi-user.target把上面的内容存到/etc/systemd/system/emqx.service,然后执行:
systemctl daemon-reload systemctl enable emqx systemctl start emqx这样就算服务器重启,EMQX也会自动恢复。我以前图省事直接手动启动,结果一次服务器维护重启后整个采集链路断了,业务方找上门来才意识到进程托管的重要性,血的教训。
2.4 部署后的初始化设置
启动完成后,建议按我下面的顺序做初始化设置:
打开Dashboard,在“系统设置”里先把默认密码改掉,这是最基本的安全操作。然后去“管理者”标签页创建一个普通用户,用于日常登录查看,避免长期使用admin操作。
接着检查Dashboard上的“监听器”页面,确认1883端口监听的是0.0.0.0还是127.0.0.1。如果只在本机测试,监听127.0.0.1没有问题;如果要对外提供服务,必须监听0.0.0.0。默认就是0.0.0.0,但这一步确认一下,很多连不上的问题都是监听地址不对导致的。
还有个容易被忽略的点:云服务器默认有安全组规则。就算EMQX本身监听没问题,安全组没放行1883端口,外部设备一样连不进来。我调试时经常碰到“本地能连、远程连不上”,十有八九是云控制台的安全组规则漏配了。
3. 基础功能拆解:从连接会话到规则引擎
3.1 Dashboard界面怎么用
EMQX 5.x的Dashboard做得相当直观,左侧菜单分为监控、访问控制、管理、集成等模块。
“监控”页面能实时看到当前连接数、订阅数、消息收发速率、丢弃消息数等指标。我测试时会一直盯着这个页面,客户端一接入,连接数立刻变化,消息发布也能看到收发速率跳动,整体反馈非常实时。
“客户端”页面列出了当前所有连接上的设备和它们的IP、Client ID、协议版本、Keep Alive等信息。排查设备掉线、客户端异常时,这个页面就是第一排查入口。
“主题”页面可以查看当前被订阅的主题和订阅关系。我曾经遇到过一种情况:设备数据一直上报,但服务端就是收不到,排查半天发现是设备和平台上订阅的主题不一致,差了末尾一个层级。在Dashboard的主题页上对照着一看就明白了。
3.2 连接认证与访问控制
默认情况下,EMQX允许任何人连接,只要有Broker地址和端口就能发布订阅任意主题。这在公网上等于裸奔,所以认证必须开。
EMQX 5.x的认证方式有内置数据库认证、MySQL认证、JWT认证等。最简单的就是内置数据库认证,在Dashboard的“访问控制 -> 认证”里添加一个认证器,选择“内置数据库”,然后手动添加用户名和密码。客户端连接时必须携带这对账号密码,认证不通过直接拒绝连接。
更细的权限控制叫ACL,可以控制某个用户只能发布或只能订阅某些主题。比如生产环境我给采集网关配了“只允许发布到device/{clientid}/data”这个主题的ACL,这样即使设备被攻破,影响范围也有限。ACL规则在“访问控制 -> 授权”里配置,语法很直白:选动作(发布/订阅)、填主题过滤器、选结果(允许/拒绝)。
提示:认证和授权是两个层面。认证管“你是谁”,授权管“你能干什么”。只做认证不做授权,内部用户越权操作的风险依然存在。物联网设备安全风险很高,这两步建议都配上。
3.3 消息背后的机制:QoS、遗嘱、保留消息与共享订阅
这部分是MQTT协议最核心也最容易理解偏的地方,我挑重点讲。
QoS(服务质量)是消息可靠性的等级,分0、1、2三档。QoS 0是尽力而为,发出去不管,可能丢;QoS 1保证消息至少到达一次,但可能重复;QoS 2保证只到达一次,最严格但开销最大。实际项目里,设备状态上报用QoS 1足够,控制指令如果想保证不丢又不怕重复处理,也选QoS 1,在业务层做去重就行。QoS 2一般用在订阅关系变更这类关键控制消息上。
遗嘱消息(Last Will)是MQTT一个非常实用的机制。设备连接时可以在connect报文里带上遗嘱主题和遗嘱内容,当设备异常断开(网络超时、心跳过期)时,Broker会代替设备向遗嘱主题发送这条消息。其他订阅者收到后就知道设备掉线了,可以触发告警。
保留消息(Retained)也很常用。发布消息时如果设置retain标志为true,Broker会保存这条消息的主题和内容,以后任何客户端订阅这个主题时,第一件事就是收到这条保留消息。这正好解决“新订阅者错过历史消息”的问题,比如设备状态、App端最新一条配置,都可以用保留消息让新订阅者立刻拿到当前状态。
共享订阅是另一个实用功能。多个客户端订阅同一个共享订阅主题($share/{group}/{topic}),消息会被分发给其中一个客户端,实现负载均衡。我在产品里用共享订阅,服务端部署多个消费者实例处理同一批设备数据,既提升了吞吐量,又实现了一定程度的高可用。
3.4 规则引擎和数据集成简述
EMQX 5.x的规则引擎是我提到选型理由时的重点能力。它可以在消息流经Broker时做实时处理,然后转发到外部系统。比如一条温度传感器数据发到device/1/data,规则引擎可以抽出JSON里的温度字段,写入InfluxDB时序数据库,同时转发到Kafka分区。这个能力非常实用,因为物联网项目里消息进来之后通常都要落库分析,有了规则引擎就不需要自己再写一个消费程序转存数据了。
配置路径是Dashboard里的“集成 -> 规则”,创建规则时先用SQL语句筛选和提取消息字段,然后添加动作,选数据转发到哪个连接器(数据库、Kafka、HTTP服务等)。它还有个调试功能,可以模拟输入测试规则,这个对新手非常友好,至少省了我不少来回测试的时间。
需要说明的是,规则引擎虽然方便,但别一上来在单机上接几十个规则,每条消息都过一遍规则是有CPU开销的。我见过有同事把每条设备消息同时转发到MySQL、MongoDB、Redis,单机立马撑不住。合理的设计是只把真正需要落库或者转发的消息接出来,能省则省。
4. Python代码测试:完整的发布订阅实战
4.1 环境准备:安装paho-mqtt
终于到代码部分了。Python连EMQX最常用的库是paho-mqtt,它是Eclipse Paho项目下的Python客户端实现,API稳定、资料丰富。安装一条命令搞定:
pip install paho-mqtt如果服务器上有多个Python版本,记得用pip3,或者python3 -m pip install paho-mqtt,避免装到错误的解释器上。
测试前我先确认一下EMQX是否在运行:现在假定Broker地址是127.0.0.1,端口是1883,认证暂时没开。如果之前开了认证,代码里就要加上username_pw_set("用户名", "密码"),下面代码里我会留出位置。
注意:请确保代码运行环境和EMQX所在主机的1883端口网络是通的。本地测试用127.0.0.1没问题,远程测试记得检查防火墙和安全组。
4.2 写一个订阅端脚本
先写订阅端,它会连接Broker并订阅主题test/hello,实时打印收到的消息。把下面的代码保存为sub.py:
import paho.mqtt.client as mqtt BROKER_HOST = "127.0.0.1" BROKER_PORT = 1883 TOPIC = "test/hello" # 如果开了认证,解开下面这行并换成真实账号密码 # AUTH_USERNAME = "admin" # AUTH_PASSWORD = "password123" def on_connect(client, userdata, flags, rc): if rc == 0: print("连接成功,准备订阅主题") client.subscribe(TOPIC, qos=1) else: print(f"连接失败,返回码:{rc}") def on_message(client, userdata, msg): print(f"收到消息 -> 主题: {msg.topic} | QoS: {msg.qos} | 载荷: {msg.payload.decode()}") def on_disconnect(client, userdata, rc): print("连接已断开") client = mqtt.Client(client_id="python_sub_001") client.on_connect = on_connect client.on_message = on_message client.on_disconnect = on_disconnect # client.username_pw_set(AUTH_USERNAME, AUTH_PASSWORD) client.connect(BROKER_HOST, BROKER_PORT, keepalive=60) client.loop_forever()这段代码有三个回调函数,分别是连接成功回调on_connect、收到消息回调on_message、断开连接回调on_disconnect。MQTT客户端是事件驱动的,你注册回调,库在对应的时机自动调用。loop_forever()会阻塞当前线程,自动维持心跳、处理收包和重连。
有个细节需要注意:on_connect里最好只做连接成功后的订阅操作,不要在连接之前调用subscribe,因为此时客户端和Broker之间的会话还没建立,订阅在协议层面是不成立的。很多新手在这里踩坑,我一开始也犯过。
4.3 写一个发布端脚本
再写发布端,保存为pub.py:
import paho.mqtt.client as mqtt BROKER_HOST = "127.0.0.1" BROKER_PORT = 1883 TOPIC = "test/hello" # AUTH_USERNAME = "admin" # AUTH_PASSWORD = "password123" client = mqtt.Client(client_id="python_pub_001") # client.username_pw_set(AUTH_USERNAME, AUTH_PASSWORD) client.connect(BROKER_HOST, BROKER_PORT, keepalive=60) client.loop_start() result = client.publish(TOPIC, "hello emqx, this is python test", qos=1) status = result.rc if status == mqtt.MQTT_ERR_SUCCESS: print("消息发布成功") else: print(f"消息发布失败,错误码:{status}") client.disconnect() client.loop_stop()这里我用了loop_start()而不是loop_forever(),区别在于loop_start()会在后台开一个新线程维护网络循环,主线程可以继续往下走。发完消息后直接disconnect()再loop_stop(),脚本就退出了。如果是发布短消息的测试场景,这样写比loop_forever()更干净。
publish()的返回值result.rc等于MQTT_ERR_SUCCESS(即0)时说明消息已经成功交给底层socket发送。需要提醒的是,这个返回只代表发送动作成功,不代表Broker确认收到。真正要确认QoS 1消息到达Broker,需要监听on_publish回调,里面消息的mid和publish()返回的mid对应上才算数。
4.4 联调:观察消息、Dashboard和回调
现在开始联调。先启动订阅端:
python3 sub.py终端输出“连接成功,准备订阅主题”,然后脚本就挂在那里等消息。再开一个终端跑发布端:
python3 pub.py发布端输出“消息发布成功”,同时订阅端的窗口会打印:
收到消息 -> 主题: test/hello | QoS: 1 | 载荷: hello emqx, this is python test看到这行输出,说明发布-订阅链路已经完全打通了。这个消息从Python客户端发出,经过EMQX路由,再投递给另一个Python客户端,整个链路没有任何其他中间层。此时打开Dashboard的“监控”页面,你会看到消息收发计数在上涨,“客户端”页面会同时出现python_sub_001和python_pub_001两个Client ID的连接。
我还建议你在这基础上多玩几个功能:把发布端的消息保留标志设为retain=True,然后重新启动订阅端,看看是不是一订阅就马上收到这条旧消息。或者设置遗嘱消息,手动断开设备连接,看看遗嘱主题是否能触发。这些测试做完,你对MQTT协议的理解会有一个质的提升。
4.5 同步订阅、批量发布等进阶测试方向
当基础收发没问题后,还可以再扩展几个测试场景,这些也是我实际项目里经常用到的。
同步订阅方式:client.subscribe对于QoS 1的订阅,在5.x里有一个client.subscribe_callback的快捷方式,也可以直接用message_callback_add给不同主题注册不同的回调。这样多个主题可以走进不同的处理函数,类似路由功能,比在on_message里自己写一长串if-else要灵活得多。
批量模拟设备:用循环发消息是常规操作,但这里有个坑:在短连接场景下,每次循环都connect和disconnect,如果间隔特别短,Broker侧可能会出现TIME_WAIT连接堆积。更好的做法是长连接批量发布,或者用异步客户端paho.mqtt.client.Client加上loop_start()跑一个长任务。我曾经写过一个脚本,开20个线程,每个线程一个客户端,模拟20台设备同时上报数据,压测EMQX转发能力,验证结果相当直观。
Python异步场景还可以考虑asyncio-mqtt这个库,它是基于paho-mqtt封装的支持async/await语法的库。如果你的服务端是FastAPI或者Tornado这类异步框架,用它接入EMQX会舒服很多,代码风格一致,不会阻塞事件循环。
5. 常见问题与排查技巧实录
5.1 连接失败排查
连接不上是出现频率最高的异常。按我经验,先看报错类型:如果是超时,多半是网络不通;如果立即拒绝,多半是端口或地址不对。
| 现象 | 可能原因 | 排查方法 |
|---|---|---|
| 连接超时 | 防火墙拦截、安全组没放行、Broker地址错误 | 在客户端机器上执行telnet 服务器IP 1883测试连通性 |
| 立即断开 | 客户端ID冲突、认证失败 | 换一个Client ID重试,检查用户名密码和ACL规则 |
| 偶尔能连偶尔连不上 | Keep Alive设置过短,网络抖动 | 把keepalive适当调大,比如60秒 |
| 能连但没消息 | 订阅主题不一致、QoS设置不对 | Dashboard“主题”页面查看实际订阅关系 |
还有一条很隐蔽的问题:同一Client ID的客户端同时在线时,后连接的会把先连接的踢下线。原因就是MQTT协议规定同一Client ID只能有一个会话。排查时如果发现设备频繁掉线,可以去Dashboard“客户端”页面看是否有两个相同的Client ID反复上下线。
5.2 收不到消息排查
收不到消息十有八九是主题问题,而不是配置问题。MQTT主题是层级结构,用/分隔。比如设备发布到device/1/data,订阅端必须订阅device/1/data才能收到,订阅device/#或device/+/data也能收到,但订阅device/1就什么都收不到。
我整理了一个排查清单:
- 发布端和订阅端是否都在同一个Broker上,没有连错服务器;
- 主题字符串是否完全一致,包括大小写和尾部斜杠;
- 发布时是否设置了
retain=True,如果是,新订阅者应该立刻收到,但收到的是历史消息,可能被误认成“收不到”; - 订阅时QoS和发布时QoS取两者的最小值,比如订阅QoS 0、发布QoS 2,实际收到QoS 0。如果程序里判断QoS等级不匹配就不处理,就会像“没收到”一样;
- Dashboard“主题”页面里有没有对应的订阅关系,没有就是订阅没成功。
5.3 回调不执行的问题
paho-mqtt新手最容易忽略的是loop_forever()或者loop_start()的调用。没有启动网络循环,客户端虽然底层socket连上了,但事件处理线程不运行,on_connect、on_message永远不会触发。我见过不少人在脚本里写完多行connect后直接print等消息,结果消息迟迟不来。
正确写法就是client.loop_forever()阻塞运行,或者client.loop_start()后台启动。如果你的程序在回调里做耗时操作,比如写数据库、调用HTTP接口,注意回调线程会被阻塞,这会影响后续消息接收。稳妥的做法是把耗时操作丢到队列或者线程池里异步处理。
5.4 高连接数下的优化方向
如果你测试时发现连接数上来之后CPU占用偏高或者消息延迟变大,可以从这几个方向调整。
操作系统层,修改文件描述符上限,因为每一条TCP连接都对应一个fd。用ulimit -n查看当前限制,生产环境建议调大到100000以上。EMQX配置文件里也有一个listener.tcp.external.max_connections默认值,5.x版本默认已经很高,但如果你用Docker部署,宿主机的端口范围和fd限制反而可能成为瓶颈。
消息层面,如果客户端场景不要求每次都实时转发,可以开启EMQX的延迟发布功能,或者用外部规则引擎做批量聚合。不过这些属于调优范畴,基础阶段把QoS等级选择合理、避免不必要的大消息体,就能避免大部分性能问题。
5.5 数据持久化和升级的注意事项
EMQX 5.x把配置和数据都放在/opt/emqx/data目录下。用Docker时务必挂载出来,否则docker rm之后等于重装了一次Broker,所有认证用户、ACL规则、数据集成配置全部丢失。裸机部署则建议定期备份data目录。
升级时也要先备份,再停服,解压新版本替换旧版本,然后启动。虽然EMQX支持热升级,但跨大版本升级最好还是走完整的备份-迁移流程。我吃过一次亏,从4.x升5.x时没仔细看迁移文档,配置格式不兼容导致认证全部失效,那个下午光调权限就花了不少时间。
最后分享一个小技巧:EMQX的命令行工具emqx_ctl里藏了很多实用的管理命令,比如./bin/emqx_ctl broker publish可以直接在服务器本机发一条测试消息,./bin/emqx_ctl listeners可以查看所有监听器状态。排查问题时不用急着写脚本,先在这些命令里找找答案,往往更快。这也是我这两年用EMQX攒下的最实在的经验。