Галоўная / Артыкулы / Выкантовка роботы, яка сильна на CPU, у Node.js за дапамою worker_threads і Pools

Выкантовка роботы, яка сильна на CPU, у Node.js за дапамою worker_threads і Pools

Дазвольце дазнацься, як worker_threads у Node.js падтрымліяе чутлівасць цыклу здарэнняў: стварэнне працоўнікаў, абмен рашчынаннямі-запытамі, перадавальныя буферы, пулі і ўсе проблемы, якія з’яўляюцца па ўзмацненні.

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, які паказаны нижэй.) Паўстаноўкі з процэсамі, стварэнне і вялічанне з рабочымі элементамі є дышэўней, што робіць іх натуральным спосабам выкарыстоўвання колькох ядер у аднай інстанцыі прыкладнага програму.

Што коштае вас заблокаваны цыкл змагання

Галоўная задача працоўніка — падтрымліваць вольны цыкл запуску. Калі на главнай ніті ведзецца сінхронны вычыслення, сервер не можа прыймаць новыя запиты, не можа працаваць над відкрытымі з’ѐеднаннямі і нават не можа адпавядаць на перагляд стану. Практычныя наследкі выглядаюць так:

  • Падышанне якосці адчування корыстніка: адпаведзі AI прыходзяць пазней, а функцыі у рэальны час, такія як жывыя апдэты, працюють з перашкодамі.
  • Знижанне праўоўдзімасці: адна дліга задача займае ўсюедыну ніту JavaScript, таму колькасць запитоў за секунду зменшыцца.
  • Нестабільнасць: дзяўнага блакування можа перасягнуць таймауты балансіра навантажэння або кліента, і гэтыя неудачы могу распрастарыцца на іншыя службы.

Перанесенне важкага вычыслення на працоўніка памогае галоўнаму цыклу залишацца вольным для обслужвання сокетаў і таймераў, такім чынам аплікацыя застаецца рэагуючай нават калі ведзецца складнае вычыслення.

Ёсць падробнейшая інфармацыя пра тое, як супрацоўнае цыклі, libuv і ўсё яго пулы задачаў супрацоўнае адносны, — пагляньце на тое, як фактычна працюе паралельная обработка ў Node.js.

Іспытанне всіх ядоў

Серверы зазвычай маюць многа ядоў CPU. Без працоўнікаў, або без запуску кальколькі процэсаў через child_process чы ўправалювача процэсаў, такога як PM2, адзін процэс Node.js можа выкарыстоўваць прыблізна адны ядро для виконання коду на JavaScript, які залежыць ад моці CPU. Рэшта ядоў застаюцца незайнятымі. Працоўнікі дазволяюць адной аплікацыі паралельна виконваць кальколькі тэжкіх задач на всіх доступных ядоў, чым падвышаецца абщая праўоцэнная продуктывасць.

Задачы, якія стаюць практычнымі

Паралельная обработка адкрывае можлівасці ў сферах, якія ранейш выкарыстоўваліся толькі мовамі з першасортным падтрымкам задач:

  • Обробка дадзейнаў: аналіз у рэальны час, маніпуляцыя зображэннямі і відэаматеріяламі, а таксама трансфармацыя большых набораў дадзейнаў.
  • Выводы з машыннага навучэння: адклёкванне запрацоўваўных модэляў, калі рэшта прыемлівае часткі програмы застаецца інтэрактыўной.
  • Крыптаграфія: хэшаванне, шифраванне і дэшифраванне.
  • Сімуляцыя і научны вычысленні: паралельныя сімуляцыі і складныя матэматычныя модэлі.

Ключовы момант заключаецца ў тым, што нічга з гэтага не пазбавляе вас асінхроннага модэлю I/O. Рабочыя задачы дадаюць паралельныя вычысленні на яго аднойчыне.

Кампаненты модуля

worker_threads адкрывае невялікі набор прымітіў:

  • Worker: клас, які выкарыстоўваецца галоўным потокам для запуску новай рабочай задачы з скрыпты.
  • isMainThread: булева значэнне, якое паведамляе, чы рэальны код выконваецца на галоўнай ніті; ёна корыстная, калі адной файл магчыма выплываць обе ролі.
  • parentPort: доступна ўнутры рабочага процэсу, яна ёсьць канал назад да ніты, якая яго створыла. На галоўнай ніті яна равна null.
  • workerData: доступна ўнутры рабочага процэсу, яна зберагчыца клон усіх дадзеных, якія абмовіцель стварыў рабочы процэс.
  • MessagePort і MessageChannel: выкорыстоўваюцца для стварэння дадатковых незалежных двухканальных каналаў між нітямі.
  • SharedArrayBuffer і Atomics: справжня спільная памяць плюс прымітывы сінхронізаціі, неабходныя для ўпэўненага яе выкарыстоўвання. Це найэфектывнейшы варыянт, але таксама і найболей складны, таму што ў яму знова выступаюць умовы канкурэнціі.
  • Першы працоўнік: вычысленне простых чисоў паza галоўным потокам

    Хорашым першым завданнем є навмысна сповольнены тэлеграф: вычысленне всіх простых чисоў да вялікага ліміту. У прыкладзе выкарыстоўваюцься два файлы, якія знаходзяцца ў той самай папকе.

    Скрыпт галоўнага потока

    Галоўны скрыпт стварае працоўнік, даўае яму даны і чакае на адпаведзь. Перш чым яго прачытаць, зверніце увагу на тры моменты. Працоўнік ствараецца з адной з іншых дапаможных файлавых шляхоў. Даны надаюцца через опцыю 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 задаецца на дастаткова вялікую значэнне (20 мільйонаў), так што вырахунак займае значны час.
  • 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.");
    }
    

    Што звернуць увагу:

    • У яе запрашваюцца parentPort і workerData з модуля.
    • Умова if (parentPort) гарантуе, што логіка будзе выконвалася толькі тады, калі файл завантажваецца як рабочы працоўнік; якщо вы ўвядзеце яе безпосередна за дапамогай node, parentPort будзе null, і тады будзе выведана паведамленне.
    • const { limit } = workerData; чытае данні, якія адправіў родны працоўнік.
    • parentPort.postMessage(result) адправляе простыя числы на галоўную нітку, дзе працоўнік-обработчык message іх прыме.
  • Блок try...catch фіксуе кожную памятку падення падзе час вычыслаў і дапавядае яе роднім працэсу, замест таго каб адпаведны потак заканчыў сваю роботу.
  • Адна тонкасць: у разе памятакі працоўнік падае { error: ... } у той самы канал message, які викорыстоўваецца для рэзультатаў, але галоўны скрыпт спрацьвоўвае кожную паведамленне як масэвы элемент з прымовымі числамі. У рэальным кодзе неабходна надаць паведамленням чысткі формат (напрыклад, поле status або type), каб родні працэс могў разлічыць рэзультат ад памяткі падення.

    Запуск і чытанне выходных дадзенняў

    Зберажыце оба файлы падле аднаго іншага і запусціце програму за дапамой node main.js. Вы пабачыце, як галоўны потак анансуе сваю роботу, працоўнік паведаміць пра свой запуск, а чыраз калькі часу — рэзультат працоўніка разам з перыядам часу, які прыйшоў.

    Зверніце увагу на порядак выведання рэчы «Галоўны поток продаўжае выкананне іншых задач...». У гэтым скрыпце яна выводзіцца ўнутрь обрабоўчыка message, таму яна з’являецца толькі пасля таго, як працоўнік завершыў свою роботу. Це паказвае, што галоўны поток змог даць прыем і обробіць рэзультат, а не тое, што ён выканаў іншую роботу пры гэтым. Ёсць можлівасць побачыць, як галоўны поток застаецца рэагіруючым пад час вычыслаў, якщо запусціць якую-небудзь дзеянне ў галоўным потоку да таго, калі прыйдзе рэзультат, напрыклад setInterval, який фіксуе стан потока кожныя кальколькі сотых мілісекунд. Стан потока будзе фіксавацца, пакуль працоўнік выканае вычыслення, што абоцу не будзе відбывацца, якшо findPrimes будзе выкананы ў галоўным потоку.

    Двунаправленая пераказваўка з протакалам запыт-адказ

    Першы прыклад адправляе адна вхідная даны і отрымае адны рэзультат. Большасць рэальных застосоўванняяў выклікаюць патрэбу ў працоўніку, які дзейнае трывалы час і обрабоцвае многа запитоў. Комунікацыя ўсё цэла двухнаправленая: об’ект Worker у галоўнай ніті і parentPort унутры працоўніка абодва задаюць .postMessage(...) і .on('message', ...).

    Працоўнік, показаны нижэй, застаецца жывым і чакае на паведамленні. Калі ён атрымле запит calculateFibonacci, ён вырахоўвае значэнне за дапамой простай рекурзіўной функцыі (навмысна медленна) і адпавядае або паведамленнем fibonacciResult, або fibonacciError, падаючы ідэнтыфікатор запиту, ўпорам з тым, каб вызываючы адэнт могаў паспрацаваць з адпаведзю.

    // 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, зберагае функціі resolve і reject абяцанні ў Map пад гэтым ID і адправляе запит. Калі прыходзіць адпаведзь, працоўнік message шукае ID, вычысляе відпаведную ентрэту і задовольняе абяцанне. У прыкладзе для стварэння ID выкарыстоўваецца пакет uuid; у ныях версіях Node.js функцыя crypto.randomUUID() з вбудованага модуля node:crypto выконвае тую ж задачу без додатковых залежнасцяў.

    // 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 зберагае функцыі resolve і reject пакуль працоўнік не дае адпаведзі.
  • Працоўнік обробляе паведамленні calculateFibonacci і адпавядае з fibonacciResult або fibonacciError, пры чым перадае первісны requestId.
  • Якщо вы адаптаваеце гэты код, atsargавайцеся з тым, як формуецца паведамленне. Галоўны адбег ставіе requestId на верхні ўрадзень паведамлення, рыжыячы з payload, тады калі працоўнік распаковвае яго з message.payload. У такім варыянце працоўнік бачыць requestId як undefined, адпаведна не можа знайсці адпаведны ўтвар, і обяўленая праміс ніколі не будзе выпанавана. Зберагайце ID у аднам і тым жа месцаў на обох сторонах (напрыклад, payload: { n, requestId }). Таксама варта дадаць таймаут для кожнага чакаючага запиту, ўпэўніваючыся, што загублены адпаведны ўтвар будзе спрацоўваць як адказ на адпаведную памылку, а не будзе спрычыняць застоў. Нарэшце, зверніце увагу на блок finally: калі робота завершылася, worker.terminate() зупініць адбег, і процес сможа завершыцца.

    Перадача вялікых дадзенняў без ўсунення іх

    Усе, што адправляецца через postMessage, за замовчанням клонуецца у структураваны формат, тое значыць, ў прыёмніка практычна атрымліваецца копія. Для маленькых паведамленняў гэта нічога страшнага. Але для вялікіх бінарных дадзеных сама копія можа стаць прычынай спаду выдатков. Рашэнняя — пераказаць права на владу базовай памяцю замест адтворэння яе: об’ект стае непрыдатным для адправцы і негайна прыдатным для прыёмніка, без жадных дублікацый.

    ArrayBuffer і MessagePort — гэта аб’екты, якія найчащэй пераказваюцца. SharedArrayBuffer — гэта адзінаковы случай: яго ўжо не пераказваюць, таму што обе задачы можа адразу доступвацца да той самай памяці.

    У прыкладе нижэй знаходзяцца два скрыпты ў аднай спісаванні. Рабочы скрыпт (buffer-worker.js) прымеў буфер, удвойвае кожны байт на месцы і адправляе буфер знова як об’ект для перадачы. Галоўны скрыпт (main-buffer.js) заполняе буфер размерам 1 МБ, перадае яго рабочаму скрыпту і чытае обробленыя даны, калі той вярнуецца.

    // 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 МБ даных ніколі не копіюецца ні в однам, ні ў другым напрамку.
  • Адносны недзея — пасля перадачы неабходна прыпыніць викорыстоўвання первіснага канкэрента. Якщо оба потакі справды патрабуюць даных адночасна, трэба викорыстоўваць копію або SharedArrayBuffer.

    Раскрытка адказоў на адказы і чыставанне

    Рабочыя процесы можу злучыцца, як і будзь-які іншы код, і пакуль яны працуюць, яны займаюць памяць і потак. Доступныя інструменты:

    • worker.on('error', handler) у главным потаку: прымае адказы, якія былі выклеканы ў рабочым процесе і не падбіраны.
  • worker.on('exit', handler) на галоўнай вясі: выклікаецца кожны раз, калі працоўнік зупыняецца, незалежна ад таго, чы справілася ён нормальна чы зламаўся. Код выходу разлічвае гэтыя два варыянты, пры чым 0 абазаваны на успеху.
  • process.on('uncaughtException', handler) усередзіне працоўніка: даўае можлівасць фіксаваць чы выкарыстоўваць дакладную інфармацыю з усередзіння працоўніка першым, чым ён зупыняецца, паўтараючы тое, што бачыць родны процес.
  • worker.terminate(): зупыняе працоўнік якомога шырэй і вяртае прабаму, якая рашаецца з кодам выходу. Ён не чакае, пакількі завершыцца выконванне задач, таму, якщо вам патрэбна грацыязная зупінка, надаюце працоўніку паведамленне "stop" і дазволяеце яму самаму выйсці. Заўжды старайцеся, каб працоўнікі, якія больш не патрэбны, былі зупыненыя якім-небудзь спосабам.
  • Выкананне больш заўданняў через пул працоўнікаў

    Стварэнне новага працоўніка для кожнай задачы є марнасцю, а запуск большага колічынства працоўнікаў, якія сильна выкарыстоўваюць процэсор, чым є ядры, толькі паўяршвае ситуацыю з конкурэнціей. Аблака працоўнікаў рашаюць оба гэтыя проблемы: яны запускаюць фіксаванае колькіцтва працоўнікаў, кладуць новыя задачы у чергу і перадаюць кожную з яных наступнаму вільнаму працоўніку.

    У спроставаным аблаку працоўнікаў, показаным нижэй, зберагаецца масіва з рэкордамі працоўнікаў, спіс 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, для кожнага прыйнятага паведамлення і адпавядае з status і payload. Цикл сумавання выступае адпраўнікам для будь-якай рэальной навантажэння, такой як обробка зяўроў або шифраванне.

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

    У канцы галоўны скрыпт насталявае розмер пулу на колькість ядраў CPU, якія вярнюе os.cpus(), адправляе шасць заведамаў і чакае на ўсія з іх за дапамогою 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 разв’язвае запавет з payload, незалежна ад status; таму задача, яка зазнала нявезення ў рамках працэйніка, все роўна лічыцца як успех. Пераканаўся ў status і адхіліў задачу, якшо ёй 'error'.
    • worker.terminate() выклікае завершэнне працэйніка з кодам, не равным нулю. Адтако адміністратар exit выклікае terminateWorker, який стварае новага працэйніка, таму вызыв close() спрычынае цыкл стварэння новых працэйнікаў у замяне на завершэнне пула. Кроме таго, error у поўнасці з ўзрошаным імпартаментам exit можа два разы запрацаваць для таго ж месца і колькаразова дадаць яго ID у список вільных працэйнікаў. Настоямы пул патрабуе флага „закрыцця“ і должен заменяць толькі тыя працэйнікі, якія неспадзевана зламаліся.
  • Заданне, якое было ў процесе выконання, калі зламаўся яго працоўнік, адхіляецца толькі через шлях error; ненормальны выход без змагання error заставляе яго абяву ў стане очаквання.
  • У дэманстрацыі сумы, якія досягаюць двух мільярдаў, перавышаюць Number.MAX_SAFE_INTEGER, таму выдрукаваныя рэзультаты втрачаюць точнасць. Якщо важны точныя значэння, выкарыстоўваюце BigInt.
  • Лог „Doing other stuff“ запускаецца пасля таго, як усі заданні будуць завершаны, таму, як і ў першым прыкладзе, ён не паказвае паралельнасці з заданнямі.
  • Болей завершаныя пулі таксама дадаюць таймауты для бездзейнальных працоўнікаў, дынамічнае масштабаванне і болей рэтельную вяснавання. Якісна падтрымваныя бібліятэкі, такія як Piscina, реалізуюць гэтыя деталі, і яны зазвычай ўжо кращы выбар, чым самастоятельна створаны пул.

    Калі працоўнікі-адрасы дапамагаюць і калі ні

    Хорашы варыянты

    • Задачы, які зависяць ад CPU: все тое, што спрацоўвае значную колькасць часу CPU, такія як складныя вычысленні, стисненне, шифраванне чы адобработка зяўроцэў.
    • Захаванне цыклу здарэнняў: будзь-яя операцыя, якая інакша могла бы заблокаваць галоўную ніт на перыяд, які не можна прымусіць.

    Няпасуючы варыянты

    • Робота, якая завісіць ад I/O: запыты сеті, запыты базы дадзенаў і операцыі сістэмы файлаў у Node.js ўжо являюцься неблакуючымі. Їх пакрыцце рабочым процэсам прыносіць толькі дадатковую навантажэння без жадных выгод.
    • Мілкія задачы: запуск рабочага процэсу і перадача паведамленняў або тое, або іншае спрацоўваюць за час. Для швайных вычысленняў такая дадатковая навантажэнне можа перакрыць практычныя выгоды, што ўжо адной з прычын трэба трывалыя рабочыя процэсы зберагаць у пуле.
  • Адмініструемы стан, які можна зменшваць: SharedArrayBuffer дазваляе гэта зрабіць, але скоординаванне запісаў за дапамою Atomics є надзвычайна складным. Якщо вы сапраўды не патрабуеце гэтага і не розумееце наследків, краща викорыстоўваць перадачу паведамленняў.
  • Практыкі, якія падтрымліваюць эфектывнае працаванне воркероў

    1. Робіце скрыпты воркероў маленькімі. Завантажвайце толькі логіку і залежнасці, якія патрабуе задача. Імпортаванне всіяго фрэймворку вашай аплікацыі ў кожны воркер паўзначае збільшэнне часу запуску і викорыстоўвання памяці.
    2. Выбірайце найдешавшы спосаб перадачы дадзенняў. Структураванае клонаванне через postMessage падходзіць для маленькіх паведамленняў; вялікія ArrayBuffer-ы трэба перадаваць через спецыяльны список перадач, а глыбокія чы роўназначныя графы об’ектаў краща утрымваць, таму што ўсё гэта сповалічвае процес клонавання.
  • Зупніце працоўнікамі работу. Заканчыце ўжотыя працоўнікі або дазвольце ім выйсці, калі ўжо завершылася ўсё іхняе заданне. Бездзейны працоўнікі продавжвають займать памяць.
  • Керуйце адмовамі з обох бакоў. Следзіце за сигналамі error і exit у галоўнай ніті, а таксама абгорніце работу працоўніка ў блок try...catch, каб адмовы вярталіся у вигляде чыстых паведамленняў пра адмову.
  • Стежыце за викорыстоўваннем рэсурсаў. Занадта много працоўнікаў спрычынае зміну контекста і вышэйшы расход памяці, што можа зробіць працэсы медленнейшымі, а не шырэйшымі. os.cpus().length — граматычны пачатковы параметр для такога пула; у новейшых версіях Node.js рэкамендуецца вжываць os.availableParallelism() для атрыбутавання гэтага значэння.
  • Не блакуйце сабеяшні цыкл рабочага. У рабочага таксама ёсць цыкл здарэнняў. Якшо яго працоўнік з обробкай message выканае дзейсненне, якое трывае дазго і ў сінхронным режыме, наступныя паведамленні заставаюцца ў черзе пасля яго. Рабочы, які должен адрабатваць многа запросоў, павінен робіць кожную частку работы короткай, альбо, у рэдкіх случаях, перадаваць яе своім падрабочым.
  • Памылкі, якія часта ствараюць проблемы

    • Аспакоўваннеся спяльнымі глобальнымі зменнымі. Даныя, якія падаюць у new Worker(), доходзяць да рабочага толькі через workerData (альбо пазнейшыя паведамленні). Зменны на рэвэле модуля ў галоўнай ніткі не ёсць видныя ўнутры рабочага, таму што той выканаецца у сваім абсалютна ізольаваным сераўере. Ядзіныя справжнія способы раздзелення дадзеных — через SharedArrayBuffer.
  • Неправыя выкладкі ўчынка isMainThread. Ён апісвае контекст выконання коду, який зараз працюе. Адны файл можа выкорыстоўвацца як галоўны скрыпт, так і як рабочы процес, але окольныя галоўныя та рабочыя файлы зазвычай є болей зрозумелымі.
  • Прыпуск, што адлагоджэнне працюе як зазвычай. Інспектар Node.js можа прабачвацца да рабочых нитэй, але процес пачынаецца з флагамі інспектара, такімі як --inspect-brk, і трэба перыходзіць между контекстамі нитэй у адлагоджвачы. Редактары, такія як VS Code, адчыняюць вялікую частку ціх задач, пры тым як гэта все жа менш безпосередняя форма, чым адлагоджэнне галоўной ніті.
  • Занятак рабочых процесаў. Забутыя рабочыя процесы з часам накапліваюць память та ніті і можу выкарыстоўваць усе ресурсы ў дужа дзейнічых сервісах.
  • Ключовыя выводы

    • Рабочыя процесы запускаюць JavaScript паралельна ў адзінаковых ізолятных сэрвісах у межах адного процесу; яны вырашаюць пытанні взаімнае абмену інфармацыяючы праз перадачу зместоў і дзеліцца памяцю толькі через SharedArrayBuffer.
    • Їх трэба выкарыстоўваць для задач, якія залежаць ад працы CPU і якія інакш заставілі бы на пазі цыкл адбывання запускаемых задач, а не для доступу да сеті чы ўстройства зберагання, якія Node.js вже обрабоцвае асінхронна.
    • Неабходна, каб зместы паведамленняяў мелі чыста визначаную структуру, каб рэзультаты, адказы пра памялкі і ID запитоў былі одназначныя для обох сторон.
    • Кращэ перадаваць вялікія об’екты ArrayBuffer, а не клонаваць іх, і памятаць, што апарат, які ўпрабаваў іх перадацю, пасля цього тыятроўвае да яных доступ.
    • Рабочыя процесы трэба перадаўаць знову праз пул, які складаецца з колькасці елементаў, роўнай колькасці ядраў апарата, і трэба чытка раздзеляць процесы выключэння і вяснавання пасля збою, але можна таксама выкарыстоўваць перавярэную бібліятэку для керування пулам.

    Супаўзеяныя матэрыялы

  • Заканчэнне процеса обработкі заявок у SvelteKit за дапамою функцый паралельнага завантажэння — Дазвольце дазнацца, чаму последовыя вызывы awaits спрычыняюць медленную роботу сторанак, і як SvelteKit, функцыі load, Promise.all, а таксама рэштыльныя вызывы parent() і стрімаваныя promises пазбегаюць гэтага явы.