【问题标题】:How to process websocket messages sequentially如何按顺序处理 websocket 消息
【发布时间】:2021-04-21 14:25:39
【问题描述】:

每个 WebSocket 接收到数十条消息,它们的到达时间可能相差几毫秒。我需要使用有时需要一些时间的操作来处理这些数据(例如插入数据库)。 为了处理收到的新消息,必须先处理完之前的消息。

我的第一个想法是用 node.js Bull ( 用 Redis ) 准备一个队列,但恐怕运行时间太长。这些消息的处理必须保持快速。

我尝试使用 JS 迭代器/生成器(直到现在我从未使用过的东西)并且我测试了这样的东西:

const ws = new WebSocket(`${this.baseUrl}${this.path}`)
const duplex = WebSocket.createWebSocketStream(ws, { encoding: 'utf8' })
    
const messageGenerator = async function* (duplex) {
    for await (const message of duplex) {
      yield message
    }
  }


for await (let msg of messageGenerator(socketApi.duplex)) {
    console.log('start process')
    await this.messageHandler.handleMessage(msg, user)
    console.log('end process')
}

日志:

  1. 开始进程
  2. 开始进程
  3. 结束进程
  4. 结束进程

不幸的是,如您所见,消息会继续被处理,而无需等待前一个消息完成。你有解决这个问题的办法吗? 我最终应该使用 Redis 的队列来处理消息吗?

谢谢

【问题讨论】:

  • 很难知道发生了什么而不知道其中包含什么函数,以及它是执行两条消息还是执行一条消息的两次。毕竟,有一个浮动的 for await 似乎不属于任何东西(在您的示例中,第二个位于最外层范围)

标签: javascript node.js ecmascript-6 websocket promise


【解决方案1】:

我不是 nodeJS 人,但我曾多次用其他语言思考过同样的问题。我得出的结论是,消息处理操作的速度有多慢真的很重要,因为如果它们太慢(慢于某个阈值,取决于每秒的 msg 值),这可能会导致 websocket 连接出现瓶颈,并且当这个瓶颈建立时可能会导致未来消息的极度延迟。

如果awaitasync与python中的行为相同,如果你使用它们处理任何操作,你的处理将是异步的,这意味着它确实不会等待前一个被处理。

到目前为止,我有两种选择:

  1. 继续异步处理消息,但在处理它们的代码中编写一些额外的逻辑,以管理订单混乱。例如,在继续处理当前消息之前,确认之前的消息已经被处理。此逻辑可能既复杂又缓慢,因为它在单独的线程中运行并且不会阻塞 websocket 消息。
  2. 同步处理消息,一条一条地处理消息,但速度非常快,只需执行一项操作:将它们存储在 Redis 中。这比将它们存储在数据库中要快得多,并且在大多数情况下速度足够快,不会导致 WS 连接出现瓶颈。然后使用单独的进程从 Redis 获取这些消息并进行处理。

【讨论】:

  • 谢谢,我走了这条路,它的效果与我的担心完全相反。对于未来的读者,我使用 Bee-Queue 和 redis 来快速处理消息。
猜你喜欢
  • 2012-07-12
  • 1970-01-01
  • 1970-01-01
  • 2012-12-26
  • 2014-05-18
  • 2019-03-02
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多