尧图网络科技YAOTU DIGITAL 获取报价
获取报价
首页 / 资讯中心 / 文章详情

TDengine 对接 EMQX Broker:用 MQTT 消息零代码写入时序数据的完整实践

发布时间:2026/9/14 19:27:40

资讯中心
01
ARTICLE

TDengine 对接 EMQX Broker:用 MQTT 消息零代码写入时序数据的完整实践

TDengine 对接 EMQX Broker:用 MQTT 消息零代码写入时序数据的完整实践
TDengine 对接 EMQX Broker用 MQTT 消息零代码写入时序数据的完整实践【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine本文介绍如何利用 EMQX Broker 的规则引擎将 MQTT 主题上的消息通过 WebHook 转发到 taosAdapter 的 REST 接口实现 MQTT 数据免代码直接写入 TDengine。全文覆盖前置环境准备、EMQX 规则与资源配置、模拟数据发送程序编写以及数据落库后的验证方法读完后你可以独立搭建一条“设备MQTT 客户端→ EMQX → taosAdapter → TDengine”的物联网数据接入链路并理解其中每个环节的配置要点与底层转发原理。整体链路与原理MQTT 是物联网场景中最主流的数据传输协议EMQX 是一个开源的 MQTT Broker 软件。TDengine 本身不直接监听 MQTT 端口而是通过 taosAdapter 组件暴露的 REST 接口接收数据。整条链路的分工如下MQTT 客户端传感器/模拟程序将 JSON 报文发布到约定主题如sensor/data默认走tcp://127.0.0.1:1883EMQX 规则引擎对匹配主题的消息执行一条 SQL 选择语句抽取payload整个消息体并触发绑定到该规则上的“动作Action”WebHook 动作 → 资源Resource动作以 HTTP POST 方式把模板渲染后的 SQL 发送给指定 URLtaosAdapter REST 接口URL 即 taosAdapter 提供的http://fqdn:6041/rest/sql6041 为httpPort配置项缺省值可参考 taosAdapter 参考手册携带 Basic 认证请求头body 为完整 INSERT 语句taosd执行 SQL 并将数据写入对应数据库与表。从源码结构看taosAdapter 的 RESTful 接口是对 TDengine C 接口的 HTTP 封装因此 EMQX 发送的 body 必须是一条完整 SQLSQL 中的表名需带数据库前缀如test.sensor_data否则会因“HTTP 模块没有当前 DB 的概念”而报错——这一约束在 TDengine REST API 文档中有明确说明也是配置 EMQX 动作模板时最容易踩的坑之一。此外 EMQX 企业版提供原生 TDengine 驱动实现直接保存社区版则使用本文的“发送到 Web 服务”方式。前置条件要让 EMQX 能正常添加 TDengine 数据源需要以下几方面准备TDengine 集群已经部署并正常运行taosAdapter 已经安装并正常运行具体细节参考 taosAdapter 参考手册如果使用本文介绍的模拟写入程序需要安装合适版本的 Node.js推荐 v12。安装并启动 EMQX根据当前操作系统到 EMQX 官网下载安装包并执行安装。安装后使用sudo emqx start或sudo systemctl start emqx启动 EMQX 服务。注意本文基于 EMQX v4.4.5 版本其他版本的规则引擎配置界面、配置方法以及功能可能随版本升级有所区别请以对应版本的官方文档为准。在 TDengine 中创建数据库和表进入taosshell复制并执行以下 SQL为接收 MQTT 数据创建数据库与表结构CREATE DATABASE test; USE test; CREATE TABLE sensor_data (ts TIMESTAMP, temperature FLOAT, humidity FLOAT, volume FLOAT, pm10 FLOAT, pm25 FLOAT, so2 FLOAT, no2 FLOAT, co FLOAT, sensor_id NCHAR(255), area TINYINT, coll_time TIMESTAMP);各字段与后续 MQTT 报文、规则模板的对应关系是ts取服务端的nowtemperature/humidity/volume/pm10/pm25/so2/no2/co映射 JSON 数值字段sensor_id映射设备 IDarea映射区域编号coll_time映射报文里设备自己上报的时间戳ts注意这与表主键时间戳ts不是同一个东西命名上需仔细区分。注表结构以 TDengine 官方博客“EMQX TDengine 搭建 MQTT 物联网数据可视化平台”场景为例后续操作均以此场景进行请你根据实际业务场景修改字段。配置 EMQX 规则引擎以下操作以 EMQX v4.4.5 的 Dashboard 界面为例。登录 EMQX Dashboard使用浏览器打开http://IP:18083登录 EMQX Dashboard。初次安装的用户名为admin密码为public。创建规则Rule选择左侧“规则引擎Rule Engine”中的“规则Rule”并点击“创建Create”按钮。编辑 SQL 字段复制以下内容输入到 SQL 编辑框SELECT payload FROM sensor/data其中payload代表整个消息体sensor/data是本规则选取的消息主题。后续动作模板中通过${payload.xxx}引用消息体字段。新增“动作Action”在规则详情页进入“动作”配置为规则绑定一个 Action handler。新增“资源Resource”选择“发送数据到 Web 服务”并点击“新建资源”按钮。编辑“资源Resource”选择“WebHook”并填写“请求 URL”为 taosAdapter 提供 REST 服务的地址。如果是本地启动的 taosadapter默认地址为http://127.0.0.1:6041/rest/sql其他属性请保持默认值。如果 taosAdapter 部署在远程节点请替换为对应主机的 FQDN/IP 和httpPort。编辑“动作Action”编辑动作配置增加Authorization认证的键/值配对项。默认用户名root和密码taosdata对应的 Basic 认证值为root:taosdata经 Base64 编码后的结果Basic cm9vdDp0YW9zZGF0YQ根据 TDengine REST API 文档请求 Header 里需带有身份认证信息若使用了自定义账号请按{username}:{password}重新做 Base64 编码后填入。在消息体中输入规则引擎替换模板INSERT INTO test.sensor_data VALUES( now, ${payload.temperature}, ${payload.humidity}, ${payload.volume}, ${payload.PM10}, ${payload.pm25}, ${payload.SO2}, ${payload.NO2}, ${payload.CO}, ${payload.id}, ${payload.area}, ${payload.ts} )最后点击左下方的“Create”按钮保存规则。模板编写要点EMQX 的${payload.xxx}替换是按 JSON 键名严格区分大小写的。上面模板中PM10、SO2、NO2、CO为大写而temperature、humidity、volume、pm25、id、area、ts为小写这必须与 MQTT 报文中的实际键名一一对应键名不匹配时变量不会被替换最终 SQL 解析失败sensor_id列是NCHAR类型因此${payload.id}需要用单引号包裹成字符串字面量数值列则直接裸值替换表名必须写全test.sensor_data的库前缀见上文 REST 接口约束若希望写入的时间戳由 taosd 执行 SQL 时的当前时间决定用now即可报文中携带的ts字段则落到coll_time列可用于核对上报时间。编写模拟测试程序仓库中提供了完整的模拟客户端程序 mock.js其核心逻辑是创建 N 个 MQTT 虚拟客户端从 24 小时前到现在按 5 秒步长回补数据每个时间点为所有客户端生成一条随机数据并发布到sensor/data主题。关键参数与完整代码摘录如下// mock.js const mqtt require(mqtt) const Mock require(mockjs) const EMQX_SERVER mqtt://localhost:1883 // EMQX MQTT 接入地址 const CLIENT_NUM 10 // 虚拟客户端数量 const STEP 5000 // 数据点间隔毫秒对应“last 24h every 5s”回补步长 const AWAIT 5000 // 每轮写入后的休眠时间避免写入过快压垮硬件 const CLIENT_POOL [] startMock() async function startMock() { const now Date.now() for (let i 0; i CLIENT_NUM; i) { const client await createClient(mock_client_${i}) CLIENT_POOL.push(client) } // last 24h every 5s const last 24 * 3600 * 1000 for (let ts now - last; ts now; ts STEP) { for (const client of CLIENT_POOL) { const mockData generateMockData() const data { ...mockData, id: client.clientId, // 设备 ID对应 sensor_id 列 area: 0, ts, // 客户端时间戳对应 coll_time 列 } client.publish(sensor/data, JSON.stringify(data)) } const dateStr new Date(ts).toLocaleTimeString() console.log(${dateStr} send success.) await sleep(AWAIT) } console.log(Done, use ${(Date.now() - now) / 1000}s) } function createClient(clientId) { return new Promise((resolve, reject) { const client mqtt.connect(EMQX_SERVER, { clientId }) client.on(connect, () { console.log(client ${clientId} connected) resolve(client) }) client.on(error, (e) { console.error(e) reject(e) }) }) } function generateMockData() { return { temperature: parseFloat(Mock.Random.float(22, 100).toFixed(2)), humidity: parseFloat(Mock.Random.float(12, 86).toFixed(2)), volume: parseFloat(Mock.Random.float(20, 200).toFixed(2)), PM10: parseFloat(Mock.Random.float(0, 300).toFixed(2)), pm25: parseFloat(Mock.Random.float(0, 300).toFixed(2)), SO2: parseFloat(Mock.Random.float(0, 50).toFixed(2)), NO2: parseFloat(Mock.Random.float(0, 50).toFixed(2)), CO: parseFloat(Mock.Random.float(0, 50).toFixed(2)), area: Mock.Random.integer(0, 20), ts: 1596157444170, } }注意两点CLIENT_NUM在开始测试时先设置一个较小的值避免硬件性能不能完全处理较大并发客户端数量generateMockData输出的键名含PM10/SO2/NO2/CO的大小写必须与 EMQX 规则模板中的${payload.xxx}完全一致否则模板无法渲染出合法 SQL。执行测试模拟发送 MQTT 数据在存放mock.js的目录下执行npm install mqtt mockjs --save --registryhttps://registry.npm.taobao.org node mock.js程序会先建立CLIENT_NUM个 MQTT 连接随后每 5 秒为一个时间点批量发布数据直到回补完最近 24 小时的数据点并打印Done, use xxxs。验证数据写入结果在 EMQX Dashboard 查看规则命中统计在 EMQX Dashboard 规则引擎界面刷新可以看到“已匹配消息”等计数确认规则被触发、有多少条记录被正确接收和处理。在 taos shell 中查询落库数据使用taosshell 程序登录查询相应数据库和表验证数据是否被正确写入USE test; SELECT COUNT(*) FROM sensor_data; SELECT * FROM sensor_data LIMIT 10;若能查到与报文对应的temperature、humidity、sensor_id即mock_client_0等客户端名等字段值说明整条链路已打通。排障要点速查现象可能原因与排查依据EMQX 动作发送后 TDengine 无数据taosAdapter 未启动或httpPort非默认 6041URL 填写错误taosAdapter 参考手册、REST API 文档请求返回 401Authorization头缺失或 Base64 值错误默认root:taosdata对应Basic cm9vdDp0YW9zZGF0YQREST API 文档“身份认证”一节返回错误而非 200/落库失败默认httpCodeServerError为 false 时 C 接口出错也返回 200需看 body 中的错误码模板变量未替换、表名缺库前缀都会导致 SQL 解析错误REST API 文档“HTTP 响应码”一节数据缺字段/为 NULL模板${payload.xxx}键名与 JSON 实际键名大小写不一致或 JSON 字段缺失对照 mock.js 的键名与规则模板版本升级后界面不一致本文操作基于 EMQX v4.4.5升级后规则引擎界面与配置方式可能有差异EMQX 官方对应版本文档TDengine 与 taosAdapter 的详细使用方法可分别参考 TDengine 官方文档 中的开发指南与运维手册章节。【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
02
RELATED NEWS

相关资讯

更多网站建设与数字化升级内容

03
WHY YAOTU

想打造同款高转化官网?

懂行业、懂生意,从建站到增长一站式陪跑

场景化定制

不做模板站,围绕你的业务场景量身设计,小众不撞款。

营销型架构

以转化目标组织内容与路径,让官网真正带来询盘。

全周期服务

设计、开发、运营、运维一体,上线只是开始。

免费获取你的建站方案

留下需求,专属顾问 24 小时内为你输出方案建议。