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

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存储,仅统计当前会话内的消息数

额外优化建议

  1. 替换内存存储:你当前使用的MemoryBotStorage是临时内存存储,Bot重启后所有计数和队列状态都会丢失。建议换成AzureBotStorage或其他持久化存储方案。
  2. 队列清理:可以定期清理长时间无活动的会话队列,避免内存占用过高。
  3. 重试机制:在API请求失败时,加入有限次数的重试逻辑,减少消息丢失概率。

内容的提问来源于stack exchange,提问作者Odirlei Borgert

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 09:36:07