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

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);

问题分析与修复方案

核心问题点

  1. 请求标识冲突:用request.url作为请求唯一标识,高并发下多个请求会共用同一URL,导致resolvers映射表被覆盖,大量请求无法正确resolve。
  2. 事件监听覆盖:同一子进程处理多个请求时,后续的once('message')监听会覆盖之前的,仅最后一个请求能收到响应,其余请求挂起或报错。
  3. 监听器数量限制错误:renderer.setMaxListeners(0)将子进程事件监听器上限设为0,触发MaxListenersExceededWarning甚至直接报错。
  4. 子进程重启逻辑不安全:直接kill旧进程导致未处理请求丢失,新进程未承接旧任务。
  5. 业务代码变量错误: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);

额外优化建议

  1. 子进程资源监控:定期检查子进程CPU、内存使用,超出阈值自动重启。
  2. 请求超时处理:在ForkBalancer中为每个请求设置超时时间,避免挂起。
  3. 下游服务容错:为axios添加重试机制,针对可重试错误自动重试。
  4. 负载均衡优化:根据子进程待处理请求数分配任务,替代简单轮询。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 17:20:58