【问题标题】:Set a new socket io connection per user为每个用户设置一个新的套接字 io 连接
【发布时间】:2020-09-02 20:34:05
【问题描述】:

我在 MERN 堆栈应用程序中使用套接字 io 在客户端和服务器之间进行通信。基本上,我正在监听我的 mongodb 集群中的更改,并使用 socketio emit 函数将这些更改传达给客户端。问题是传回给客户端的数据是特定于某些用户的,但是 emit 函数会将这些详细信息“发送”给应用程序上的所有用户。似乎所有用户都在同一个套接字客户端连接上,或者至少被服务器识别为同一个客户端连接。我如何确保这种与 socketio 的服务器间客户端通信是针对每个用户的。

express&socket io服务器:

const app = express();

app.use(cors({
  origin: CLIENT_ORIGIN
}));

//body parser middleware
app.use(express.json());
app.use(cookieParser());

// app.use(express.urlencoded({extended: true}));
app.use( bodyParser.urlencoded({extended: true}));

// Passport Config
require('./config/passport')(passport);

//Database config
const db = process.env.MONGODB_URI;
// console.log(db);

//connect to MongoDB database
mongoose
.connect(db, {
    useNewUrlParser: true,
    useFindAndModify: false,
    useUnifiedTopology: true
    // useCreateIndex: true
})
.then(() => {
    console.log("Connected to MongoDB...");

    // change streams go here
    // const changeStream = Users.collection.watch();
    // changeStream.on('change', next => {
    //     console.log(next);
    // })
})
.catch(err => {
    console.log(err);
});

// Express session
app.use(
    session({
      secret: process.env.EXPRESS_SESSION_SECRET+'',
      resave: true,
      saveUninitialized: true
    })
);

// Passport middleware
app.use(passport.initialize());
app.use(passport.session());


//Server routes
//add server route that allows "home landing path"
app.use('/api/posts', require('./routes/api/posts'));
app.use('/api/tags', require('./routes/api/tags'));
app.use('/api/comments', require('./routes/api/comments'));
app.use('/api/users', require('./routes/api/users'));
app.use('/api/reactions', require('./routes/api/reactions'));
app.use('/', require('./routes/landing'));



//Serve static folder if in production

if (process.env.NODE_ENV === 'production') {
    // Set static folder
    app.use(express.static('client/build'));
    app.enable('trust proxy');

    // Catch all to handle all other requests that come into the app. 
    app.get('*', (req, res) => {
      res.sendFile(path.resolve(__dirname, 'client', 'build', 'index.html'));
    });
  }





const server = app.listen(PORT, () => {
    console.log("Server is running on port: ", PORT)
});


const io = socket.listen(server);


io.origins(CLIENT_ORIGIN);

io.on('connection', (socket) => {
  console.log("server socket connection on...");

  socket.connected ? console.log('client server linked') : socket.emit('restartListen');

  //listen for changes/inserts/deletes in all posts
  socket.on('listenForPosts', (schoolId) => {
    console.log('listening for posts...');

    const pipeline = {
      $match: {
        'fullDocument.institution': schoolId
      }
    }

    const changeStream = Post.collection.watch(pipeline);

    changeStream.on('change', next => {
      switch(next.operationType) {
        case 'insert':
          console.log('an insert happened...');

          socket.emit('insert', {
            type: 'insert',
            msg: 'New question available',
            newPost: next.fullDocument
          });
          break;

        case 'update':
          console.log('an update happened...');

          socket.emit('update', {
            type: 'update',
            postId: next.documentKey._id,
            updateInfo: next.updateDescription.updatedFields,
            msg: "Question has been updated."
          });
          break;

        case 'delete':
          console.log('a delete happened...');

          socket.emit('delete', {
            type: 'delete',
            deletedId: next.documentKey._id,
            msg: 'Question has been deleted.'
          });
          break;

        default:
          break;
      }
    });

  });

  socket.on('stopListenForPosts', () => {
    changeStream.close();
    setTimeout(() => {
      socket.disconnect(true);
    }, 3000)
  });

})

clint 端:

import io from 'socket.io-client';
const socket = io();

socket.connected ? console.log('Client Socket connected') : socket.connect();

socket.on('insert', (data) => {
     this.props.actions.postActions.doSocketInsert(data);
});

socket.on('update', (data) => {
     // console.log(data);
    this.props.actions.postActions.doSocketUpdate(data);
});

socket.on('delete', (data) => {
     this.props.actions.postActions.doSocketDelete(data);
});

socket.on('connect', () => {
     this.listen(socket);
});

socket.on('restartListen', () => {
    socket.connect();
    this.listen(socket);
});

我在想也许让每个客户端加入一个房间并将更改发送到那个特定的套接字房间?

【问题讨论】:

  • Socket.io 自动将每个用户连接到一个房间,并在连接时使用他们的 id。你可以向那个房间发送一条消息,它只会发送给一个用户。

标签: node.js express sockets socket.io mern


【解决方案1】:

为了解决此问题,您需要使用套接字保存有关已连接用户的数据。为了执行此操作,您只需在内存中创建一个数组并将用户信息保存在其中。 (连接用户的user_id和socket_id)。

这里是你如何存储它。

let users = [];

socket.on('insert', data => {
  let found = false;
  users.forEach(element => {
    if (element.user._id === data.user._id) {
      found = true;
    }
  });
  if (!found) {
    users.push({ socket_id: socket.id, user: data.user });
  }
  console.log('-------Logged In Users--------');
  console.log(users);

  // do your other stuff here
});

然后,您需要解析 user_id 以及您需要与特定用户广播的每条消息。然后,在listen 函数中,您可以检查我们在内存中创建的用户数组,并找到与特定用户相关的特定 socket_id,您可以只为该用户发出如下喜欢的消息。

for (let i = 0; i < users.length; i++) {
    if(user._id === data.user._id){
        io.to(users[i].socket_id).emit('chat_message', obj);
    }
}

【讨论】:

  • 循环访问所有连接的用户似乎不是最有效的想法;必须有另一种方式。也许一个 socket.io 函数本身就可以按需要工作?
【解决方案2】:

你这样说的:

changeStream.on('change', ...)

内部:

io.on('connection', ...)

这意味着每个套接字都在侦听changeStream.on('change', ...),因此changeStream 上发生的任何事情都将发送到每个套接字。更糟糕的是,它甚至会尝试将其发送到所有先前连接的不再连接的套接字(这会失败,但仍然是错误的做法)。

所以,我不知道你究竟想在这里做什么,但你通常不希望为每个新连接的套接字启动一个单独的外部资源侦听器。

【讨论】:

  • 我基本上是在尝试为每个用户创建一个通知系统。因此,当 changeStream 获取数据库中的更改时,这些文档被标记为更改所关注的用户,我将使用 socket.io 将这些更改传达给前端。但出于某种原因,所有登录用户都会收到所有更改的通知。
  • 有没有办法在 io.on() 中没有 changeStream 的情况下实现我在上面的评论中写的内容?谢谢。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2013-04-08
  • 2014-01-30
  • 2011-05-04
  • 2014-09-30
  • 1970-01-01
  • 2019-04-04
  • 1970-01-01
相关资源
最近更新 更多