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

MongoDB 4.4.6下Mongoose watch()未触发insert事件求助

MongoDB 4.4.6中Data.watch()无法触发change事件的问题

当前使用MongoDB 4.4.6,调用Data.watch()方法时无法触发事件。已确认changeStream的连接与创建均正常,数据流也已正确生成。尝试通过API和MongoDB客户端两种方式添加数据,但changeStream.on('change')始终未触发。

相关代码如下:

const express = require('express');
const app = express();
const http = require('http').Server(app);
const io = require('socket.io')(http);
const mongoose = require('mongoose');
app.use(express.json())

// Connect to MongoDB using Mongoose
mongoose.connect('mongodb://localhost:27017/socketdb');


// Define a schema for the data
const dataSchema = new mongoose.Schema({
  name: String,
  age: Number,
  timestamp: { type: Date, default: Date.now }
});

// Create a model for the data
const Data = mongoose.model('Data', dataSchema);

// Connect to Socket.io
io.on('connection', function (socket) {
  console.log('A user connected');
  Data.find({}).then(
    data => socket.emit('allData', data)
  ).catch(err => console.log(err))


  // Watch for changes in the MongoDB collection
  if (mongoose.connection.readyState !== 1) {
    console.log('Mongoose is not connected');
  }
  let pipeline = [
    {
      $match: {
        operationType: 'insert'
      }
    }
  ];
  
  let changeStream = Data.watch(pipeline)
  .on('change', data => console.log(data))

  console.log(changeStream)
  changeStream.on('change', function (change) {
    console.log('Change stream event:', change);

    // Send the new or updated data to the client
    if (change.operationType === 'insert') {
      socket.emit('newData', change.fullDocument);
      console.log(change.fullDocument)
    } else if (change.operationType === 'update') {
      Data.findById(change.documentKey._id, function (err, data) {
        if (err) throw err;
        socket.emit('newData', data);
      });
    }
  });

  socket.on('disconnect', function () {
    console.log('A user disconnected');
    // changeStream.close();
  });
});

// Serve the HTML page
app.get('/', function (req, res) {
  res.sendFile(__dirname + '/index.html');
});

app.post('/adddata', (req, res) => {
  let d = new Data({
    name: req.body.name,
    age: req.body.age
  })
  d.save().then(() => {
    res.send("Data Added")
  }).catch((err) => {
    res.send("ERR")
  })
  Data.watch(pipeline).on('change', data => console.log('inside the add',data))

})

// Start the server
http.listen(3000, function () {
  console.log('Listening on port 3000');
});

已查阅Mongoose官方文档,仍无法解决问题,寻求技术支持。

内容的提问来源于stack exchange,提问作者Prashant Dey

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 20:55:18