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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 05:33:57