如何在NestJS中实现MongoDB Change Stream?
在NestJS中实现MongoDB Change Stream的方案
你之前用的Mongoose post('save')钩子仅对当前应用内的Mongoose操作生效,外部应用直接修改MongoDB数据时不会触发,而MongoDB原生的Change Stream可以监听集合的所有变更操作,不管修改来自哪里。下面是在NestJS中的具体实现步骤:
1. 核心实现思路
通过NestJS的Mongoose连接获取原生MongoDB数据库实例,创建Change Stream并监听集合的变更事件,在事件回调中处理业务逻辑。
2. 代码实现
第一步:创建Change Stream服务
import { Injectable, OnModuleInit, OnModuleDestroy } from '@nestjs/common'; import { InjectConnection } from '@nestjs/mongoose'; import { Connection } from 'mongoose'; import { ChangeStream } from 'mongodb'; @Injectable() export class ChangeStreamService implements OnModuleInit, OnModuleDestroy { private changeStream: ChangeStream | null = null; constructor(@InjectConnection() private readonly connection: Connection) {} async onModuleInit() { // 获取原生MongoDB的db实例 const db = this.connection.db; // 监听指定集合(这里以Cat集合为例) this.changeStream = db.collection('cats').watch(); // 监听变更事件 this.changeStream.on('change', async (change) => { // 根据操作类型处理不同的变更 switch (change.operationType) { case 'insert': console.log('新增文档:', change.fullDocument); // 这里可以添加你的业务逻辑,比如发送消息、更新缓存等 break; case 'update': console.log('更新文档ID:', change.documentKey._id); console.log('更新内容:', change.updateDescription); break; case 'delete': console.log('删除文档ID:', change.documentKey._id); break; default: console.log('其他操作:', change.operationType); } }); // 监听错误 this.changeStream.on('error', (error) => { console.error('Change Stream出错:', error); }); } async onModuleDestroy() { // 应用关闭时关闭Change Stream if (this.changeStream) { await this.changeStream.close(); } } }
第二步:在模块中注册服务并配置Mongoose连接
import { Module } from '@nestjs/common'; import { MongooseModule } from '@nestjs/mongoose'; import { ConfigModule, ConfigService } from '@nestjs/config'; import { CatsSchema, Cat } from './schemas/cat.schema'; import { ChangeStreamService } from './change-stream.service'; @Module({ imports: [ ConfigModule.forRoot(), MongooseModule.forRootAsync({ imports: [ConfigModule], useFactory: (configService: ConfigService) => ({ uri: configService.get('MONGODB_URI'), }), inject: [ConfigService], }), MongooseModule.forFeature([{ name: Cat.name, schema: CatsSchema }]), ], providers: [ChangeStreamService], }) export class AppModule {}
3. 关键注意事项
- MongoDB Change Stream要求集群为副本集或分片集群,本地开发时可以启动单节点副本集:
# 启动单节点副本集 mongod --replSet rs0 # 进入MongoDB Shell初始化副本集 mongosh rs.initiate() - 可通过
watch()方法传入过滤条件,只监听特定变更,比如仅监听name字段的更新:this.changeStream = db.collection('cats').watch([ { $match: { 'updateDescription.updatedFields.name': { $exists: true }, }, }, ]);
内容的提问来源于stack exchange,提问作者thavorac
相关产品推荐
相关产品推荐

