Node.js worker_threads 메모리 공유 문제
Understanding SharedArrayBuffer, Atomics, and message passing patterns in Node.js worker threads for concurrent processing.
Node.js worker_threads 메모리 공유 문제
Introduction
Node.js worker threads provide true parallel execution, but sharing memory between threads introduces complex synchronization challenges. After implementing a data processing pipeline using worker threads, I encountered subtle race conditions and memory visibility issues that took significant debugging effort to resolve.
Environment
node --version
# v20.11.0
# System with multiple CPU cores
nproc
# 8Problem
A worker thread implementation for processing financial data was producing incorrect results:
// main.js
const { Worker } = require('worker_threads');
const sharedBuffer = new SharedArrayBuffer(1024);
const sharedArray = new Int32Array(sharedBuffer);
// Initialize data
for (let i = 0; i < 100; i++) {
Atomics.store(sharedArray, i, i * 10);
}
// Spawn worker
const worker = new Worker('./worker.js', {
workerData: { sharedBuffer }
});
worker.on('message', (result) => {
console.log('Sum from worker:', result.sum);
// Expected: 49500
// Getting: 24750 or other random values
});// worker.js
const { parentPort, workerData } = require('worker_threads');
const { sharedBuffer } = workerData;
const arr = new Int32Array(sharedBuffer);
let sum = 0;
for (let i = 0; i < 100; i++) {
sum += Atomics.load(arr, i);
}
parentPort.postMessage({ sum });The sums were inconsistent across runs because of memory visibility issues between threads.
Analysis
The problem is that without proper synchronization, changes made by one thread may not be visible to another thread. The Atomics module provides the synchronization primitives needed:
Atomics.store()- Write a valueAtomics.load()- Read a valueAtomics.add()- Atomic additionAtomics.compareExchange()- CAS operationAtomics.wait()/Atomics.notify()- Thread synchronization
The race condition occurs when:
Thread A: stores arr[0] = 10
Thread B: reads arr[0] → may see 0 (stale value)Solution
Solution 1: Use Atomics for all shared memory operations
// worker.js - Correct atomic operations
const { parentPort, workerData } = require('worker_threads');
const { sharedBuffer } = workerData;
const arr = new Int32Array(sharedBuffer);
let sum = 0;
for (let i = 0; i < 100; i++) {
// Atomics.load ensures we read the latest value
sum += Atomics.load(arr, i);
}
parentPort.postMessage({ sum });Solution 2: Message passing instead of shared memory
// main.js - Safer approach using message passing
const { Worker } = require('worker_threads');
const data = Array.from({ length: 100 }, (_, i) => i * 10);
const worker = new Worker(`
const { parentPort } = require('worker_threads');
parentPort.on('message', (data) => {
const sum = data.reduce((a, b) => a + b, 0);
parentPort.postMessage({ sum });
});
`, { eval: true });
worker.postMessage(data);
worker.on('message', (result) => {
console.log('Sum:', result.sum); // Always correct: 49500
});Solution 3: Worker pool for parallel processing
// worker-pool.js
const { Worker } = require('worker_threads');
const os = require('os');
class WorkerPool {
constructor(workerScript, poolSize = os.cpus().length) {
this.workers = [];
this.queue = [];
for (let i = 0; i < poolSize; i++) {
this.addWorker(workerScript);
}
}
addWorker(script) {
const worker = new Worker(script);
worker.busy = false;
worker.on('message', (result) => {
worker.busy = false;
worker.resolve(result);
this.processQueue();
});
this.workers.push(worker);
}
async execute(data) {
return new Promise((resolve, reject) => {
const availableWorker = this.workers.find(w => !w.busy);
if (availableWorker) {
availableWorker.busy = true;
availableWorker.resolve = resolve;
availableWorker.postMessage(data);
} else {
this.queue.push({ data, resolve });
}
});
}
processQueue() {
if (this.queue.length === 0) return;
const availableWorker = this.workers.find(w => !w.busy);
if (!availableWorker) return;
const { data, resolve } = this.queue.shift();
availableWorker.busy = true;
availableWorker.resolve = resolve;
availableWorker.postMessage(data);
}
}
module.exports = WorkerPool;Solution 4: SharedArrayBuffer with proper initialization
// main.js - Correct initialization
const { Worker } = require('worker_threads');
const sharedBuffer = new SharedArrayBuffer(1024);
const sharedArray = new Int32Array(sharedBuffer);
// Initialize using Atomics.store for visibility
for (let i = 0; i < 100; i++) {
Atomics.store(sharedArray, i, i * 10);
}
// Signal workers that data is ready
Atomics.store(sharedArray, 100, 1); // Ready flag
const worker = new Worker('./worker.js', {
workerData: { sharedBuffer }
});// worker.js - Wait for data readiness
const { parentPort, workerData } = require('worker_threads');
const arr = new Int32Array(workerData.sharedBuffer);
// Wait for ready signal
Atomics.wait(arr, 100, 0);
let sum = 0;
for (let i = 0; i < 100; i++) {
sum += Atomics.load(arr, i);
}
parentPort.postMessage({ sum });Lessons Learned
- Prefer message passing over shared memory - It is simpler and less error-prone
- Use Atomics for all shared memory access - Never read/write directly without atomics
- Implement proper synchronization - Use wait/notify for coordination between threads
- Consider using a worker pool - Managing thread lifecycle manually is complex
- Profile memory usage - SharedArrayBuffer counts against all threads' memory limits
This blog does not accept any external sponsorships, affiliate marketing, or ad revenue.