求助:如何让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ý
相关产品推荐
相关产品推荐

