You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用Node.js通过MQTT Broker发送多数据并存储到MySQL多字段

解决方案:MQTT多组数据存入MySQL不同字段

问题分析

当前代码存在两个核心问题:

  • 发布端拆分发送多主题消息,Broker每次仅能收到单条数据,无法关联为同一条数据库记录的不同字段
  • Broker端未区分消息主题,每次插入时将三个字段都赋值为当前消息内容,导致数据重复且无法对应到正确字段

方法一:单主题发送JSON格式组合数据(推荐)

将所有需要存储的字段打包成JSON字符串,通过单个主题发送,Broker收到后解析JSON并直接插入数据库对应字段。

修改后的发布端代码

var mqtt = require('mqtt');
var client = mqtt.connect('mqtt://192.168.0.92');

client.on('connect', function () {
    setInterval(function () {
        // 组合所有数据为JSON对象
        var data = {
            message: "message content",
            time: new Date().toISOString(), // 示例:使用当前时间
            address: "192.168.0.100"
        };
        // 将JSON转为字符串发送
        client.publish('myTopic', JSON.stringify(data));
        console.log('Message Sent:', data);
    }, 5000);
});

修改后的Broker端代码

var mosca = require('mosca');
var settings = { port: 1883}
var broker = new mosca.Server(settings);

var mysql = require('mysql');
var db = mysql.createConnection({
    host: 'localhost',
    user: 'root',
    password: 'root@123',
    database: 'abc'
});

db.connect(() => { 
    console.log("db connect");
});

broker.on('ready', () => { 
    console.log("broker is ready");
});

broker.on('published', (packet) => { 
    try {
        // 过滤系统消息
        const payloadStr = packet.payload.toString();
        if (payloadStr.startsWith('{') && packet.topic === 'myTopic') {
            // 解析JSON数据
            const data = JSON.parse(payloadStr);
            // 插入数据库,对应不同字段
            var dbSet = 'INSERT INTO broker SET ?';
            db.query(dbSet, data, (error, output) => {
                if (error) {
                    console.log('数据库插入错误:', error);
                } else {
                    console.log("数据已保存:", data);
                }
            });
        }
    } catch (err) {
        console.log('消息解析错误:', err);
    }
});

方法二:多主题缓存数据,凑齐后插入

若必须使用不同主题发送数据,Broker端需缓存各主题的消息内容,当一组数据全部收集完成后,再插入数据库。

修改后的发布端代码(调整主题名更清晰)

var mqtt = require('mqtt');
var client = mqtt.connect('mqtt://192.168.0.92');

client.on('connect', function () {
    setInterval(function () {
        var a = "message content";
        var b = new Date().toISOString();
        var c = "192.168.0.100";
        // 使用区分度更高的主题
        client.publish('data/message', a);
        client.publish('data/time', b);
        client.publish('data/address', c);
        console.log('Messages Sent');
    }, 5000);
});

修改后的Broker端代码

var mosca = require('mosca');
var settings = { port: 1883}
var broker = new mosca.Server(settings);

var mysql = require('mysql');
var db = mysql.createConnection({
    host: 'localhost',
    user: 'root',
    password: 'root@123',
    database: 'abc'
});

db.connect(() => { 
    console.log("db connect");
});

broker.on('ready', () => { 
    console.log("broker is ready");
});

// 缓存数据对象,用于暂存一组数据
var dataCache = {};

broker.on('published', (packet) => { 
    const payload = packet.payload.toString();
    // 过滤系统消息
    if (payload.startsWith('{') || payload.startsWith('mqtt')) return;

    // 根据主题分配缓存字段
    switch(packet.topic) {
        case 'data/message':
            dataCache.message = payload;
            break;
        case 'data/time':
            dataCache.time = payload;
            break;
        case 'data/address':
            dataCache.address = payload;
            break;
        default:
            return; // 忽略未知主题
    }

    // 检查是否凑齐所有字段
    if (dataCache.message && dataCache.time && dataCache.address) {
        var dbSet = 'INSERT INTO broker SET ?';
        db.query(dbSet, dataCache, (error, output) => {
            if (error) {
                console.log('数据库插入错误:', error);
            } else {
                console.log("数据已保存:", dataCache);
            }
        });
        // 清空缓存,准备下一组数据
        dataCache = {};
    }
});

内容的提问来源于stack exchange,提问作者Prakash Kumar Gupta

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.02 22:32:21