NodeJS Fork负载均衡处理Fastify并行API高并发响应错误问题
问题描述
通过Node.js Fork子进程处理高内存的并行API调用任务,成功降低了Fastify服务器的延迟,但在高并发(高TPS)场景下,约60%的请求会返回错误响应。具体场景为:Fastify处理浏览器请求时需发起100+并行axios下游调用,用Fork卸载并行逻辑后延迟问题得到缓解,但高并发下大量请求失败。
以下是相关实现代码:
ForkBalancer.js
import { ChildProcess, fork } from 'child_process'; const requestLimit = 0; interface forkResponse { kill: boolean; string?: string; } class ForkBalancer { path: string; forks: number; maxRAM?: number; args?: Array<string>; private activeFork: number; private resolvers = new Map(); private renderers: Array<ChildProcess>; constructor({ path = '', forks = 5, maxRAM = 250, args = [] }) { this.activeFork = 0; this.forks = forks; this.maxRAM = maxRAM; this.path = path; this.args = args; this.renderers = Array.from({ length: forks }, () => this.createFork()); } public getFromRenderer(params: any): Promise<forkResponse> { const { resolvers, maxRAM, activeFork, restartFork, renderers } = this; const renderer = renderers[activeFork]; return new Promise(function(resolve, reject) { try { renderer.once('message', (res: any) => { resolvers.delete(params.request.url); resolve(res); if (res.kill) restartFork(); }); if (!resolvers.has(params.request.url)) { renderer.setMaxListeners(requestLimit); resolvers.set(params.request.url, resolve); renderer.send({ ...params, maxRAM }); } } catch (error) { resolvers.delete(params.request.url); reject(error); } }); } private createFork = () => { const { path, args } = this; return fork(path, args); }; private restartFork = () => { const { activeFork, renderers, next, createFork } = this; const renderer = renderers[activeFork]; next(); renderer.kill(); this.renderers[activeFork] = createFork(); }; private next = () => { const { activeFork, forks } = this; if (activeFork === forks - 1) { this.activeFork = 0; } else { this.activeFork++; } }; } export default ForkBalancer;
ParallelAPI.js
import axiosInstance, { AxiosRequestConfig } from 'axios'; const maxRAM = 128 process.on('message', async (params: any) => { const { totalPages, offset: offsetProps = 0, PAGE_SIZE_LIMIT, request, body, testId, url } = params; const requests = []; for (let offset = offsetProps; offset <= totalPages; offset++) { requests.push( axiosInstance.post( `API_URL/search/v2?page=${offset}&limit=${PAGE_SIZE_LIMIT}`, body, { headers: { accept: 'application/json', authorization: `${request.token?.token_type} ${request.token?.access_token}`, } }, ) ); } const results = await Promise.allSettled(requests); const list: any = []; let isPartialFailed = false; results.forEach((result) => { if (result.status === 'fulfilled') { const quotesListData = result.value?.data?.quotes; if (Array.isArray(quotesListData)) { list.push(...quotesListData); } } else { isPartialFailed = true; } }); const { heapUsed } = process.memoryUsage(); if (process.send) { process.send({ key: request.url, list, url: request.url, testId, isPartialFailed, kill: heapUsed > maxRAM * 1024 * 1024, }); } });
业务实现代码
import path from 'path'; import ForkBalancer from './forkBalancer'; const forkBalancer = new ForkBalancer({ path: path.resolve(__dirname, './ParallelAPI'), }); const handler = async (req, res) => { const { body } = req.body; const { testId } = req.query; const response = await forkBalancer.getFromRenderer({ request: { token: request.token, url: request.url }, testId, PAGE_SIZE_LIMIT: 50, body, totalPages: 100 }); return res.send(response); } fastify.post('/getAllItems', handler);
问题分析与修复方案
核心问题点
- 请求标识冲突:用
request.url作为请求唯一标识,高并发下多个请求会共用同一URL,导致resolvers映射表被覆盖,大量请求无法正确resolve。 - 事件监听覆盖:同一子进程处理多个请求时,后续的
once('message')监听会覆盖之前的,仅最后一个请求能收到响应,其余请求挂起或报错。 - 监听器数量限制错误:
renderer.setMaxListeners(0)将子进程事件监听器上限设为0,触发MaxListenersExceededWarning甚至直接报错。 - 子进程重启逻辑不安全:直接kill旧进程导致未处理请求丢失,新进程未承接旧任务。
- 业务代码变量错误:handler中误用
request而非req,导致token和url未定义,直接抛出错误。
修复后的代码
ForkBalancer.js(修复版)
import { ChildProcess, fork } from 'child_process'; import { v4 as uuidv4 } from 'uuid'; interface forkResponse { kill: boolean; requestId: string; list?: any[]; isPartialFailed?: boolean; testId?: string; } interface RequestParams { request: { token?: any; url: string }; testId?: string; PAGE_SIZE_LIMIT: number; body: any; totalPages: number; } class ForkBalancer { path: string; forks: number; maxRAM?: number; args?: string[]; private activeForkIdx: number; private resolvers = new Map<string, (value: forkResponse) => void>(); private renderers: ChildProcess[]; constructor({ path = '', forks = 5, maxRAM = 250, args = [] }) { this.activeForkIdx = 0; this.forks = forks; this.maxRAM = maxRAM; this.path = path; this.args = args; this.renderers = Array.from({ length: forks }, () => this.createFork()); } public getFromRenderer(params: RequestParams): Promise<forkResponse> { const requestId = uuidv4(); const currentForkIdx = this.activeForkIdx; const renderer = this.renderers[currentForkIdx]; this.nextFork(); return new Promise((resolve, reject) => { try { const messageHandler = (res: forkResponse) => { if (res.requestId === requestId) { renderer.off('message', messageHandler); this.resolvers.delete(requestId); resolve(res); if (res.kill) this.restartFork(currentForkIdx); } }; renderer.on('message', messageHandler); this.resolvers.set(requestId, resolve); renderer.send({ ...params, maxRAM: this.maxRAM, requestId }); if (renderer.getMaxListeners() < this.forks * 2) { renderer.setMaxListeners(this.forks * 2); } } catch (error) { this.resolvers.delete(requestId); reject(error); } }); } private createFork(): ChildProcess { const forkProcess = fork(this.path, this.args); forkProcess.on('exit', (code) => { console.log(`子进程退出,代码:${code}`); }); return forkProcess; } private restartFork(forkIdx: number): void { const oldFork = this.renderers[forkIdx]; oldFork.kill('SIGTERM'); this.renderers[forkIdx] = this.createFork(); } private nextFork(): void { this.activeForkIdx = (this.activeForkIdx + 1) % this.forks; } } export default ForkBalancer;
ParallelAPI.js(修复版)
import axiosInstance from 'axios'; import { default as pLimit } from 'p-limit'; const maxRAM = 128; const concurrencyLimit = 20; const limit = pLimit(concurrencyLimit); process.on('message', async (params: any) => { const { totalPages, offset: offsetProps = 0, PAGE_SIZE_LIMIT, request, body, testId, requestId, maxRAM } = params; const requests = []; for (let offset = offsetProps; offset <= totalPages; offset++) { requests.push(limit(async () => { try { return await axiosInstance.post( `API_URL/search/v2?page=${offset}&limit=${PAGE_SIZE_LIMIT}`, body, { headers: { accept: 'application/json', authorization: `${request.token?.token_type} ${request.token?.access_token}`, }, timeout: 10000, } ); } catch (err) { return { status: 'rejected', reason: err }; } })); } const results = await Promise.all(requests); const list: any[] = []; let isPartialFailed = false; results.forEach((result) => { if (result.status !== 'rejected' && result.data?.quotes) { if (Array.isArray(result.data.quotes)) { list.push(...result.data.quotes); } } else { isPartialFailed = true; } }); const { heapUsed } = process.memoryUsage(); if (process.send) { process.send({ requestId, list, testId, isPartialFailed, kill: heapUsed > maxRAM * 1024 * 1024, }); } });
业务实现代码(修复版)
import path from 'path'; import ForkBalancer from './forkBalancer'; const forkBalancer = new ForkBalancer({ path: path.resolve(__dirname, './ParallelAPI'), forks: 8, }); const handler = async (req, res) => { try { const { body } = req.body; const { testId } = req.query; const response = await forkBalancer.getFromRenderer({ request: { token: req.token, url: req.url }, testId, PAGE_SIZE_LIMIT: 50, body, totalPages: 100 }); return res.send(response); } catch (error) { console.error('请求处理失败:', error); return res.status(500).send({ error: '服务器内部错误' }); } } fastify.post('/getAllItems', handler);
额外优化建议
- 子进程资源监控:定期检查子进程CPU、内存使用,超出阈值自动重启。
- 请求超时处理:在ForkBalancer中为每个请求设置超时时间,避免挂起。
- 下游服务容错:为axios添加重试机制,针对可重试错误自动重试。
- 负载均衡优化:根据子进程待处理请求数分配任务,替代简单轮询。
内容的提问来源于stack exchange,提问作者Srigar
相关产品推荐
相关产品推荐

