如何在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
相关产品推荐
相关产品推荐

