首页 / 文章 / 使用 worker_threads 和 Pool 在 Node.js 中卸载 CPU 密集型任务

使用 worker_threads 和 Pool 在 Node.js 中卸载 CPU 密集型任务

了解 Node.js worker_threads 如何保持事件循环的响应能力:工作进程的创建、请求-响应通信、可传输缓冲区、线程池以及相关隐患。

4814 词

Node.js 能在单线程事件循环中处理数千个并发连接,因为几乎所有工作都在等待网络或磁盘响应。但一旦请求需要实际计算——比如对大型数据包进行哈希处理、调整图像大小或分析数据集——这个模型就会失效,因为唯一的 JavaScript 线程会被占满,其他所有请求都不得不排队等待。worker_threads 模块正是为解决这一问题而设计的内置方案。读完本指南后,你将能够把依赖 CPU 的代码放到工作线程中执行,与它们高效地交换数据,以线程池的形式运行这些工作线程,并判断在哪些情况下使用工作线程并不合适。

工作线程存在的原因

worker_threads首次出现在Node.js 10.5.0版本中,当时处于实验性状态,从Node.js 12系列起被视作稳定功能。它允许单个Node.js进程同时在多个线程上运行JavaScript代码。

它与child_process在重要方面存在差异。子进程是完全独立的Node.js进程,拥有自己的内存和运行时副本。而工作线程则存在于同一个进程中,但并非简单的“共享所有资源的另一个线程”。每个工作线程都有独立的V8隔离环境、堆内存以及事件循环。没有任何东西是默认共享的;工作线程之间通过传递消息进行通信,只有当你明确将SharedArrayBuffer交给它们时才会共享内存。(有一种常见的说法认为工作线程会共享主线程的V8实例和内存,但这并不准确,而这种差异也解释了下文中的大部分API设计。)与进程相比,工作线程的启动和通信成本更低,因此成为单个应用程序实例充分利用多核资源的理想方式。

被阻塞的事件循环会带来什么代价

工作线程的核心职责是保持事件循环的畅通。当主线程上正在执行同步计算时,服务器无法接收新请求,也无法处理现有连接,甚至连健康检查都无法响应。实际后果如下:

  • 用户体验下降:API响应延迟,实时功能如动态更新也会出现卡顿。
  • 吞吐量降低:一个耗时的任务占用了唯一的JavaScript线程,导致每秒处理的请求数减少。
  • 稳定性问题:长时间的阻塞可能超出负载均衡器或客户端的超时时间,这些故障还可能波及其他服务。

将繁重的计算任务转移到工作线程后,主线程便可继续处理套接字和定时器,因此即使在执行复杂计算时,应用程序仍能保持响应能力。

如需更深入地了解事件循环、libuv及其线程池之间的协作机制,请参阅Node.js并发机制的底层原理

充分利用每一核心

服务器通常拥有众多CPU核心。如果没有工作进程,或者没有通过child_process或PM2之类的进程管理器启动多个进程,单个Node.js进程在处理CPU密集型的JavaScript任务时大约只能使用一个核心,其余核心则处于闲置状态。工作进程能够让单个应用程序在可用核心上并行执行多个繁重任务,从而提升整体计算效率。

可实际应用的工作负载

并行计算为那些传统上需要依赖一流线程支持的语言开辟了新的应用领域:

  • 数据处理:实时分析、图像与视频处理,以及大规模数据集的转换。
  • 机器学习推理:在应用程序保持交互状态的同时运行已训练好的模型。
  • 加密技术:哈希运算、加密与解密。
  • 仿真与科学计算:并行仿真及复杂的数学模型运算。

关键在于,这些功能都不会影响异步I/O模型。工作线程只是在原有基础上增加了并行计算能力。

该模块的构成要素

worker_threads提供了一小组基础功能:

  • Worker主线程用于从脚本中启动新工作线程的类。
  • isMainThread一个布尔值,用于指示当前代码是否在主线程上运行;当一个文件可能同时扮演两种角色时,这一属性非常有用。
  • parentPort仅在工作线程中可用,它是通往创建该工作线程的线程的通道。在主线程上该值为null
  • workerData仅在工作线程中可用,其中保存了创建工作线程时父线程所提供的数据的副本。
  • MessagePortMessageChannel用于在多个线程之间创建额外的独立双向通道。
  • SharedArrayBufferAtomics真正的共享内存以及安全使用它所需的同步原语。这是效率最高的方案,但同时也最为复杂,因为它们会重新引入竞态条件。
  • 第一个工作线程:在主线程之外寻找素数

    一个不错的入门练习是设计一个故意较慢的任务:计算出某个较大上限以内的所有素数。该示例使用了位于同一文件夹中的两个文件。

    主线程脚本

    主脚本会创建工作线程,为其提供输入并等待回复。在读取结果之前,请注意三点:工作线程是从独立的文件路径创建的;输入数据通过workerData选项传递;而父脚本会监听三个事件:用于获取结果的message、用于处理工作线程未能捕获的异常的error,以及用于检测线程停止时的exit

    // main.js
    const { Worker, isMainThread, workerData } = require('worker_threads');
    const path = require('path');
    
    if (isMainThread) {
        console.log(`Main thread started. PID: ${process.pid}`);
    
        const largeNumber = 20_000_000; // A large number for prime calculation
    
        console.log(`Starting CPU-intensive task (finding primes up to ${largeNumber})...`);
        const startTime = Date.now();
    
        const worker = new Worker(path.join(__dirname, 'prime-worker.js'), {
            workerData: { limit: largeNumber }
        });
    
        worker.on('message', (result) => {
            const endTime = Date.now();
            console.log(`Worker finished. Found ${result.length} primes.`);
            console.log(`Time taken by worker: ${endTime - startTime}ms`);
            // console.log('Primes found:', result.slice(0, 10), '...'); // Log first 10 for brevity
    
            // Demonstrate a simple non-blocking task on the main thread
            console.log('Main thread continuing with other tasks...');
            setTimeout(() => {
                console.log('Main thread completed a separate non-blocking task.');
            }, 100);
        });
    
        worker.on('error', (err) => {
            console.error('Worker error:', err);
        });
    
        worker.on('exit', (code) => {
            if (code !== 0)
                console.error(`Worker stopped with exit code ${code}`);
            else
                console.log('Worker thread exited normally.');
        });
    
    } else {
        // This part will not be executed in our scenario as worker.js runs separately.
        // It's here for illustrative purposes if main.js itself was to be imported as a worker.
        console.log('This code runs inside a worker thread if main.js was spawned as one.');
        console.log('WorkerData:', workerData);
    }
    

    重点代码行解析:

    • isMainThread检查确保生成逻辑仅在主线程上运行。else分支的存在只是为了说明如果将同一文件作为工作线程加载会发生什么;在这种配置下它永远不会被执行。
    • largeNumber被设置得足够大(2000万),以便让计算过程需要较长的时间。
  • new Worker(path.join(__dirname, 'prime-worker.js'), { workerData: { limit: largeNumber } }) 是核心调用语句。第一个参数指定了新线程要运行的脚本,该脚本必须是独立的文件;第二个参数是一个选项对象,其中的 workerData 会成为工作线程的初始输入数据。
  • workerData 是通过 HTML 的结构化克隆算法进行复制的。工作线程收到的是自己的副本,而非对父对象的引用。
  • message 监听器会接收工作线程发送的任何数据,在此例中即为质数数组;error 监听器会在工作线程内部出现未捕获的异常时触发;而 exit 监听器则会收到退出码,其中 0 表示正常结束。
  • 工作线程脚本

    该工作进程仅包含昂贵的计算任务以及用于输出结果的代码。findPrimes函数是一个简单的试除法循环,故意未做优化以便消耗更多CPU时间。

    // prime-worker.js
    const { parentPort, workerData } = require('worker_threads');
    
    // A simple (not highly optimized) function to find prime numbers
    function findPrimes(limit) {
        const primes = [];
        for (let i = 2; i <= limit; i++) {
            let isPrime = true;
            for (let j = 2; j <= Math.sqrt(i); j++) {
                if (i % j === 0) {
                    isPrime = false;
                    break;
                }
            }
            if (isPrime) {
                primes.push(i);
            }
        }
        return primes;
    }
    
    // Ensure this code only runs if it's indeed a worker thread being executed
    if (parentPort) {
        const { limit } = workerData;
        console.log(`Worker thread started to find primes up to ${limit}.`);
    
        // Simulate an error for demonstration purposes sometimes
        // if (Math.random() < 0.2) {
        //     throw new Error('Simulated worker error!');
        // }
    
        try {
            const result = findPrimes(limit);
            parentPort.postMessage(result); // Send the result back to the main thread
        } catch (error) {
            console.error('Error during prime calculation in worker:', error);
            parentPort.postMessage({ error: error.message }); // Send error back
        }
    } else {
        console.log("This script is intended to be run as a worker thread.");
    }
    

    需要注意以下几点:

    • 它从模块中获取parentPortworkerData
    • if (parentPort)条件确保该逻辑仅在文件作为工作进程加载时执行;如果直接使用node命令运行,则parentPort的值为null,此时会打印相应消息。
    • const { limit } = workerData;用于读取父进程发送的输入数据。
    • parentPort.postMessage(result)将筛选出的素数发送到主线程,由message处理函数接收这些数据。
  • try...catch块会记录计算过程中的任何失败,并将其上报给父进程,而不会让线程终止。
  • 有一个需要注意的细节:当发生错误时,工作线程会在用于传递结果的同一message通道上发送{ error: ... },但主脚本会将每条消息都视为质数数组。在实际代码中,应为消息指定明确的结构(例如添加statustype字段),这样父进程就能区分结果与错误。

    运行程序并查看输出

    将两个文件保存在相邻位置,然后使用node main.js启动程序。你会看到主线程宣布自己已启动,工作线程报告它也已开始运行,过一段时间后则会显示工作线程的计算结果以及所花费的时间。

    注意“主线程继续执行其他任务...”这一行的输出顺序。在这段脚本中,它是在message处理函数内部打印的,因此只有在工作线程完成任务后才会显示。这表明主线程能够接收并处理结果,而非在此期间正在执行其他任务。若要真正看到主线程在计算过程中仍保持响应状态,应在结果返回之前就在主线程上启动某些操作,比如使用setInterval每隔几百毫秒记录一次心跳信息。在工作线程计算期间,心跳信号会持续发送,而这正是如果findPrimes在主线程上运行时不会出现的情况。

    基于请求-响应协议的双向通信

    第一个示例发送一个输入并得到一个结果。大多数实际应用需要一个长期运行的工作线程来处理大量请求。通信是完全双向的:主线程中的Worker对象与工作线程内的parentPort都提供了.postMessage(...).on('message', ...)方法。

    下面的工作线程会持续运行并等待消息。当收到calculateFibonacci请求时,它会使用一种简单的递归函数来计算结果(故意设计得较慢),然后通过fibonacciResultfibonacciError消息进行回复,同时附上请求的标识符,以便调用方能够对应上该回复。

    // request-response-worker.js
    const { parentPort } = require('worker_threads');
    
    function fibonacci(n) {
        if (n <= 1) return n;
        return fibonacci(n - 1) + fibonacci(n - 2);
    }
    
    if (parentPort) {
        parentPort.on('message', (message) => {
            if (message.type === 'calculateFibonacci') {
                const { n, requestId } = message.payload;
                console.log(`Worker: Calculating Fibonacci(${n}) for request ID ${requestId}`);
                try {
                    const result = fibonacci(n);
                    parentPort.postMessage({ type: 'fibonacciResult', requestId, payload: result });
                } catch (error) {
                    console.error(`Worker: Error calculating Fibonacci(${n}):`, error);
                    parentPort.postMessage({ type: 'fibonacciError', requestId, payload: error.message });
                }
            }
        });
        console.log('Worker ready to receive Fibonacci requests.');
    }
    

    在主线程中,关键在于将消息传递转换为承诺对象。每次调用都会生成一个唯一的requestId,将该承诺的resolvereject方法存储在对应ID下的Map中,随后发起请求。当收到回复后,message处理函数会根据该ID查找对应的记录并移除它,从而完成对应承诺的解决。示例中使用了uuid包来生成ID;在当前版本的Node.js中,内置的node:crypto模块中的crypto.randomUUID()函数也能实现相同功能且无需额外依赖。

    // main-request-response.js
    const { Worker, isMainThread } = require('worker_threads');
    const path = require('path');
    const { v4: uuidv4 } = require('uuid'); // npm install uuid
    
    if (isMainThread) {
        const worker = new Worker(path.join(__dirname, 'request-response-worker.js'));
        const pendingRequests = new Map();
    
        worker.on('message', (message) => {
            const { type, requestId, payload } = message;
            if (pendingRequests.has(requestId)) {
                const { resolve, reject } = pendingRequests.get(requestId);
                pendingRequests.delete(requestId);
    
                if (type === 'fibonacciResult') {
                    resolve(payload);
                } else if (type === 'fibonacciError') {
                    reject(new Error(payload));
                }
            }
        });
    
        worker.on('error', (err) => {
            console.error('Worker error:', err);
        });
    
        worker.on('exit', (code) => {
            if (code !== 0) console.error(`Worker stopped with exit code ${code}`);
        });
    
        async function calculateFibonacciInWorker(n) {
            const requestId = uuidv4();
            return new Promise((resolve, reject) => {
                pendingRequests.set(requestId, { resolve, reject });
                worker.postMessage({ type: 'calculateFibonacci', payload: { n }, requestId });
            });
        }
    
        (async () => {
            console.log('Main: Sending Fibonacci requests to worker...');
            try {
                const fib40 = await calculateFibonacciInWorker(40);
                console.log(`Main: Fibonacci(40) = ${fib40}`);
    
                const fib35 = await calculateFibonacciInWorker(35);
                console.log(`Main: Fibonacci(35) = ${fib35}`);
    
                // This will block the main thread if done directly, but not here
                // console.log(`Main: Local Fibonacci(40) = ${fibonacci(40)}`);
    
            } catch (error) {
                console.error('Main: Failed to get Fibonacci result from worker:', error);
            } finally {
                worker.terminate(); // Terminate worker when done
            }
        })();
    }
    

    总结来说,这种模式的特点是:

    • 主线程为每个请求生成一个requestId,这样即使回复按不同顺序到达,也能将其与对应的承诺关联起来。
  • pendingRequests 映表用于存储 resolvereject 函数,直到工作进程给出响应。
  • 工作进程负责处理 calculateFibonacci 消息,并携带原始的 requestIdfibonacciResultfibonacciError 的形式回复。
  • 如果修改此代码,请仔细注意消息的结构。主线程会将requestId放在消息的最顶层,与payload相邻,而工作线程则通过message.payload来获取它。按照当前的写法,工作线程会将requestId视为undefined,从而导致回复无法匹配,等待的承诺也永远不会得到解决。应在双方约定的同一位置保存该ID(例如payload: { n, requestId })。此外,还为每个待处理请求添加超时机制,这样丢失的回复会转化为错误而非程序挂起。最后,请注意finally块:一旦任务完成,worker.terminate()就会终止线程,从而使进程能够退出。

    无需复制即可传输大量数据

    通过 postMessage 发送的所有内容默认都会进行结构化克隆,这意味着接收方会得到一个副本。对于较小的消息来说这没有问题,但对于较大的二进制数据而言,复制本身就可能成为瓶颈。解决方法是转移底层内存的所有权而非进行复制:这样发送方无法再使用该对象,而接收方则可以立即使用它,且不会产生重复。

    ArrayBufferMessagePort 是最常被传输的对象。SharedArrayBuffer 则情况不同:它根本不会被传输,因为两个线程可以同时访问同一块内存。

    下面的示例在一个列表中包含了两个脚本。工作线程(buffer-worker.js)接收一个缓冲区,就地将每个字节的值翻倍,然后以可传输对象的形式将其返回。主脚本(main-buffer.js)则填充一个1 MB的缓冲区,将其传递给工作线程,并在工作线程返回后读取处理后的数据。

    // buffer-worker.js
    const { parentPort } = require('worker_threads');
    
    if (parentPort) {
        parentPort.on('message', (message) => {
            if (message.type === 'processBuffer') {
                const { buffer } = message.payload; // This is now the ArrayBuffer
                const uint8 = new Uint8Array(buffer);
    
                // Modify the buffer in the worker
                for (let i = 0; i < uint8.length; i++) {
                    uint8[i] = uint8[i] * 2;
                }
                console.log('Worker: Buffer processed. First 5 elements:', uint8.slice(0, 5));
    
                // Send it back as a transferable
                parentPort.postMessage({ type: 'bufferProcessed', payload: buffer }, [buffer]);
            }
        });
    }
    // main-buffer.js
    const { Worker, isMainThread } = require('worker_threads');
    const path = require('path');
    
    if (isMainThread) {
        const worker = new Worker(path.join(__dirname, 'buffer-worker.js'));
        const bufferSize = 1024 * 1024; // 1MB
        let myBuffer = new ArrayBuffer(bufferSize);
        let uint8 = new Uint8Array(myBuffer);
    
        // Initialize buffer
        for (let i = 0; i < uint8.length; i++) {
            uint8[i] = i % 256;
        }
        console.log('Main: Original buffer (first 5 elements):', uint8.slice(0, 5));
    
        worker.postMessage({ type: 'processBuffer', payload: { buffer: myBuffer } }, [myBuffer]);
    
        // After postMessage with transfer, myBuffer becomes detached/empty in main thread
        // Attempting to access it will result in an error or empty view.
        // console.log('Main: Buffer after transfer (should be detached):', uint8.slice(0, 5)); // This would likely show zeros or error.
    
        worker.on('message', (message) => {
            if (message.type === 'bufferProcessed') {
                const receivedBuffer = message.payload;
                const receivedUint8 = new Uint8Array(receivedBuffer);
                console.log('Main: Received processed buffer (first 5 elements):', receivedUint8.slice(0, 5));
                worker.terminate();
            }
        });
    }
    

    关键细节:

    • postMessage的第二个参数,即一边的[myBuffer]和另一边的[buffer],就是传输列表。
    • myBuffer被发送后,其所有权便转移给工作线程。在主线程中,该缓冲区会被视为已断开连接:其byteLength值变为0,而基于它的任何类型化数组视图,比如之前的uint8,其长度也会变为0。
  • 工作线程修改数据后将其传回,这样主线程就能再次收到缓冲区内容。
  • 这1 MB的数据在传输过程中不会被双向复制。
  • 其代价是传输完成后必须停止使用原始引用。如果两个线程确实需要同时访问这些数据,就需要创建副本或使用SharedArrayBuffer

    错误处理与资源清理

    工作线程和其他代码一样可能出错,而且在运行期间会占用内存和线程资源。可用的处理工具包括:

    • 主线程中的worker.on('error', handler)用于接收在工作线程中抛出但未被捕获的异常。
  • 主线程中的 worker.on('exit', handler)无论工作进程是正常结束还是崩溃,只要其停止运行就会触发该事件。退出码可用于区分这两种情况,0 表示成功。
  • 工作进程内部的 process.on('uncaughtException', handler)除了父进程能够看到的信息外,它还提供了在工作进程终止前记录日志或上报详细信息的途径。
  • worker.terminate()会尽快停止工作进程,并返回一个承诺对象,该对象会在进程退出时解析为对应的退出码。该方法不会等待正在执行的任务完成,因此如果需要优雅关闭,请向工作进程发送“停止”消息,让其自行退出。务必确保不再需要的工作进程能够以某种方式被终止。
  • 通过工作进程池执行大量任务

    为每个任务创建一个新的工作进程是浪费资源的行为,而且运行的CPU密集型工作进程数量超过核心数只会加剧竞争。工作进程池可以解决这两个问题:它启动固定数量的工作进程,将传入的任务放入队列中,然后把每个任务分配给下一个空闲的工作进程。

    下面这个简化的进程池包含一个工作进程记录数组、一份空闲工作进程ID列表以及一个待处理任务队列。runTask函数会返回一个承诺对象并将任务加入队列;processQueue函数则将最旧的任务与空闲的工作进程配对;当某个工作进程完成任务并反馈后,其对应的承诺对象就会得到解决,该工作进程也会重新加入空闲列表。如果某个工作进程出现错误或异常退出,terminateWorker函数会用一个新的工作进程替换它,以保持进程池的规模不变。

    // workerPool.js
    const { Worker } = require('worker_threads');
    const path = require('path');
    
    class WorkerPool {
        constructor(workerPath, numWorkers) {
            this.workerPath = workerPath;
            this.numWorkers = numWorkers;
            this.workers = [];
            this.freeWorkers = [];
            this.queue = [];
    
            this.initWorkers();
        }
    
        initWorkers() {
            for (let i = 0; i < this.numWorkers; i++) {
                const worker = new Worker(this.workerPath);
                worker.id = i; // Assign an ID for easier debugging
                worker.on('message', (message) => {
                    // Resolve the promise associated with this worker's task
                    this.workers[worker.id].resolve(message.payload);
                    this.returnWorkerToPool(worker.id);
                    this.processQueue();
                });
                worker.on('error', (err) => {
                    this.workers[worker.id].reject(err);
                    console.error(`Worker ${worker.id} error:`, err);
                    this.terminateWorker(worker.id); // Re-initialize or handle as appropriate
                });
                worker.on('exit', (code) => {
                    if (code !== 0) console.error(`Worker ${worker.id} exited with code ${code}`);
                    this.terminateWorker(worker.id); // Handle crashed worker
                });
                this.freeWorkers.push(worker.id);
                this.workers.push({ instance: worker, busy: false, resolve: null, reject: null });
            }
            console.log(`Worker Pool initialized with ${this.numWorkers} workers.`);
        }
    
        runTask(taskData) {
            return new Promise((resolve, reject) => {
                this.queue.push({ taskData, resolve, reject });
                this.processQueue();
            });
        }
    
        processQueue() {
            if (this.queue.length === 0 || this.freeWorkers.length === 0) {
                return;
            }
    
            const workerId = this.freeWorkers.shift();
            const workerInfo = this.workers[workerId];
            workerInfo.busy = true;
    
            const { taskData, resolve, reject } = this.queue.shift();
            workerInfo.resolve = resolve;
            workerInfo.reject = reject;
    
            workerInfo.instance.postMessage(taskData);
        }
    
        returnWorkerToPool(workerId) {
            const workerInfo = this.workers[workerId];
            workerInfo.busy = false;
            workerInfo.resolve = null;
            workerInfo.reject = null;
            this.freeWorkers.push(workerId);
        }
    
        terminateWorker(workerId) {
            const workerInfo = this.workers[workerId];
            if (workerInfo.instance) {
                workerInfo.instance.terminate();
            }
            // Remove from current workers list and potentially replace
            this.workers[workerId] = { instance: null, busy: false, resolve: null, reject: null };
            // Optionally, re-create the worker to maintain pool size
            console.log(`Worker ${workerId} terminated. Re-initializing...`);
            const newWorker = new Worker(this.workerPath);
            newWorker.id = workerId;
            newWorker.on('message', (message) => {
                this.workers[newWorker.id].resolve(message.payload);
                this.returnWorkerToPool(newWorker.id);
                this.processQueue();
            });
            newWorker.on('error', (err) => {
                this.workers[newWorker.id].reject(err);
                console.error(`Worker ${newWorker.id} error:`, err);
                this.terminateWorker(newWorker.id);
            });
            newWorker.on('exit', (code) => {
                if (code !== 0) console.error(`Worker ${newWorker.id} exited with code ${code}`);
                this.terminateWorker(newWorker.id);
            });
            this.workers[workerId] = { instance: newWorker, busy: false, resolve: null, reject: null };
            this.freeWorkers.push(workerId); // Add the new worker to the pool
        }
    
        close() {
            for (const workerInfo of this.workers) {
                if (workerInfo.instance) {
                    workerInfo.instance.terminate();
                }
            }
            console.log('Worker Pool closed.');
        }
    }
    
    module.exports = WorkerPool;
    

    该任务池使用的通用工作线程会为接收到的每条消息执行一个依赖 CPU 的函数,然后返回statuspayload作为响应。这里的求和循环代表了诸如图像处理或加密之类的实际工作负载。

    // pool-worker.js
    const { parentPort } = require('worker_threads');
    
    function performHeavyCalculation(data) {
        // Example heavy calculation: sum of numbers up to 'limit'
        // This could be anything CPU-bound: image processing, encryption, etc.
        let sum = 0;
        for (let i = 0; i <= data.limit; i++) {
            sum += i;
        }
        return sum;
    }
    
    if (parentPort) {
        parentPort.on('message', (taskData) => {
            try {
                const result = performHeavyCalculation(taskData);
                parentPort.postMessage({ status: 'success', payload: result });
            } catch (error) {
                parentPort.postMessage({ status: 'error', payload: error.message });
            }
        });
    }
    

    最后,主脚本根据os.cpus()返回的 CPU 核心数来设定任务池的大小,提交六项任务,并使用Promise.allSettled等待所有任务完成。每个任务的承诺对象都会被转换为易于理解的成功或失败字符串。

    // main-pool.js
    const WorkerPool = require('./workerPool');
    const path = require('path');
    
    const numCores = require('os').cpus().length;
    const pool = new WorkerPool(path.join(__dirname, 'pool-worker.js'), numCores);
    
    async function runExample() {
        const tasks = [
            { limit: 1_000_000_000 },
            { limit: 500_000_000 },
            { limit: 1_500_000_000 },
            { limit: 750_000_000 },
            { limit: 2_000_000_000 },
            { limit: 250_000_000 }
        ];
    
        console.log('Main: Submitting tasks to the worker pool...');
        const results = await Promise.allSettled(tasks.map((task, index) =>
            pool.runTask(task)
                .then(res => `Task ${index} completed with result: ${res}`)
                .catch(err => `Task ${index} failed: ${err.message}`)
        ));
    
        results.forEach(res => console.log(res.value));
    
        // Demonstrate main thread responsiveness
        console.log('Main: Tasks submitted. Doing other stuff...');
        await new Promise(resolve => setTimeout(resolve, 100)); // Simulate async work
        console.log('Main: Other stuff done.');
    
        pool.close();
        console.log('Main: Worker pool closed.');
    }
    
    runExample();
    

    请将此任务池严格视为教学示例;在类似系统投入实际使用之前,它还存在若干亟需解决的缺陷。

    • 池的message处理程序会忽略status的值,直接用payload来满足承诺,因此在工作线程中失败的任务仍会被视为成功。应检查status值,若为'error'则拒绝处理。
    • worker.terminate()会令工作线程以非零代码退出。由于exit处理程序会调用terminateWorker来生成新的工作线程,因此调用close()只会引发新一轮的工作线程创建,而无法真正关闭整个池。此外,error事件及其后续的exit事件可能导致同一个位置被重复占用两次,其ID也会被多次加入空闲列表。一个真正的池需要一个“关闭”标志,并且只应替换那些意外崩溃的工作线程。
  • 当任务的处理线程崩溃时,只有通过error路径才能拒绝该任务;如果没有error事件而发生异常退出,其承诺状态将保持待处理状态。
  • 在演示中,总和超过20亿时会超出Number.MAX_SAFE_INTEGER的范围,因此输出的结果会失去精度。如果需要精确值,请使用BigInt
  • “正在执行其他操作”的日志是在所有任务都完成等待之后才输出的,所以和第一个示例一样,它无法体现任务的并发性。
  • 更完善的线程池还会添加空闲超时机制、动态扩展功能以及更完善的恢复策略。像Piscina这样维护良好的库已经实现了这些功能,通常比自行实现的线程池更佳。

    何时工作线程有帮助,何时没有帮助

    适用场景

    • CPU密集型任务:那些会消耗大量CPU时间的操作,比如复杂的计算、压缩、加密或图像处理。
    • 保护事件循环:任何可能使主线程阻塞超过可容忍时间的操作。

    不合适的应用场景

    • I/O密集型工作:在Node.js中,网络请求、数据库查询和文件系统操作本身就是非阻塞的。将它们封装在worker中只会增加开销而没有任何好处。
    • 微小任务:启动worker以及传递消息都需要时间。对于简单的计算,这些开销可能超过其带来的收益,这也是为何应将长期运行的worker保留在工作池中的另一个原因。
  • 共享可变状态:SharedArrayBuffer可以实现这一功能,但使用Atomics协调写入操作却极难做到正确。除非你确实需要它并了解其后果,否则建议采用消息传递方式。
  • 保持工作线程健康的实践

    1. 保持工作线程脚本的简洁。仅加载任务所需的逻辑和依赖项。将整个应用框架导入每个工作线程会增加启动时间和内存占用。
    2. 选择成本最低的数据传输方式。对于小型数据,通过postMessage进行结构化克隆即可;大型ArrayBuffer则应通过传输列表传递;而深度或复杂的对象结构最好避免,因为克隆它们速度很慢。
  • 终止工作线程。在任务完成后将其结束或让其退出,否则空闲的工作线程仍会占用内存。
  • 处理双方出现的错误。在主线程中监听errorexit信号,同时在工作线程内的操作中使用try...catch结构,以便将故障以明确的错误信息形式返回。
  • 监控资源使用情况。过多的工作线程会导致上下文切换增加和内存占用上升,反而可能使程序运行变慢。os.cpus().length是确定线程池规模的一个合理起点;在较新的Node.js版本中,推荐使用os.availableParallelism()来获取该数值。
  • 不要阻塞 Worker 自身的事件循环。Worker 也有自己的事件循环。如果其 message 处理函数执行耗时的同步任务,后续的消息就会在后面排队等待处理。需要同时处理大量请求的 Worker 应尽量让每项任务的时间较短,或者在极少数情况下将任务委托给子 Worker 处理。
  • 常让人犯错的误区

    • 误以为可以共享全局变量。传递给 new Worker() 的数据只能通过 workerData(或后续的消息)才能到达 Worker。主线程中的模块级变量在 Worker 内部是不可见的,因为 Worker 在独立的执行环境中运行。真正能够实现共享的只有 SharedArrayBuffer
  • 误解了isMainThread的含义。它描述的是当前正在运行的代码的执行上下文。同一个文件既可以作为主脚本使用,也可以作为工作线程脚本,但通常分开设置主脚本和工作线程文件会更清晰。
  • 假设调试方式与平常相同。Node.js调试器可以附加到工作线程上,但你需要使用如--inspect-brk这样的标志来启动进程,并且在调试器中需要在不同线程上下文之间切换。虽然VS Code等编辑器能处理很多相关操作,但这仍然不如直接调试主线程那么便捷。
  • 让工作线程持续运行。被遗忘的工作线程会随着时间积累内存和线程资源,从而在长时间运行的服务中耗尽可用资源。
  • 关键要点

    • Worker会在同一个进程内的不同隔离环境中并行执行JavaScript代码;它们通过消息传递进行通信,仅能通过SharedArrayBuffer共享内存。
    • 应将它们用于那些会阻塞事件循环的CPU密集型任务,而不适用于Node.js已经能够异步处理的网络或磁盘操作。
    • 需为消息设定明确且一致的格式,以确保双方对结果、错误和请求ID的理解一致。
    • 应传输大型ArrayBuffer对象而非复制它们,并记住发送方在传输后会失去对该对象的访问权。
    • 可通过与机器核心数相对应的池来重用Worker,同时要确保关闭操作与崩溃恢复功能被清晰区分开,或者使用经过验证的池管理库。

    相关阅读

  • 利用并行加载函数消除 SvelteKit 中的请求瀑布效应 — 了解为何顺序的 await 会导致页面加载缓慢,以及 SvelteKit 如何通过加载函数、Promise.all、恰当的 parent() 调用和流式承诺来消除这种瀑布效应。