【发布时间】: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')
}
日志:
- 开始进程
- 开始进程
- 结束进程
- 结束进程
不幸的是,如您所见,消息会继续被处理,而无需等待前一个消息完成。你有解决这个问题的办法吗? 我最终应该使用 Redis 的队列来处理消息吗?
谢谢
【问题讨论】:
-
很难知道发生了什么而不知道其中包含什么函数,以及它是执行两条消息还是执行一条消息的两次。毕竟,有一个浮动的
for await似乎不属于任何东西(在您的示例中,第二个位于最外层范围)
标签: javascript node.js ecmascript-6 websocket promise