【后端开发实习】用MongoDB和Redis实现消息队列搭建分布式邮件消息系统
创始人
2025-01-11 09:06:22
0

用Redis实现消息队列并搭建分布式邮件消息系统

  • 系统介绍
  • Redis实现消息队列
    • 思路分析
    • 代码实现
  • MongoDB监听数据变化
    • 思路分析
    • 代码实现
      • Mongoose测试连接
      • 监听mongodb数据变化
  • 注意点

系统介绍

本次要实现的是一个能够实现实时监控Mongodb中数据变化的系统,要能够在数据发生变动的时候实时将变动消息发送给指定的邮箱。

  • Node.js:用于开发的语言,既能用于前端开发,又能用来做后端开发。
  • Redis:用于搭建消息队列,实现消息的分布式。
  • MongoDB:持久化数据,同时实现触发条件的监听,当MongoDB中有新增数据的时候发送新增数据的邮件消息。

Redis实现消息队列

思路分析

主要使用的就是Redis-smq这个库,下面展示的就是主要使用的消息队列类,其中包括了很多队列种类,有先进先出、优先级先出等方式。
在这里插入图片描述
整个库的原理如下结构图,本次使用到的只有主线,就是发送和接收:
在这里插入图片描述

代码实现

const { transemail } = require('../email_list/email.js'); const redis = require('promise-redis-client'); const redisHost = 'localhost'; const redisPort = 6379;  // 配置 Redis 客户端 const createRedisClient = () => {     return new Promise((resolve, reject) => {         let client = redis.createClient({ host: redisHost, port: redisPort });         client.on('error', err => {             console.log('Redis 连接出错');             reject(err);         });         client.on('ready', () => {             console.log('Redis ready');             resolve(client);         });     }); };  async function startWaitMsg(redisClient) {     while (true) {         let res = null;         try {             res = await redisClient.brpop('bookChanges', 0);             console.log('收到消息', res);         } catch (err) {             console.log('brpop 出错,重新 brpop');             continue;         }         res = res.toString();         transemail(res);     } }  async function listenredis() {     try {         // 启动生产者         // startProducer();          // 创建 Redis 客户端         const redisClient = await createRedisClient();          // 启动消息监听         startWaitMsg(redisClient);     } catch (error) {         console.error('Error:', error);     } } //测试的时候使用的代码 listenredis().catch(console.error);  // 处理退出信号以关闭客户端 process.on('SIGINT', async () => {     console.log('Closing clients...');     process.exit(0); });  

MongoDB监听数据变化

思路分析

由于要实现实时检测,经过分析以后使用mongoose中的数据流监控最为合适,但是要实现这个方法需要用到watch方法,这个方法只有在mongodb有副本集的时候才能使用,因此还需要提前配置好mongodb才能进行这里下一步的操作,如果没有配置过mongodb的副本集的可以参考我的这篇博客。

  1. 用mongoose中的watch连接mongodb副本集数据库获取数据变化。
  2. 将数据变化发送到redis消息队列中。

首先在命令行中将服务启动:
在这里插入图片描述

代码实现

Mongoose测试连接

const mongoose = require('mongoose');  mongoose.connect('mongodb://localhost/test', {   useNewUrlParser: true,   useUnifiedTopology: true }).then(() => {   console.log('Successfully connected to MongoDB');    const bookSchema = new mongoose.Schema({     title: String,     author: String   });    const Book = mongoose.model('Book', bookSchema);    const bookChangeStream = Book.watch();    bookChangeStream.on('change', (change) => {     console.log('Collection changed:', change);     if (change.operationType === 'insert') {       console.log('New book added:', change.fullDocument);     }   }); }).catch((error) => {   console.log('Error connecting to MongoDB:', error); });  

在这里插入图片描述
测试结果:
在Mongo Campass中添加数据以后,在终端中出现如下消息:
在这里插入图片描述
证明测试成功,可以进行下一步操作啦!

监听mongodb数据变化

const redis = require('redis'); const mongoose = require('mongoose'); // 创建 Redis 客户端 const redisClient = redis.createClient({ 	host: 'localhost', 	port: 6379   });      // 连接到 Redis redisClient.connect();    //连接mongodb数据库并检测变化发送到redis消息队列 async function connectAndMonitorMongoDB(redisClient) {   try {     await mongoose.connect('mongodb://localhost/test', {       useNewUrlParser: true,       useUnifiedTopology: true     });     console.log('Successfully connected to MongoDB');      const bookSchema = new mongoose.Schema({       title: String,       author: String     });      const Book = mongoose.model('Book', bookSchema);      const bookChangeStream = Book.watch(); 	try{ 		bookChangeStream.on('change', (change) => { 			console.log('Collection changed:', change); 			console.log("type of change:",typeof(change)); 			msg = JSON.stringify(change.fullDocument); 			msg = msg.replace(/{|}/g, ''); 			msg = "New message received:"+msg; 			console.log("massage:",msg); 			console.log("type of message:",typeof(msg)); 			if (change.operationType === 'insert') { 			  console.log('New book added:', msg); 			  redisClient.lPush('bookChanges', msg, function(err, reply) { 				if (err) { 				  console.log('Error storing JSON to Redis:', err); 				} else { 				  console.log('JSON stored successfully, list length:', reply); 				}}) 			} 		  }); 	}catch (err){ 		console.log("error while loading data into redis:", err) 	}   } catch (error) {     console.log('Error connecting to MongoDB:', error);   } }  // module.exports = { connectAndMonitorMongoDB }; async function main() {   try {     await connectAndMonitorMongoDB(redisClient);     console.log('Monitoring MongoDB changes...');   } catch (error) {     console.error('Failed to start monitoring:', error);   } }  main(); 

注意点

在nodejs中将JSON对象转换成字符串的JSON.Stringify函数并不是严格的转换成字符串而是带有一个大括号,然而这个在进行redis进队列的时候会有问题,因此需要用正则表达式去掉大括号:

msg = JSON.stringify(change.fullDocument); msg = msg.replace(/{|}/g, ''); msg = "New message received:"+msg; 

相关内容

热门资讯

总算了解!广西老友有破解吗,开... 总算了解!广西老友有破解吗,开心泉州免费辅助器,攻略教程(有挂分享)1、开心泉州免费辅助器有没有辅助...
无独有偶!微信小程序免费黑科技... 无独有偶!微信小程序免费黑科技透视,微乐小程序多功能修改器法子教程(有挂分享)1、这是跨平台的微乐小...
避坑细节!宝宝游戏辅助,决战卡... 避坑细节!宝宝游戏辅助,决战卡五星游戏辅助器,积累教程(有挂规律)亲,关键说明,决战卡五星游戏辅助器...
据权威媒体报道!微信小程序免费... 据权威媒体报道!微信小程序免费黑科技透视,钱塘十三水其实真的有挂手段教程(有挂详细)1、钱塘十三水其...
实测揭晓!同乡游辅助工具制作,... 实测揭晓!同乡游辅助工具制作,永胜联盟会封号吗,模块教程(有挂教学)1、同乡游辅助工具制作模拟器是什...
实测必看!微信小程序免费黑科技... 实测必看!微信小程序免费黑科技开挂,皮皮胡子辅助举措教程(有挂详细)1、起透看视 微信小程序免费黑科...
一分钟揭秘!微乐小程序免费黑科... 一分钟揭秘!微乐小程序免费黑科技安卓,新超凡软件辅助练习教程(有挂存在)1、微乐小程序免费黑科技脚本...
必备辅助推荐!蜜瓜大厅可以装挂... 必备辅助推荐!蜜瓜大厅可以装挂吗,友友联盟辅助脚本,指南书教程(有挂技术)1、进入游戏-大厅左侧-新...
今日焦点!微信小程序免费黑科技... 今日焦点!微信小程序免费黑科技透视,川南久久辅助模块教程(存在有挂)运微信小程序免费黑科技辅助工具,...
技术分享!山西扣点点辅助工具,... 技术分享!山西扣点点辅助工具,福州十八扑外卦视频,诀窍教程(有挂方法)1、山西扣点点辅助工具公共底牌...