Docker Compose 构建 EMQX 集群 实现mqqt 和websocket
EMQX 集群化管理mqqt真香
目录
#目录 /usr/emqx
容器构建
vim docker-compose.yml
version: '3'services:emqx1:image: emqx:5.8.3container_name: emqx1environment:- "EMQX_NODE_NAME=emqx@node1.emqx.io"- "EMQX_CLUSTER__DISCOVERY_STRATEGY=static"- "EMQX_CLUSTER__STATIC__SEEDS=[emqx@node1.emqx.io,emqx@node2.emqx.io]"healthcheck:test: ["CMD", "/opt/emqx/bin/emqx", "ctl", "status"]interval: 5stimeout: 25sretries: 5networks:emqx-bridge:aliases:- node1.emqx.ioports:- 1883:1883- 8083:8083- 8084:8084- 8883:8883- 18083:18083 # volumes:# - $PWD/emqx1_data:/opt/emqx/dataemqx2:image: emqx:5.8.3container_name: emqx2environment:- "EMQX_NODE_NAME=emqx@node2.emqx.io"- "EMQX_CLUSTER__DISCOVERY_STRATEGY=static"- "EMQX_CLUSTER__STATIC__SEEDS=[emqx@node1.emqx.io,emqx@node2.emqx.io]"healthcheck:test: ["CMD", "/opt/emqx/bin/emqx", "ctl", "status"]interval: 5stimeout: 25sretries: 5networks:emqx-bridge:aliases:- node2.emqx.io# volumes:# - $PWD/emqx2_data:/opt/emqx/datanetworks:emqx-bridge:driver: bridge
启动
docker-compose up -d
集群状态
#查看集群状态
docker exec -it emqx1 sh -c "emqx ctl cluster status"#验证
telnet 192.168.0.15 1883
#内网
nc -zv 192.168.0.15 1883 #账户
admin
#默认密码
public
服务开放端口
1883,8083,8084,8883,18083
端口占用
EMQX 默认使用以下端口,请确保这些端口未被其他应用程序占用,并按照需求开放防火墙以保证 EMQX 正常运行。
| 端口 | 协议 | 描述 |
| 1883 | TCP | MQTT over TCP 监听器端口,主要用于未加密的 MQTT 连接。 |
| 8883 | TCP | MQTT over SSL/TLS 监听器端口,用于加密的 MQTT 连接。 |
| 8083 | TCP | MQTT over WebSocket 监听器端口,使 MQTT 能通过 WebSocket 进行通信。 |
| 8084 | TCP | MQTT over WSS (WebSocket over SSL) 监听器端口,提供加密的 WebSocket 连接。 |
| 18083 | HTTP | EMQX Dashboard 和 REST API 端口,用于管理控制台和 API 接口。 |
| 4370 | TCP | Erlang 分布式传输端口,根据节点名称不同实际端口可能是 BasePort (4370) + Offset。 |
| 5370 | TCP | 集群 RPC 端口(在 Docker 环境下为 5369),根据节点名称不同实际端口可能是 BasePort (5370) + Offset。 |
前端js示
<!DOCTYPE html>
<html>
<head><title>MQTT WebSocket Test</title><script src="https://unpkg.com/mqtt/dist/mqtt.min.js"></script>
</head>
<body><script>// 使用提供的客户端 ID 或生成一个唯一 IDconst clientId = 'emqx_NjI4MT2';// 配置 WebSocket MQTT broker 地址const host = 'ws://127.0.0.1:8083/mqtt';// MQTT 连接选项const options = {keepalive: 60, // 心跳时间间隔clientId: clientId,protocolId: 'MQTT', // 协议 IDprotocolVersion: 5, // 使用 MQTT 5 协议clean: true, // 是否清除会话reconnectPeriod: 1000, // 重连间隔时间 (ms)connectTimeout: 30 * 1000, // 连接超时时间 (ms)username: 'admin', // 设置用户名password: 'public', // 设置密码will: {topic: 'pushRanking/1',payload: 'Connection Closed abnormally..!',qos: 0,retain: false},};console.log('Connecting mqtt client');// 连接到 MQTT Brokerconst client = mqtt.connect(host, options);// 连接成功回调client.on('connect', () => {console.log('Connected to MQTT broker');// 订阅主题 pushRanking/#,支持通配符client.subscribe('pushRanking/1', { qos: 0 }, (err) => {if (!err) {console.log('Subscribed to topic: pushRanking/1');} else {console.error('Failed to subscribe:', err);}});});// 处理接收到的消息client.on('message', (topic, message) => {console.log(`Received message from topic "${topic}": ${message.toString()}`);});// 连接错误回调client.on('error', (err) => {console.log('Connection error:', err);client.end();});// 重新连接回调client.on('reconnect', () => {console.log('Reconnecting...');});// 连接关闭回调client.on('close', () => {console.log('Connection closed');});// 模拟消息发布以测试接收setTimeout(() => {client.publish('pushRanking/1', JSON.stringify({ msg: 'hello' }), { qos: 0 });}, 5000);</script>
</body>
</html>
后端
package emqximport ("fmt"mqtt "github.com/eclipse/paho.mqtt.golang""testing""time"
)func TestMQTT(t *testing.T) {// 创建 EMQX 客户端实例client := NewEMQXClient("tcp://127.0.0.1:1883", "test-client", "admin", "QfTzLy3cop9NOGWj")// 连接到 EMQXif err := client.Connect(); err != nil {fmt.Printf("Failed to connect: %v\n", err)return}defer client.Disconnect()// 订阅主题client.Subscribe("testtopic/#", 1, func(client mqtt.Client, msg mqtt.Message) {fmt.Printf("Message received: %s\n", msg.Payload())})// 发布消息client.Publish("testtopic/1", 1, false, "Hello from Golang!")// 保持连接一段时间以接收消息time.Sleep(10 * time.Second)
}/*长连接的场景 DEMOfunc main() {client := emqxclient.NewEMQXClient("tcp://broker.emqx.io:1883", "test-client", "", "")if err := client.Connect(); err != nil {fmt.Printf("Failed to connect: %v\n", err)return}// 使用 defer 确保程序退出时断开连接defer client.Disconnect()// 订阅主题client.Subscribe("test/topic", 1, func(client mqtt.Client, msg mqtt.Message) {fmt.Printf("Received message: %s\n", msg.Payload())})// 发布消息client.Publish("test/topic", 1, false, "Hello from Golang!")// 捕获退出信号signalChan := make(chan os.Signal, 1)signal.Notify(signalChan, os.Interrupt, syscall.SIGTERM)fmt.Println("Running... Press Ctrl+C to exit.")<-signalChanfmt.Println("Exiting...")
}*//*JS 调用如下
import mqtt from 'mqtt';const brokerURL = 'ws://broker.emqx.io:8083/mqtt'; // WebSocket 连接地址
const clientID = `mqttjs_${Math.random().toString(16).substr(2, 8)}`;// 创建客户端
const client = mqtt.connect(brokerURL, {clientId: clientID,username: '', // 如需要认证,填入用户名password: '', // 如需要认证,填入密码
});// 连接事件
client.on('connect', () => {console.log('Connected to EMQX');// 订阅主题client.subscribe('test/topic', (err) => {if (!err) {console.log('Subscribed to topic: test/topic');} else {console.error('Failed to subscribe:', err);}});// 发布消息client.publish('test/topic', 'Hello from JavaScript!');
});// 接收消息事件
client.on('message', (topic, message) => {console.log(`Received message on topic "${topic}": ${message.toString()}`);
});// 错误事件
client.on('error', (err) => {console.error('Connection error:', err);
});*/
common封装调用
package emqximport ("fmt""time"mqtt "github.com/eclipse/paho.mqtt.golang"
)type EMQXClient struct {client mqtt.Client
}// NewEMQXClient 初始化 EMQX 客户端
func NewEMQXClient(broker string, clientID string, username string, password string) *EMQXClient {opts := mqtt.NewClientOptions().AddBroker(broker).SetClientID(clientID).SetUsername(username).SetPassword(password).SetKeepAlive(60 * time.Second).SetDefaultPublishHandler(func(client mqtt.Client, msg mqtt.Message) {fmt.Printf("Received message on topic: %s, message: %s\n", msg.Topic(), msg.Payload())}).SetPingTimeout(1 * time.Second)client := mqtt.NewClient(opts)return &EMQXClient{client: client}
}// Connect 连接到 EMQX
func (c *EMQXClient) Connect() error {token := c.client.Connect()if token.Wait() && token.Error() != nil {return token.Error()}fmt.Println("Connected to EMQX broker")return nil
}// Publish 发布消息
func (c *EMQXClient) Publish(topic string, qos byte, retained bool, payload interface{}) error {token := c.client.Publish(topic, qos, retained, payload)token.Wait()return token.Error()
}// Subscribe 订阅主题
func (c *EMQXClient) Subscribe(topic string, qos byte, callback mqtt.MessageHandler) error {token := c.client.Subscribe(topic, qos, callback)token.Wait()return token.Error()
}// Unsubscribe 取消订阅
func (c *EMQXClient) Unsubscribe(topics ...string) error {token := c.client.Unsubscribe(topics...)token.Wait()return token.Error()
}// Disconnect 断开连接
func (c *EMQXClient) Disconnect() {c.client.Disconnect(250)fmt.Println("Disconnected from EMQX broker")
}

相关文章:
Docker Compose 构建 EMQX 集群 实现mqqt 和websocket
EMQX 集群化管理mqqt真香 目录 #目录 /usr/emqx 容器构建 vim docker-compose.yml version: 3services:emqx1:image: emqx:5.8.3container_name: emqx1environment:- "EMQX_NODE_NAMEemqxnode1.emqx.io"- "EMQX_CLUSTER__DISCOVERY_STRATEGYstatic"- …...
Spring 过滤器:OncePerRequestFilter 应用详解
在Web应用中,过滤器(Filter)是一个强大的工具,它可以在请求到达目标资源之前或响应返回客户端之前对请求或响应进行拦截和处理。然而,在某些情况下,我们可能希望确保过滤器逻辑在一次完整的HTTP请求中仅执行…...
3.CSS字体属性
3.1字体系列 CSS使用font-family属性定义文本的字体系列。 p{font-family:"微软雅黑"} div{font-family:Arial,"Microsoft Yahei",微软雅黑} 3.2字体大小 css使用font-size属性定义字体大小 p{ font-size:20px; } px(像素)大小是我们网页的最常用的单…...
微信小程序 单选多选radio/checkbox 纯代码分享
单选按钮 <radio-group class"radiogroup" bindchange"radioChange"> <label class"radio" wx:for"{{items}}"> <radio value"{{item.name}}" checked"{{item.checked}}" /> {{item.value}} &…...
k8s 部署meilisearch UI
https://github.com/riccox/meilisearch-ui 拉取镜像 sudo docker pull riccoxie/meilisearch-ui:latestk8s 部署 apiVersion: v1 kind: Service metadata:name: meilisearch-uinamespace: meilisearch spec:type: NodePortselector:app: meilisearch-uiports:- port: 24900…...
gitlab 还原合并请求
事情是这样的: 菜鸡从 test 分支切了个名为 pref-art 的分支出来,发布后一机灵,发现错了,于是在本地用 git branch -d pref-art 将该分支删掉了。之后切到了 prod 分支,再切出了一个相同名称的 pref-art 分支出来&…...
ChatGPT最新版本“o3”的概要
o3简介 o3于2024年12月20日发布——也就是OpenAI 12天直播的最后一天。目前处于安全性测试阶段。它是o1的继任者,旨在处理更复杂的推理任务。o3特别针对数学、科学和编程等领域进行了优化。 o3在多项基准测试中表现出色。例如,在ARC-AGI基准测试中&…...
uniapp——App下载文件,保存、打开文件(二)
uniapp如何下载文件、保存、打开文件 时光荏苒,2024即将过去! 迈向2025,祝大家新的一年工作顺利、万事如意,少一点BUG,涨一点工资…↖(ω)↗ 文章目录 uniapp如何下载文件、保存、打开文件下载文件保存并打开文件处理 …...
Postman接口测试05|实战项目笔记
目录 一、项目接口概况 二、单接口测试-登录接口:POST 1、正例 2、反例 ①姓名未注册 ②密码错误 ③姓名为空 ④多参 ⑤少参 ⑥无参 三、批量运行测试用例 四、生成测试报告 1、Postman界面生成 2、Newman命令行生成 五、token鉴权(“…...
【paddle】初次尝试
张量 张量是 paddlepaddle, torch, tensorflow 等 python 主流机器学习包中唯一通货变量,因此应当了解其基本的功能。 张量 paddle.Tensor 与 numpy.array 的转化 import paddle as paddle import matplotlib.pyplot as plt apaddle.to_t…...
01-2023年上半年软件设计师考试java真题解析
1.真题内容 在某系统中,类 Interval(间隔) 代表由下界(lower bound(边界))上界(upper bound )定义的区间。 要求采用不同的格式显示区间范围。 如[lower bound , upper bound ]、[ lower bound … upper bound ]、[ lower bou nd - upper bound &#x…...
一文讲清楚CSS3新特性
文章目录 一文讲清楚CSS3新特性1. 新增选择器特性2. 新增的样式3. 新增布局方式 一文讲清楚CSS3新特性 1. 新增选择器特性 层次选择器(div~p)选择前面有div的p元素伪类选择器 :first-of-type 表示⼀组同级元素中其类型的第⼀个元素:last-of-type 表示⼀组同级元素中其类型的最…...
系统设计案例:设计 Spotify
https://levelup.gitconnected.com/system-design-interview-question-design-spotify-4a8a79697dda 这是一道系统设计面试题,即设计 Spotify。在真正的面试中,你通常会关注应用程序的一两个主要功能,但在本文中,我想从高层次概述…...
太速科技-633-4通道2Gsps 14bit AD采集PCie卡
4通道2Gsps 14bit AD采集PCie卡 一、板卡概述 二、性能指标 板卡功能 参数 内容 ADC 芯片型号 AD9689 路数 4路ADC, 采样率 2Gsps 数据位 14bit 数字接口 JESD204B 模拟接口 交流耦合 模拟输入 1V 连接器 6路 SMA 输入阻抗 50Ω 模拟指…...
图片叠加拖拽对比展示效果实现——Vue版
图片叠加拖拽对比展示效果实现——Vue版 项目中遇见一个需求:2张图片按竖线分割,左右两侧分别展示对应图片,通过滚动条拖动对应展示图片区域;; 网上搜索了下,没有找到直接可用的组件,这里自己封装了一个次功…...
结合长短期记忆网络(LSTM)和无迹卡尔曼滤波器(UKF)的技术在机器人导航和状态估计中的应用前景
结合长短期记忆网络(LSTM)和无迹卡尔曼滤波器(UKF)的技术在机器人导航和状态估计中具有广泛的应用前景。如有滤波、导航方面的代码定制需求,可通过文末卡片联系作者获得帮助 文章目录 结合LSTM和UKF的背景结合LSTM和UKF的优势应用实例研究现状MATLAB代码示例结论结合LSTM和…...
【MATLAB APP Designer】小波阈值去噪(第一期)
代码原理及流程 小波阈值去噪是一种信号处理方法,用于从信号中去除噪声。这种方法基于小波变换,它通过将信号分解到不同的尺度和频率上来实现。其基本原理可以分为以下几个步骤: (1)小波变换:首先对含噪信…...
ClickHouse副本搭建
一. 副本概述 副本的目的主要是保障数据的高可用性,ClickHouse中的副本没有主从之分。所有的副本都是平等的。 副本写入流程: 二. 副本搭建 1. 实验环境 hadoop1(192.168.47.128) hadoop2(192.168.47.129)2. 修改配置文件 修改两台主机/etc/click…...
K3知识点
提示:文章 文章目录 前言一、顺序队列和链式队列题目 顺序队列和链式队列的定义和特性实际应用场景顺序表题目 链式队列 二、AVL树三、红黑树四、二叉排序树五、树的概念题目1左子树右子树前序遍历、中序遍历,后序遍历先根遍历、中根遍历左孩子右孩子题目…...
cocos creator 3.x版本如何添加打开游戏时首屏加载进度条
前言 项目有一个打开游戏时添加载入进度条的需求。这个功能2.X版本是自带的,不知为何在3.X版本中移除了。 实现 先说一下解决思路,就是在引擎源码加载场景的位置插入一个方法,然后在游戏入口HTML处监听即可。 1.找到对应源码脚本 在coco…...
白细胞介素-6(IL-6)在临床疾病中的作用机制与靶向治疗研究进展
白细胞介素-6(Interleukin-6, IL-6)是一种由多种细胞(如单核/巨噬细胞、T细胞、成纤维细胞等)分泌的多效性细胞因子,参与免疫调节、炎症反应、代谢稳态及组织修复等生理过程。在病理状态下,IL-6过度表达与感…...
AI 智能体 8 层架构:生产级系统构建指南
AI 智能体(Agentic AI)革命的关键不在更好的提示词,而在于系统化的架构设计。随着企业竞相部署能够自主感知、推理、规划和行动的 AI 智能体(AI Agent),真正的挑战已经从"我们能构建吗?“转变为"…...
端侧AI算力瓶颈与优化企业格局解析
一、引言:端侧AI的发展困境与研究核心1.1 端侧AI的产业价值与普及现状端侧AI作为边缘计算的核心落地形态,正深度渗透工业制造、智能终端、车载电子、安防监控等领域。据IDC数据,2025年全球端侧AI芯片市场规模突破180亿美元,工业端…...
嵌入式工程师高薪进阶指南:从软硬兼通到系统思维的跨越
1. 嵌入式行业的现状与人才困境最近几年,和不少同行、猎头以及企业招聘负责人聊下来,一个共识越来越清晰:嵌入式这个行当,正在经历一场深刻的“冰火两重天”。一方面,得益于树莓派、Arduino这类高度集成、生态友好的开…...
别再手画电路图了!用Fritzing快速搞定Arduino项目接线图(附传感器库文件下载)
告别手绘时代:Fritzing高效绘制Arduino接线图的完整指南 在Arduino项目开发中,清晰的接线图不仅是项目文档的重要组成部分,更是团队协作和后期维护的关键参考。传统的手绘方式不仅效率低下,还容易出错,尤其当项目涉及多…...
别再死记硬背了!用这5个真实案例,彻底搞懂NumPy的einsum函数
别再死记硬背了!用这5个真实案例,彻底搞懂NumPy的einsum函数 当你第一次看到np.einsum(ij,jk->ik, A, B)这样的表达式时,是不是感觉像在破译外星密码?作为NumPy中最强大却也最令人困惑的函数之一,einsum(…...
D1016UK,1MHz至1GHz宽带适用的低噪声高效率射频功率晶体管
简介今天我要向大家介绍的是 TT Electronics/Semelab 的DMOS RF FET晶体管——D1016UK。这是一款专为VHF/UHF通信频段(1 MHz至1GHz)设计的金金属化多用途硅RF功率场效应管,采用推挽式架构,在28V工作电压下可提供40W的输出功率。作…...
AI大模型Agent面试,超详细(附答案)!
AI大模型Agent面试,超详细(➕答案)!假如你从2026年开始学大模型,按这个步骤走准能稳步进阶。 接下来告诉你一条最快的邪修路线, 3个月即可成为模型大师,薪资直接起飞。阶段1:大模型基础阶段2:RA…...
告别繁琐配置!用EB和S32DS快速搭建AutoSar MCAL基础工程(附完整文件结构解析)
从零构建AutoSar MCAL工程:EB与S32DS高效协作实战指南 当第一次打开AutoSar MCAL的官方示例工程时,多数工程师都会被密密麻麻的文件夹和配置文件淹没。Base、Platform、ECUC、MemIf等模块交织在一起,而EB生成的generate文件夹里又充斥着大量看…...
深入YOLOv8损失函数:为什么自带的Focal Loss会报错?一次完整的源码调试与修复记录
深入YOLOv8损失函数:为什么自带的Focal Loss会报错?一次完整的源码调试与修复记录 在目标检测领域,YOLOv8凭借其卓越的性能和易用性赢得了广泛关注。然而,当开发者尝试深入模型内部机制时,往往会遇到一些意料之外的挑战…...
