前言

在Node.js学习中,事件循环(Event Loop)和异步编程是最核心且最难掌握的概念之一。很多开发者在使用Node.js多年后,仍然对事件循环的工作机制一知半解。本文将深入剖析Node.js的异步特性和事件循环机制,帮助你真正理解Node.js的并发模型。

什么是事件循环?

基本概念

事件循环是Node.js实现非阻塞I/O操作的关键机制。尽管JavaScript是单线程的,但通过事件循环和异步I/O,Node.js可以处理大量并发连接。

为什么需要事件循环?

在传统的多线程服务器模型中,每个连接都需要创建一个新的线程,当连接数增加时,线程间的上下文切换会消耗大量资源。而Node.js使用单线程加事件循环的模式,避免了线程创建和切换的开销。

Node.js 的线程架构

误解澄清:Node.js真的是单线程吗?

这是一个常见的误解。实际上,Node.js在以下方面使用多线程:

JavaScript主线程(单线程执行环境)
    事件循环、回调函数执行、非CPU密集型任务
        |
        V
libuv线程池(默认4个线程,可配置)
    文件I/O、DNS解析、部分加密操作、其他阻塞操作

事件循环阶段详解

事件循环的六个阶段

Node.js事件循环分为六个不同的阶段,每个阶段都有一个先进先出(FIFO)的回调队列:

开始

V
timers阶段:执行setTimeout和setInterval的回调

V
pending callbacks阶段:执行延迟到下一个循环迭代的I/O回调

V
idle, prepare阶段:仅内部使用

V
poll阶段:检索新的I/O事件,执行相关回调

V
check阶段:执行setImmediate()的回调

V
close callbacks阶段:执行关闭事件的回调,如socket.on('close')

V
回到timers阶段,继续循环

各阶段详细说明

1. Timers 阶段

console.log('开始');

setTimeout(() => {
    console.log('setTimeout回调');
}, 0);

setImmediate(() => {
    console.log('setImmediate回调');
});

console.log('结束');

输出顺序可能是:
开始
结束
setTimeout回调
setImmediate回调
或者:
开始
结束
setImmediate回调
setTimeout回调

2. Pending Callbacks 阶段

处理上一轮循环中延迟的执行I/O回调。

3. Poll 阶段(最重要的阶段)

const fs = require('fs');

console.log('开始读取文件');

I/O操作在poll阶段处理
fs.readFile('example.txt', 'utf8', (err, data) => {
    if (err) throw err;
    console.log('文件内容:', data);
});

setImmediate在poll阶段后立即执行
setImmediate(() => {
    console.log('setImmediate在I/O回调后执行');
});

console.log('继续执行其他代码');

4. Check 阶段

const fs = require('fs');

fs.readFile('example.txt', 'utf8', (err, data) => {
    I/O回调在poll阶段执行
    console.log('I/O回调');
    
    setImmediate(() => {
        这个setImmediate在check阶段执行
        console.log('在I/O回调中的setImmediate');
    });
});

setImmediate(() => {
    这个setImmediate也在check阶段执行
    console.log('外部的setImmediate');
});

异步编程模式

1. 回调函数(Callback)

const fs = require('fs');

回调地狱示例
fs.readFile('file1.txt', 'utf8', (err, data1) => {
    if (err) throw err;
    fs.readFile('file2.txt', 'utf8', (err, data2) => {
        if (err) throw err;
        fs.writeFile('result.txt', data1 + data2, (err) => {
            if (err) throw err;
            console.log('文件合并完成');
        });
    });
});

2. Promise

const fs = require('fs').promises;

使用Promise避免回调地狱
fs.readFile('file1.txt', 'utf8')
    .then(data1 => {
        return fs.readFile('file2.txt', 'utf8')
            .then(data2 => data1 + data2);
    })
    .then(combinedData => {
        return fs.writeFile('result.txt', combinedData);
    })
    .then(() => {
        console.log('文件合并完成');
    })
    .catch(err => {
        console.error('出错:', err);
    });

使用async/await更简洁
async function mergeFiles() {
    try {
        const data1 = await fs.readFile('file1.txt', 'utf8');
        const data2 = await fs.readFile('file2.txt', 'utf8');
        await fs.writeFile('result.txt', data1 + data2);
        console.log('文件合并完成');
    } catch (err) {
        console.error('出错:', err);
    }
}

3. Async/Await

class FileProcessor {
    constructor() {
        this.fs = require('fs').promises;
    }
    
    async processFiles() {
        try {
            并行处理文件读取
            const [data1, data2] = await Promise.all([
                this.fs.readFile('file1.txt', 'utf8'),
                this.fs.readFile('file2.txt', 'utf8')
            ]);
            
            串行处理
            const processedData = await this.processData(data1, data2);
            await this.saveResult(processedData);
            
            return '处理完成';
        } catch (error) {
            throw new Error('文件处理失败: ' + error.message);
        }
    }
    
    async processData(data1, data2) {
        模拟数据处理
        return new Promise(resolve => {
            setTimeout(() => {
                resolve(data1.toUpperCase() + data2.toUpperCase());
            }, 1000);
        });
    }
    
    async saveResult(data) {
        await this.fs.writeFile('result.txt', data);
    }
}

使用类
const processor = new FileProcessor();
processor.processFiles()
    .then(result => console.log(result))
    .catch(error => console.error(error));

高级异步模式

1. 控制并发

class ConcurrentController {
    constructor(maxConcurrent = 3) {
        this.maxConcurrent = maxConcurrent;
        this.queue = [];
        this.running = 0;
    }
    
    async add(task) {
        return new Promise((resolve, reject) => {
            this.queue.push({ task, resolve, reject });
            this.run();
        });
    }
    
    async run() {
        if (this.running >= this.maxConcurrent || this.queue.length === 0) {
            return;
        }
        
        this.running++;
        const { task, resolve, reject } = this.queue.shift();
        
        try {
            const result = await task();
            resolve(result);
        } catch (error) {
            reject(error);
        } finally {
            this.running--;
            this.run();
        }
    }
}

使用并发控制器
const controller = new ConcurrentController(2);

模拟多个异步任务
const tasks = Array.from({ length: 10 }, (_, i) => 
    () => new Promise(resolve => {
        setTimeout(() => {
            console.log('任务 ' + (i + 1) + ' 完成');
            resolve(i + 1);
        }, Math.random() * 1000);
    })
);

并行执行,但最多同时执行2个任务
Promise.all(tasks.map(task => controller.add(task)))
    .then(results => console.log('所有任务完成:', results));

2. 错误处理模式

高级错误处理策略
class RobustAsyncHandler {
    static async retry(fn, retries = 3, delay = 1000) {
        try {
            return await fn();
        } catch (error) {
            if (retries > 0) {
                console.log('重试中... 剩余 ' + retries + ' 次');
                await new Promise(resolve => setTimeout(resolve, delay));
                return this.retry(fn, retries - 1, delay * 2); // 指数退避
            }
            throw error;
        }
    }
    
    static async timeout(fn, ms) {
        return new Promise(async (resolve, reject) => {
            const timeoutId = setTimeout(() => {
                reject(new Error('操作超时: ' + ms + 'ms'));
            }, ms);
            
            try {
                const result = await fn();
                clearTimeout(timeoutId);
                resolve(result);
            } catch (error) {
                clearTimeout(timeoutId);
                reject(error);
            }
        });
    }
    
    static async allSettled(tasks) {
        const results = await Promise.allSettled(tasks);
        const successful = results.filter(r => r.status === 'fulfilled');
        const failed = results.filter(r => r.status === 'rejected');
        
        return {
            successful: successful.map(r => r.value),
            failed: failed.map(r => r.reason),
            total: results.length
        };
    }
}

使用示例
async function demo() {
    重试机制
    const result = await RobustAsyncHandler.retry(
        () => fetch('https://api.example.com/data'),
        3
    );
    
    超时控制
    const data = await RobustAsyncHandler.timeout(
        () => fetch('https://slow-api.com/data'),
        5000
    );
    
    处理多个可能失败的任务
    const tasks = [
        fetch('https://api1.com'),
        fetch('https://api2.com'),
        fetch('https://api3.com')
    ];
    
    const results = await RobustAsyncHandler.allSettled(tasks);
    console.log('成功: ' + results.successful.length + ', 失败: ' + results.failed.length);
}

性能优化实践

1. 避免阻塞事件循环

错误示例:阻塞事件循环
function calculateSum(n) {
    let sum = 0;
    for (let i = 0; i < n; i++) {
        sum += i; // CPU密集型任务会阻塞事件循环
    }
    return sum;
}

正确示例:分解任务
async function calculateSumAsync(n, chunkSize = 1000000) {
    let sum = 0;
    
    for (let i = 0; i < n; i += chunkSize) {
        使用setImmediate让出事件循环
        await new Promise(resolve => setImmediate(resolve));
        
        const end = Math.min(i + chunkSize, n);
        for (let j = i; j < end; j++) {
            sum += j;
        }
    }
    
    return sum;
}

使用Worker Threads处理CPU密集型任务
const { Worker, isMainThread, parentPort } = require('worker_threads');

if (isMainThread) {
    主线程
    function calculateWithWorker(n) {
        return new Promise((resolve, reject) => {
            const worker = new Worker(__filename, { 
                workerData: n 
            });
            
            worker.on('message', resolve);
            worker.on('error', reject);
            worker.on('exit', (code) => {
                if (code !== 0) {
                    reject(new Error('Worker stopped with exit code ' + code));
                }
            });
        });
    }
} else {
    Worker线程
    const n = require('worker_threads').workerData;
    let sum = 0;
    for (let i = 0; i < n; i++) {
        sum += i;
    }
    parentPort.postMessage(sum);
}

实际应用案例

构建高效的Web服务器

const http = require('http');
const { URL } = require('url');

class AsyncWebServer {
    constructor() {
        this.routes = new Map();
        this.middlewares = [];
    }
    
    use(middleware) {
        this.middlewares.push(middleware);
    }
    
    get(path, handler) {
        this.routes.set('GET:' + path, handler);
    }
    
    post(path, handler) {
        this.routes.set('POST:' + path, handler);
    }
    
    async handleRequest(req, res) {
        const url = new URL(req.url, 'http://' + req.headers.host);
        const routeKey = req.method + ':' + url.pathname;
        const handler = this.routes.get(routeKey);
        
        if (!handler) {
            res.writeHead(404);
            res.end('Not Found');
            return;
        }
        
        执行中间件
        for (const middleware of this.middlewares) {
            await new Promise((resolve, reject) => {
                middleware(req, res, (err) => {
                    if (err) reject(err);
                    else resolve();
                });
            });
        }
        
        处理请求
        try {
            await handler(req, res);
        } catch (error) {
            res.writeHead(500);
            res.end('Internal Server Error');
            console.error('处理请求出错:', error);
        }
    }
    
    start(port = 3000) {
        const server = http.createServer((req, res) => {
            this.handleRequest(req, res).catch(console.error);
        });
        
        return server.listen(port, () => {
            console.log('服务器运行在 http://localhost:' + port);
        });
    }
}

使用示例
const server = new AsyncWebServer();

添加中间件
server.use(async (req, res, next) => {
    console.log(new Date().toISOString() + ' - ' + req.method + ' ' + req.url);
    next();
});

添加路由
server.get('/api/data', async (req, res) => {
    模拟异步数据获取
    const data = await fetchDataFromDatabase();
    res.writeHead(200, { 'Content-Type': 'application/json' });
    res.end(JSON.stringify(data));
});

server.post('/api/data', async (req, res) => {
    处理POST请求
    const body = await parseRequestBody(req);
    await saveToDatabase(body);
    res.writeHead(201);
    res.end('Data created');
});

async function fetchDataFromDatabase() {
    return new Promise(resolve => {
        setTimeout(() => {
            resolve({ message: 'Hello from async server!' });
        }, 100);
    });
}

async function parseRequestBody(req) {
    return new Promise((resolve) => {
        let body = '';
        req.on('data', chunk => body += chunk);
        req.on('end', () => resolve(JSON.parse(body)));
    });
}

启动服务器
server.start(3000);

总结

1. 编写更高效的代码:避免阻塞事件循环,合理利用异步特性
2. 更好地调试问题:理解代码执行顺序,快速定位性能瓶颈
3. 构建可扩展应用:利用Node.js的并发模型处理高负载场景

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐