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

求助:如何让Node.js服务器串行处理单用户GraphQL请求

解决Apollo Server中单用户请求串行化的方案

针对离线优先应用中mutation队列的时序问题,以下是几种规范的服务器端串行处理单用户请求的方案:

1. 基于Apollo Server插件的内存队列实现

利用Apollo Server的插件系统,为每个用户维护一个Promise链,确保同一用户的请求按顺序串行执行,无需依赖外部服务:

// 存储每个用户的当前请求处理Promise链
const userRequestQueue = new Map();

const serialProcessingPlugin = {
  async requestDidStart(requestContext) {
    const userId = requestContext.context.userId; // 从context中提取已解析的用户ID
    if (!userId) return {};

    // 获取用户当前的队列Promise,无则初始化为已完成状态
    const currentQueue = userRequestQueue.get(userId) || Promise.resolve();

    // 构建新的Promise链,等待前一个请求完成后再执行当前请求
    const newQueue = currentQueue
      .then(async () => {
        try {
          // 等待当前请求的执行阶段完成
          await requestContext.execute();
        } catch (err) {
          // 捕获错误避免队列中断,可根据需求记录日志
          console.error(`用户${userId}请求执行失败:`, err);
        }
      })
      .finally(() => {
        // 队列完成后清理,避免内存泄漏
        if (userRequestQueue.get(userId) === newQueue) {
          userRequestQueue.delete(userId);
        }
      });

    userRequestQueue.set(userId, newQueue);

    return {
      async executionDidStart() {
        // 强制当前请求等待队列中的前序请求完成
        await newQueue;
      }
    };
  }
};

// 初始化Apollo Server时挂载插件
const server = new ApolloServer({
  typeDefs,
  resolvers,
  plugins: [serialProcessingPlugin],
  context: ({ req }) => {
    // 从请求头解析token,提取用户ID(示例逻辑,需替换为实际认证逻辑)
    const token = req.headers.authorization?.split(' ')[1];
    const userId = token ? verifyToken(token).userId : null;
    return { userId };
  }
});

注意事项:

  • 该方案基于内存存储,若服务器为多实例部署,需改用分布式存储(如Redis)实现跨实例的队列同步
  • 需确保用户ID的准确性,避免不同用户的请求被错误串行化

2. 基于Redis的分布式请求队列

对于多实例部署场景,使用Redis的List结构实现每个用户的请求队列,配合消费进程串行处理:

步骤1:请求入队

在Apollo Server的resolver或插件中,将用户的mutation请求参数存入Redis队列:

const redis = require('redis');
const client = redis.createClient();

const resolvers = {
  Mutation: {
    createObject: async (_, { input }, { userId }) => {
      // 将请求信息存入用户专属队列
      await client.rPush(`user:${userId}:mutationQueue`, JSON.stringify({
        operation: 'createObject',
        input,
        clientMutationId: input.clientMutationId
      }));
      // 立即返回客户端请求已接收,后续由消费进程处理
      return { success: true, message: '请求已加入同步队列' };
    },
    updateObject: async (_, { input }, { userId }) => {
      await client.rPush(`user:${userId}:mutationQueue`, JSON.stringify({
        operation: 'updateObject',
        input,
        clientMutationId: input.clientMutationId
      }));
      return { success: true, message: '请求已加入同步队列' };
    }
  }
};

步骤2:串行消费队列

编写独立的Node.js进程,轮询Redis队列并串行处理每个用户的请求:

const redis = require('redis');
const client = redis.createClient();
const { executeMutation } = require('./mutationExecutor'); // 自定义执行mutation的逻辑

async function processUserQueue(userId) {
  const queueKey = `user:${userId}:mutationQueue`;
  while (true) {
    // 阻塞式取出队列头部请求
    const requestStr = await client.blPop(queueKey, 0);
    if (!requestStr) continue;
    const request = JSON.parse(requestStr[1]);
    try {
      await executeMutation(request.operation, request.input);
      // 可选:记录已处理的mutation ID,实现幂等性
      await client.setEx(`mutation:${userId}:${request.clientMutationId}`, 86400, 'processed');
    } catch (err) {
      console.error(`处理用户${userId}的${request.operation}请求失败:`, err);
      // 可选:将失败请求重新放回队列尾部重试
      await client.rPush(queueKey, JSON.stringify(request));
    }
  }
}

// 监听新队列的创建(可选,或定时扫描所有用户队列)
client.on('message', (channel, userId) => {
  processUserQueue(userId);
});
client.subscribe('newUserQueue');

优势:

  • 天然支持多实例部署,队列状态全局一致
  • 可灵活配置重试机制、失败兜底策略

3. 辅助:mutation幂等性与存在性校验

即使实现了串行化,也建议为每个mutation添加幂等性保障,并在修改类请求中增加对象存在性校验,避免极端情况下的时序问题:

const resolvers = {
  Mutation: {
    updateObject: async (_, { input }, { userId, redis }) => {
      const { clientMutationId, objectId, ...data } = input;
      
      // 检查是否已处理过该mutation
      const isProcessed = await redis.get(`mutation:${userId}:${clientMutationId}`);
      if (isProcessed) {
        return { success: true, message: '请求已处理' };
      }
      
      // 检查对象是否存在,不存在则重试几次
      let object = await ObjectModel.findById(objectId);
      if (!object) {
        for (let i = 0; i < 3; i++) {
          await new Promise(resolve => setTimeout(resolve, 150));
          object = await ObjectModel.findById(objectId);
          if (object) break;
        }
        if (!object) {
          throw new Error('关联对象不存在,请稍后重试');
        }
      }
      
      // 执行更新操作
      object = await ObjectModel.findByIdAndUpdate(objectId, data, { new: true });
      // 标记请求已处理
      await client.setEx(`mutation:${userId}:${clientMutationId}`, 86400, 'processed');
      
      return { success: true, object };
    }
  }
};

内容的提问来源于stack exchange,提问作者Michal Bachratý

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 23:50:34