Ответ
Для эффективного распараллеливания CPU-интенсивных задач или обработки больших объемов данных из очереди в Node.js я использую комбинацию библиотеки bull (или agenda) для управления очередями и Worker Threads для выгрузки тяжелых вычислений из основного потока event loop.
Архитектурное решение на примере обработки изображений:
- Установка и настройка очереди (Bull с Redis):
// queue.js const Queue = require('bull'); const imageProcessingQueue = new Queue('image-processing', { redis: { port: 6379, host: 'redis' }, defaultJobOptions: { attempts: 3, backoff: 5000 } // Повторные попытки });
module.exports = imageProcessingQueue;
2. **Добавление задач в очередь (например, из Express-роута):**
```javascript
app.post('/upload', async (req, res) => {
const imageUrls = req.body.urls; // Массив URL изображений
const jobs = imageUrls.map(url => ({
data: { imageUrl: url, options: req.body.options },
}));
await imageProcessingQueue.addBulk(jobs); // Пакетное добавление
res.json({ message: `${jobs.length} задач добавлено в очередь` });
});
- Создание обработчика с Worker Threads для CPU-интенсивных операций:
// worker-processor.js const { Worker, isMainThread, parentPort } = require('worker_threads');
if (!isMainThread) { // Код, выполняемый в воркере const { processImage } = require('./heavy-image-lib'); parentPort.on('message', async (jobData) => { try { const result = await processImage(jobData.imageUrl, jobData.options); parentPort.postMessage({ success: true, result }); } catch (error) { parentPort.postMessage({ success: false, error: error.message }); } }); }
module.exports = (jobData) => { return new Promise((resolve, reject) => { const worker = new Worker(__filename); worker.on('message', resolve); worker.on('error', reject); worker.postMessage(jobData); }); };
4. **Запуск параллельных обработчиков очереди:**
```javascript
// worker.js
const imageProcessingQueue = require('./queue');
const processWithWorker = require('./worker-processor');
// Запускаем 4 параллельных воркера (по числу ядер CPU)
imageProcessingQueue.process(4, async (job) => {
console.log(`Обработка задачи ${job.id}`);
// Делегируем тяжелую работу в отдельный поток
return await processWithWorker(job.data);
});
Ключевые преимущества такого подхода:
- Устойчивость: Задачи сохраняются в Redis и не теряются при перезапуске воркеров.
- Контроль нагрузки: Легко регулировать количество параллельных обработчиков (
queue.process(N)). - Мониторинг: Bull предоставляет UI или API для отслеживания выполнения задач.
- Защита Event Loop: CPU-интенсивный код выполняется в Worker Threads, не блокируя основной поток.