Node.js Bot Framework消息乱序及消息计数实现方案咨询
解决Bot Framework消息同步乱序与消息计数问题
针对你遇到的两个核心问题,我来一步步帮你梳理解决方案:
一、解决消息同步到自定义API的乱序问题
消息乱序的核心原因是异步HTTP请求的网络延迟不可控——即使你在中间件里用了await,如果前一个请求因为网络波动变慢,后一个请求可能会先到达你的API。
最优解决方案:为每个会话维护串行请求队列
我们可以给每个会话(通过conversationId标识)创建一个独立的异步任务链,确保消息严格按照Bot接收/发送的顺序,串行发送到API,从根源避免并发导致的乱序。
修改你的中间件代码,加入队列管理逻辑:
// 在Bot初始化前,添加队列管理工具 const conversationRequestQueues = new Map(); // key: conversationId, value: Promise链 // 辅助函数:将API请求加入串行队列 const enqueueApiCall = async (conversationId, apiCallFn) => { // 为新会话初始化空队列 if (!conversationRequestQueues.has(conversationId)) { conversationRequestQueues.set(conversationId, Promise.resolve()); } // 将新请求追加到队列末尾,确保串行执行 const currentQueue = conversationRequestQueues.get(conversationId); const updatedQueue = currentQueue .then(apiCallFn) .catch(err => console.error('API请求失败:', err.message)); conversationRequestQueues.set(conversationId, updatedQueue); await updatedQueue; }; // 修改原有的bot.use中间件 bot.use({ // 用户发送消息到Bot的处理逻辑 receive: async (event, next) => { if ((event.type === 'message' || event.type === 'conversationUpdate') && process.env.API) { // 用队列串行发送消息到API await enqueueApiCall(event.address.conversation.id, () => api.post.receive(event)); } next(); }, // Bot发送消息到用户的处理逻辑 send: async (event, next) => { if (event.type === 'message' && process.env.API) { await enqueueApiCall(event.address.conversation.id, () => api.post.send(event)); } next(); } });
这个方案的核心逻辑:
- 每个会话对应一个独立的Promise链,避免跨会话的消息干扰
- 新的API请求必须等待前一个请求完成后才会执行,严格保证顺序
- 加入错误捕获,避免单个请求失败导致整个队列阻塞
二、为消息添加计数功能(Bot Framework无原生支持)
Bot Framework本身没有提供原生的消息计数能力,需要我们基于Bot的存储系统自行实现。我们可以分别统计用户发送的消息数和Bot发送的消息数,存储维度可以选择「用户维度」或「会话维度」,根据业务需求决定。
实现示例(结合你的代码)
我们在中间件中更新计数,并将计数附加到发送给API的消息中:
bot.use({ receive: async (event, next) => { if (event.type === 'message' && process.env.API) { // 更新用户发送消息计数(用户维度:同一个用户所有会话累加) const userData = await bot.getUserData(event.address); userData.userMessageCount = (userData.userMessageCount || 0) + 1; await bot.setUserData(event.address, userData); // 携带计数发送到API await enqueueApiCall(event.address.conversation.id, () => api.post.receive({ ...event, userMessageCount: userData.userMessageCount }) ); } else if (event.type === 'conversationUpdate' && process.env.API) { await enqueueApiCall(event.address.conversation.id, () => api.post.receive(event)); } next(); }, send: async (event, next) => { if (event.type === 'message' && process.env.API) { // 更新Bot发送消息计数(会话维度:仅当前会话内累加) const conversationData = await bot.getConversationData(event.address); conversationData.botMessageCount = (conversationData.botMessageCount || 0) + 1; await bot.setConversationData(event.address, conversationData); // 携带计数发送到API await enqueueApiCall(event.address.conversation.id, () => api.post.send({ ...event, botMessageCount: conversationData.botMessageCount }) ); } next(); } });
计数维度说明:
- 用户维度:用
getUserData/setUserData存储,同一个用户在所有会话中的消息都会累加 - 会话维度:用
getConversationData/setConversationData存储,仅统计当前会话内的消息数
额外优化建议
- 替换内存存储:你当前使用的
MemoryBotStorage是临时内存存储,Bot重启后所有计数和队列状态都会丢失。建议换成AzureBotStorage或其他持久化存储方案。 - 队列清理:可以定期清理长时间无活动的会话队列,避免内存占用过高。
- 重试机制:在API请求失败时,加入有限次数的重试逻辑,减少消息丢失概率。
内容的提问来源于stack exchange,提问作者Odirlei Borgert
相关产品推荐
相关产品推荐

