Node.js: event loop bị nghẽn trông như thế nào?
Câu hỏi: một request CPU có làm request I/O chậm theo, và chuyển sang worker phải trả chi phí gì?
Cần biết trước: JavaScript Promise, HTTP và process. Lab dùng Node.js 24.18.0, chỉ built-in modules, dữ liệu giả và HTTP loopback. Client chạy ở process riêng với server; đồng hồ client không bị CPU server chặn. Chưa đo production hoặc mạng xa.
Chốt phép so
Mỗi lượt cùng 4 CPU job, mỗi job 60.000.000 vòng tính integer; hai job gửi đồng thời, hai wave. Trong mỗi wave gửi 8 request I/O, endpoint chờ 5 ms. Cả hai cách phải trả đúng mọi kết quả. Inline chạy CPU trên loop server; phương án kia có pool 2 worker tái sử dụng, không queue chờ trong pool. Khi cả hai bận, request CPU thứ ba trả 503.
Chạy 3 lượt mỗi cách, đổi thứ tự inline/worker; mỗi lượt server mới. Worker/process startup được ghi riêng, không cộng vào request latency. Không gọi đây là pool warm production hoặc benchmark throughput tối đa. Vòng CPU async/Promise vẫn có thể chặn loop nếu chạy JavaScript đồng bộ.
Worker phù hợp CPU; I/O bất đồng bộ thường không cần chuyển vào worker. Pool tránh khởi tạo một worker/request nhưng thêm messaging và lifecycle. worker_threads.
Lưu ba file trong thư mục trống
import { parentPort } from 'node:worker_threads';
parentPort.on('message', ({ id, n }) => {
let value = 0;
for (let i = 0; i < n; i++) value += (i + id) % 997;
parentPort.postMessage(value);
});
Server giữ worker capacity 2 và timer I/O của chính nó. Không có job chạy vô hạn. Job số 93 dành riêng cho phép thử timeout, chạy 600.000.000 vòng để quan sát rõ worker vẫn bận sau khi client ngắt; job này nằm ngoài mẫu benchmark. Shutdown đóng endpoint và worker thuộc fixture; không tác động service khác.
import assert from 'node:assert/strict';
import http from 'node:http';
import { monitorEventLoopDelay, performance } from 'node:perf_hooks';
import { Worker } from 'node:worker_threads';
assert.equal(process.version, 'v24.18.0');
const mode = process.argv[2];
assert.ok(['inline', 'worker'].includes(mode));
const slots = [];
const startup = performance.now();
if (mode === 'worker') {
for (let i = 0; i < 2; i++) {
const worker = new Worker(new URL('./worker.mjs', import.meta.url));
const slot = { worker, job: null, gone: false };
worker.on('message', (value) => {
const job = slot.job;
slot.job = null;
if (job) job.resolve(value);
});
const fail = (error) => {
slot.gone = true;
if (slot.job) slot.job.reject(error);
slot.job = null;
};
worker.on('error', fail);
worker.on('exit', () => fail(new Error('worker exit')));
slots.push(slot);
}
await Promise.all(slots.map(({ worker }) => new Promise((resolve, reject) => {
worker.once('online', resolve);
worker.once('error', reject);
})));
}
const workerStartupMs = performance.now() - startup;
const delay = monitorEventLoopDelay({ resolution: 10 });
delay.enable();
const timers = new Map();
let stopping = false;
const active = () => slots.filter((slot) => slot.job !== null).length;
const cpu = (id, n) => {
if (mode === 'inline') {
let value = 0;
for (let i = 0; i < n; i++) value += (i + id) % 997;
return Promise.resolve(value);
}
const slot = slots.find((item) => !item.job && !item.gone);
if (!slot) throw new Error('busy');
return new Promise((resolve, reject) => {
slot.job = { resolve, reject };
slot.worker.postMessage({ id, n });
});
};
const server = http.createServer(async (req, res) => {
const url = new URL(req.url, 'http://127.0.0.1');
const send = (status, body) => {
if (!res.destroyed) {
res.writeHead(status, { 'Content-Type': 'application/json' });
res.end(JSON.stringify(body));
}
};
if (stopping) return send(503, { error: 'stopping' });
if (url.pathname === '/state') return send(200, { active: active() });
if (url.pathname === '/metrics') return send(200, {
meanMs: delay.mean / 1e6, maxMs: delay.max / 1e6,
p99Ms: delay.percentile(99) / 1e6, samples: delay.count, active: active(),
});
if (url.pathname === '/io') {
await new Promise((resolve) => {
const timer = setTimeout(() => { timers.delete(timer); resolve(); }, 5);
timers.set(timer, resolve);
});
return send(200, { ok: true });
}
if (url.pathname !== '/cpu') return send(404, { error: 'unknown route' });
const id = Number(url.searchParams.get('id'));
if (!Number.isInteger(id) || id < 0 || id > 99) return send(400, { error: 'id' });
try {
send(200, { value: await cpu(id, id === 93 ? 600000000 : 60000000) });
} catch (error) {
send(error.message === 'busy' ? 503 : 500, { error: error.message });
}
});
await new Promise((resolve) => server.listen(0, '127.0.0.1', resolve));
const address = server.address();
console.log(JSON.stringify({ ready: true, port: address.port, workerStartupMs }));
process.once('SIGTERM', async () => {
stopping = true;
const closed = new Promise((resolve) => server.close(resolve));
server.closeAllConnections();
for (const [timer, resolve] of timers) { clearTimeout(timer); resolve(); }
timers.clear();
await Promise.all(slots.map((slot) => slot.worker.terminate()));
await closed;
delay.disable();
assert.equal(server.listening, false);
assert.ok(slots.every(({ worker }) => worker.threadId === -1));
console.log(JSON.stringify({ stopped: true, listening: false, workers: slots.length, timers: timers.size }));
});
Client giữ timeout hữu hạn, đối chiếu output, đo wall time từng request và tắt child qua SIGTERM. Dùng một Agent riêng mỗi lượt; không tái dùng socket của service khác.
import assert from 'node:assert/strict';
import { spawn } from 'node:child_process';
import { writeFile } from 'node:fs/promises';
import http from 'node:http';
import { availableParallelism } from 'node:os';
import { performance } from 'node:perf_hooks';
import readline from 'node:readline';
import { setTimeout as sleep } from 'node:timers/promises';
assert.equal(process.version, 'v24.18.0');
const expected = (id) => {
const sum = (n) => Math.floor(n / 997) * 997 * 996 / 2 + (n % 997) * (n % 997 - 1) / 2;
return sum(60000000 + id) - sum(id);
};
async function run(mode, repeat) {
const started = performance.now();
const child = spawn(process.execPath, ['server.mjs', mode], { stdio: ['ignore', 'pipe', 'pipe'] });
let stderr = '';
child.stderr.setEncoding('utf8').on('data', (data) => { stderr += data; });
const lines = readline.createInterface({ input: child.stdout });
const messages = [];
lines.on('line', (line) => messages.push(JSON.parse(line)));
const exited = new Promise((resolve) => child.once('close', (code) => resolve(code)));
const agent = new http.Agent({ keepAlive: true, maxSockets: 20 });
let port;
const get = (path, timeout = 5000) => new Promise((resolve, reject) => {
const begun = performance.now();
const req = http.get({ hostname: '127.0.0.1', port, path, agent }, (res) => {
let body = '';
res.setEncoding('utf8').on('data', (data) => { body += data; });
res.on('end', () => resolve({ status: res.statusCode, body: JSON.parse(body), ms: performance.now() - begun }));
res.on('error', reject);
});
req.setTimeout(timeout, () => req.destroy(new Error('deadline')));
req.on('error', reject);
});
const wait = async (predicate) => {
const deadline = performance.now() + 5000;
while (performance.now() < deadline) {
if (await predicate()) return;
assert.equal(child.exitCode, null, stderr);
await sleep(5);
}
throw new Error('barrier timeout');
};
try {
await wait(() => messages.some((message) => message.ready));
const ready = messages.find((message) => message.ready);
port = ready.port;
const startupMs = performance.now() - started;
await sleep(30);
const cpuLatencies = [];
const ioLatencies = [];
for (let wave = 0; wave < 2; wave++) {
const ids = [wave * 2, wave * 2 + 1];
const jobs = ids.map((id) => get(`/cpu?id=${id}`));
await sleep(10);
const probes = Array.from({ length: 8 }, () => get('/io'));
const cpuResults = await Promise.all(jobs);
const ioResults = await Promise.all(probes);
cpuResults.forEach((result, i) => {
assert.equal(result.status, 200);
assert.equal(result.body.value, expected(ids[i]));
cpuLatencies.push(result.ms);
});
ioResults.forEach((result) => {
assert.equal(result.status, 200);
assert.equal(result.body.ok, true);
ioLatencies.push(result.ms);
});
}
await sleep(20);
const metrics = (await get('/metrics')).body;
assert.ok(metrics.samples > 0 && metrics.maxMs > 0);
assert.equal(metrics.active, 0);
if (mode === 'worker') {
const occupied = [get('/cpu?id=90'), get('/cpu?id=91')];
await wait(async () => (await get('/state')).body.active === 2);
assert.equal((await get('/cpu?id=92')).status, 503);
const occupiedResults = await Promise.all(occupied);
assert.ok(occupiedResults.every((result) => result.status === 200));
const timed = get('/cpu?id=93', 50).then(
() => new Error('unexpected completion'), (error) => error,
);
await wait(async () => (await get('/state')).body.active === 1);
assert.match((await timed).message, /deadline/);
assert.equal((await get('/state')).body.active, 1);
await wait(async () => (await get('/state')).body.active === 0);
}
return { mode, repeat, startupMs, workerStartupMs: ready.workerStartupMs, cpuLatencies, ioLatencies, loop: metrics, overload: mode === 'worker' ? 503 : 'not tested', timeout: mode === 'worker' ? 'client deadline; finite worker finished' : 'not tested' };
} finally {
agent.destroy();
child.kill('SIGTERM');
const timer = setTimeout(() => child.kill('SIGKILL'), 5000);
const code = await exited;
clearTimeout(timer);
lines.close();
assert.equal(code, 0, stderr);
assert.ok(messages.some((message) => message.stopped && !message.listening && message.timers === 0));
assert.equal(stderr, '');
}
}
const samples = [];
for (let repeat = 1; repeat <= 3; repeat++) {
for (const mode of repeat % 2 ? ['inline', 'worker'] : ['worker', 'inline']) samples.push(await run(mode, repeat));
}
const payload = { metadata: { measuredAt: new Date().toISOString(), node: process.version, arch: process.arch, platform: process.platform, cpus: availableParallelism(), cpuIterations: 60000000, cpuJobs: 4, ioRequests: 16, ioDelayMs: 5, workerCapacity: 2, loopResolutionMs: 10, repeats: 3, startupExcludedFromRequestLatency: true }, samples };
await writeFile('eventloop-results.json', JSON.stringify(payload, null, 2) + '\n');
console.log('NODE_RESULT ' + JSON.stringify(payload));
console.log('runs=6 outputs=equal capacity=2 overload=503 timeout=checked shutdown=closed');
set -euo pipefail
node --check worker.mjs
node --check server.mjs
node --check study.mjs
node study.mjs
runs=6 outputs=equal capacity=2 overload=503 timeout=checked shutdown=closed
set -euo pipefail
test -f eventloop-results.json
Đọc request latency cùng loop delay
Đo lúc 08:13:58 UTC ngày 2026-10-03, Node.js 24.18.0, macOS arm64, 10 logical CPU.
Đây là wall time đo được trên máy này; không phải cam kết latency cho máy khác.
JSON eventloop-results.json giữ từng request và window của mỗi lượt. Bảng dùng ms,
mean CPU từ 4 request và mean/max I/O từ 16 request mỗi lượt:
| Cách | Lượt | Startup ms | CPU mean ms | I/O mean ms | I/O max ms | Loop max ms | Mẫu loop |
|---|---|---|---|---|---|---|---|
| inline | 1 | 39.150 | 64.810 | 80.171 | 82.669 | 47.907 | 9 |
| worker | 1 | 56.243 | 45.679 | 8.758 | 9.873 | 12.190 | 12 |
| worker | 2 | 48.754 | 43.005 | 8.657 | 11.902 | 12.034 | 12 |
| inline | 2 | 36.692 | 62.649 | 77.545 | 80.326 | 46.498 | 9 |
| inline | 3 | 37.670 | 62.040 | 76.170 | 79.109 | 45.842 | 9 |
| worker | 3 | 49.757 | 44.165 | 8.049 | 10.765 | 12.329 | 12 |
I/O mean giảm rõ trong cả ba lượt khi CPU sang worker. Worker startup riêng là 18.367, 11.416 và 11.538 ms; startup process cộng worker nằm ở cột Startup, được loại khỏi request latency. Loop max và request max đo hai đại lượng khác nhau: I/O có thể phải chờ nhiều công việc CPU liên tiếp trước khi callback được phục vụ. Worker CPU mean thấp hơn ở workload này; chưa đo chi phí truyền payload lớn, nhiều worker hoặc CPU contention với service khác nên không kết luận worker luôn nhanh hơn.
monitorEventLoopDelay trả nanosecond, lab đổi sang ms. Resolution 10 ms là nhịp lấy
mẫu, không phải latency request và không là CPU%. P99 histogram với ít sample không
có độ tin cậy của một production SLO. perf_hooks.
Timeout client không dừng CPU đã gửi vào worker. Lab quan sát slot bận trước và sau
socket timeout 50 ms, rồi chờ công việc hữu hạn hoàn tất. req.setTimeout đo thời gian
socket không hoạt động; không phải deadline tổng cho mọi dạng HTTP streaming.
Không trả slot về idle sớm rồi giao task khác cho worker còn bận. Pool không có queue,
503 là tín hiệu quá tải để caller giảm concurrency/retry có ngân sách, không retry vô hạn.
Shutdown gọi server.close/closeAllConnections với fixture đã kiểm riêng, terminate
worker và await child close sau khi các stream đã đóng. Đây không là production drain policy: request đang chạy có thể
bị ngắt nếu shutdown cưỡng chế. HTTPclose.
Học tiếp: Python concurrency,
idempotent retry.