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

如何在Node.js中借助MQTT消息中间件实现数据双库同步存储

实现MQTT消息同时存入两个数据库的可行方案

以下提供两种低耦合、易维护的实现方案,适配你当前基于Mosca的MQTT架构:


方案一:Broker端直接处理双库写入

适合需要在消息到达Broker时立即完成双库存储的场景,直接在Broker逻辑中添加双数据库连接与并行写入操作。

修改后的Broker代码

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

// 第一个数据库连接(原有MySQL)
var mysql = require('mysql');
var db1 = mysql.createConnection({
      host: 'localhost',
      user: 'root',
      password: 'root@123',
      database: 'abc'
});

// 第二个数据库连接(示例为同类型数据库,可替换为MongoDB等其他类型)
var db2 = mysql.createConnection({
      host: 'localhost',
      user: 'root',
      password: 'root@123',
      database: 'def' // 目标数据库名
});

// 初始化数据库连接
db1.connect(() => console.log("db1 connected"));
db2.connect(() => console.log("db2 connected"));

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

// 封装数据库插入逻辑,避免重复代码
function insertToDB(db, data) {
    const sql = `INSERT INTO mqttjs (message, time, add) VALUES (?, ?, ?)`;
    return new Promise((resolve, reject) => {
        db.query(sql, [data.message, data.time, data.add], (err, result) => {
            err ? reject(err) : resolve(result);
        });
    });
}

broker.on('published', async (packet) => { 
    const payload = packet.payload.toString();
    console.log("Received message:", payload);

    // 过滤系统消息与不符合格式的内容
    if (payload.slice(0, 1) !== '{' && payload.slice(0, 4) !== 'mqtt') {
        // 构造存储数据(若消息为JSON格式,可替换为JSON.parse(payload)解析)
        const saveData = {
            message: payload,
            time: new Date().toISOString(),
            add: "自定义字段内容"
        };

        try {
            // 并行写入两个数据库,提升处理效率
            await Promise.all([
                insertToDB(db1, saveData),
                insertToDB(db2, saveData)
            ]);
            console.log("数据已成功存入两个数据库");
        } catch (err) {
            console.error("存储失败:", err);
        }
    }
});

关键优化点

  • 新增第二个数据库连接配置,支持跨类型数据库(只需替换对应驱动与连接逻辑)
  • 用Promise+async/await处理异步写入,避免回调嵌套
  • 通过Promise.all并行执行双库写入,提升处理效率
  • 修复原代码中message/time/add字段值重复的问题

方案二:独立订阅者解耦存储逻辑

符合微服务设计思想,Broker仅负责消息转发,存储逻辑拆分到独立的订阅服务中,降低Broker耦合度。

步骤1:精简Broker代码(仅做消息转发)

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

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

// 仅保留消息日志,移除数据库相关逻辑
broker.on('published', (packet) => { 
    console.log("Forward message:", packet.payload.toString());
});

步骤2:编写第一个数据库订阅服务(存abc库)

var mqtt = require('mqtt')
var mysql = require('mysql');

var client = mqtt.connect('mqtt://192.168.0.92')
var db = mysql.createConnection({
      host: 'localhost',
      user: 'root',
      password: 'root@123',
      database: 'abc'
});

db.connect(() => console.log("db1 connected"));

client.on('connect', () => client.subscribe('myTopic'));

client.on('message', (topic, message) => {
    const payload = message.toString();
    console.log("Saving to db1:", payload);

    if (payload.slice(0, 1) !== '{' && payload.slice(0, 4) !== 'mqtt') {
        const sql = `INSERT INTO mqttjs (message, time, add) VALUES (?, ?, ?)`;
        db.query(sql, [payload, new Date().toISOString(), "自定义字段"], (err) => {
            err ? console.error("db1 save error:", err) : console.log("data saved to db1");
        });
    }
});

步骤3:编写第二个数据库订阅服务(存def库)

var mqtt = require('mqtt')
var mysql = require('mysql');

var client = mqtt.connect('mqtt://192.168.0.92')
var db = mysql.createConnection({
      host: 'localhost',
      user: 'root',
      password: 'root@123',
      database: 'def'
});

db.connect(() => console.log("db2 connected"));

client.on('connect', () => client.subscribe('myTopic'));

client.on('message', (topic, message) => {
    const payload = message.toString();
    console.log("Saving to db2:", payload);

    if (payload.slice(0, 1) !== '{' && payload.slice(0, 4) !== 'mqtt') {
        const sql = `INSERT INTO mqttjs (message, time, add) VALUES (?, ?, ?)`;
        db.query(sql, [payload, new Date().toISOString(), "自定义字段"], (err) => {
            err ? console.error("db2 save error:", err) : console.log("data saved to db2");
        });
    }
});

优势

  • Broker职责单一,仅做消息转发,便于后续扩展
  • 存储服务独立部署,单个服务故障不影响整体消息链路
  • 新增数据库存储只需添加对应订阅服务,无需修改Broker代码

内容的提问来源于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.03 17:51:30