Node.js 异步编程实践
util.promisify 与回调转换
基本用法
javascript
const util = require('util');
const fs = require('fs');
// 将回调风格的函数转换为 Promise
const readFile = util.promisify(fs.readFile);
const writeFile = util.promisify(fs.writeFile);
async function example() {
try {
const data = await readFile('input.txt', 'utf8');
await writeFile('output.txt', data.toUpperCase());
console.log('文件处理完成');
} catch (err) {
console.error('错误:', err.message);
}
}自定义 promisify 行为
javascript
const util = require('util');
// 对于特殊回调签名,可以使用 util.promisify.custom
function customCallbackApi(data, callback) {
// 回调签名:(result) 而不是 (error, result)
setTimeout(() => callback(data.toUpperCase()), 100);
}
// 自定义 promisify 实现
customCallbackApi[util.promisify.custom] = (data) => {
return new Promise((resolve) => {
customCallbackApi(data, resolve);
});
};
const promisifiedCustom = util.promisify(customCallbackApi);
const result = await promisifiedCustom('hello'); // 'HELLO'批量转换
javascript
const util = require('util');
const fs = require('fs');
// 批量 promisify
const fsPromises = {
readFile: util.promisify(fs.readFile),
writeFile: util.promisify(fs.writeFile),
stat: util.promisify(fs.stat),
readdir: util.promisify(fs.readdir),
};
// 或直接使用 fs.promises(Node.js 10+)
const { promises: fsPromisesAPI } = require('fs');fs/promises 模块
文件操作
javascript
const fs = require('fs/promises');
// 读取文件
async function readFile() {
try {
const data = await fs.readFile('data.txt', 'utf8');
console.log(data);
} catch (err) {
console.error('读取失败:', err.message);
}
}
// 写入文件
async function writeFile() {
await fs.writeFile('output.txt', 'Hello World', 'utf8');
console.log('写入完成');
}
// 追加内容
async function appendFile() {
await fs.appendFile('log.txt', '新的日志行\n');
}
// 复制文件
async function copyFile() {
await fs.copyFile('source.txt', 'dest.txt');
}
// 删除文件
async function deleteFile() {
await fs.unlink('old-file.txt');
}目录操作
javascript
const fs = require('fs/promises');
// 创建目录
async function createDir() {
await fs.mkdir('new-dir', { recursive: true });
}
// 读取目录
async function readDir() {
const files = await fs.readdir('./');
console.log('文件列表:', files);
}
// 读取目录详情
async function readDirWithDetails() {
const files = await fs.readdir('./', { withFileTypes: true });
files.forEach(file => {
console.log(`${file.name} - ${file.isDirectory() ? '目录' : '文件'}`);
});
}
// 删除目录
async function removeDir() {
await fs.rm('old-dir', { recursive: true, force: true });
}文件信息
javascript
const fs = require('fs/promises');
async function getFileStats() {
const stats = await fs.stat('data.txt');
console.log('文件大小:', stats.size, '字节');
console.log('是否文件:', stats.isFile());
console.log('是否目录:', stats.isDirectory());
console.log('创建时间:', stats.birthtime);
console.log('修改时间:', stats.mtime);
console.log('访问时间:', stats.atime);
}Stream 与异步迭代
可读流的异步迭代
javascript
const fs = require('fs');
// 使用 for await...of 迭代可读流
async function readStreamAsync() {
const stream = fs.createReadStream('large-file.txt', 'utf8');
for await (const chunk of stream) {
console.log('收到数据块:', chunk.length, '字节');
// 处理数据块
}
console.log('流读取完成');
}处理大文件
javascript
const fs = require('fs');
const readline = require('readline');
// 逐行读取大文件
async function processLargeFile() {
const fileStream = fs.createReadStream('large.log');
const rl = readline.createInterface({
input: fileStream,
crlfDelay: Infinity
});
let lineCount = 0;
for await (const line of rl) {
lineCount++;
// 处理每一行
if (line.includes('ERROR')) {
console.log(`第 ${lineCount} 行: ${line}`);
}
}
console.log(`总共 ${lineCount} 行`);
}管道操作
javascript
const fs = require('fs');
const { pipeline } = require('stream/promises');
const zlib = require('zlib');
// 使用 pipeline 处理流
async function compressFile() {
await pipeline(
fs.createReadStream('input.txt'),
zlib.createGzip(),
fs.createWriteStream('input.txt.gz')
);
console.log('压缩完成');
}
// 解压文件
async function decompressFile() {
await pipeline(
fs.createReadStream('input.txt.gz'),
zlib.createGunzip(),
fs.createWriteStream('output.txt')
);
console.log('解压完成');
}自定义可读流
javascript
const { Readable } = require('stream');
// 创建自定义可读流
class NumberStream extends Readable {
constructor(options) {
super(options);
this.current = 0;
this.max = options.max || 10;
}
_read(size) {
if (this.current >= this.max) {
this.push(null); // 结束流
return;
}
// 模拟异步数据生成
setTimeout(() => {
this.push(`数字: ${this.current++}\n`);
}, 100);
}
}
async function useCustomStream() {
const numberStream = new NumberStream({ max: 5 });
for await (const chunk of numberStream) {
console.log(chunk.toString().trim());
}
}Worker Threads 处理 CPU 密集任务
基本用法
javascript
// main.js
const { Worker, isMainThread, parentPort, workerData } = require('worker_threads');
if (isMainThread) {
// 主线程
const worker = new Worker(__filename, {
workerData: { number: 42 }
});
worker.on('message', (result) => {
console.log('计算结果:', result);
});
worker.on('error', (err) => {
console.error('工作线程错误:', err);
});
worker.on('exit', (code) => {
if (code !== 0) {
console.error(`工作线程异常退出,代码: ${code}`);
}
});
} else {
// 工作线程
const { number } = workerData;
// 执行 CPU 密集计算
function heavyComputation(n) {
let result = 0;
for (let i = 0; i < n * 1000000; i++) {
result += Math.sqrt(i);
}
return result;
}
const result = heavyComputation(number);
parentPort.postMessage(result);
}线程池实现
javascript
const { Worker } = require('worker_threads');
const path = require('path');
class WorkerPool {
constructor(numThreads = 4) {
this.numThreads = numThreads;
this.workers = [];
this.taskQueue = [];
this.availableWorkers = [];
// 初始化工作线程
for (let i = 0; i < numThreads; i++) {
const worker = new Worker(path.join(__dirname, 'worker.js'));
this.workers.push(worker);
this.availableWorkers.push(worker);
}
}
runTask(task) {
return new Promise((resolve, reject) => {
if (this.availableWorkers.length === 0) {
// 没有空闲线程,加入队列
this.taskQueue.push({ task, resolve, reject });
return;
}
const worker = this.availableWorkers.pop();
const messageHandler = (result) => {
worker.off('message', messageHandler);
this.availableWorkers.push(worker);
resolve(result);
// 处理队列中的下一个任务
if (this.taskQueue.length > 0) {
const next = this.taskQueue.shift();
this.runTask(next.task).then(next.resolve).catch(next.reject);
}
};
worker.on('message', messageHandler);
worker.on('error', reject);
worker.postMessage(task);
});
}
terminate() {
this.workers.forEach(worker => worker.terminate());
}
}
// 使用示例
const pool = new WorkerPool(4);
async function processTasks() {
const tasks = [1, 2, 3, 4, 5, 6, 7, 8];
const results = await Promise.all(
tasks.map(n => pool.runTask({ compute: n }))
);
console.log('所有结果:', results);
pool.terminate();
}异步错误处理最佳实践
全局错误处理
javascript
// 处理未捕获的 Promise 拒绝
process.on('unhandledRejection', (reason, promise) => {
console.error('未处理的 Promise 拒绝:', reason);
// 记录日志、发送告警
process.exit(1); // 建议退出进程
});
// 处理未捕获的异常
process.on('uncaughtException', (error) => {
console.error('未捕获的异常:', error);
// 记录日志、发送告警
process.exit(1); // 必须退出进程
});错误边界模式
javascript
// 封装错误处理逻辑
function withErrorHandling(fn) {
return async function(...args) {
try {
return await fn.apply(this, args);
} catch (error) {
console.error(`函数 ${fn.name} 执行失败:`, error.message);
throw error; // 重新抛出或返回默认值
}
};
}
// 使用示例
const safeReadFile = withErrorHandling(async (filename) => {
const fs = require('fs/promises');
return await fs.readFile(filename, 'utf8');
});
await safeReadFile('data.txt'); // 自动处理错误重试机制
javascript
// 带指数退避的重试
async function retryWithBackoff(fn, options = {}) {
const {
maxRetries = 3,
initialDelay = 1000,
maxDelay = 30000,
factor = 2
} = options;
let lastError;
for (let attempt = 0; attempt <= maxRetries; attempt++) {
try {
return await fn();
} catch (error) {
lastError = error;
if (attempt === maxRetries) {
break;
}
const delay = Math.min(
initialDelay * Math.pow(factor, attempt),
maxDelay
);
console.log(`第 ${attempt + 1} 次重试,等待 ${delay}ms`);
await new Promise(resolve => setTimeout(resolve, delay));
}
}
throw lastError;
}
// 使用示例
const data = await retryWithBackoff(
() => fetch('https://api.example.com/data'),
{ maxRetries: 5, initialDelay: 500 }
);并发控制
限制并发数量
javascript
// 并发控制器
class ConcurrencyController {
constructor(limit) {
this.limit = limit;
this.running = 0;
this.queue = [];
}
async run(task) {
if (this.running >= this.limit) {
await new Promise(resolve => this.queue.push(resolve));
}
this.running++;
try {
return await task();
} finally {
this.running--;
const next = this.queue.shift();
if (next) next();
}
}
}
// 使用示例
async function fetchAll(urls) {
const controller = new ConcurrencyController(5);
return Promise.all(
urls.map(url => controller.run(() => fetch(url)))
);
}使用 p-limit 库模式
javascript
// 实现类似 p-limit 的功能
function pLimit(concurrency) {
const queue = [];
let activeCount = 0;
const next = () => {
activeCount--;
if (queue.length > 0) {
queue.shift()();
}
};
const run = async (fn, ...args) => {
activeCount++;
try {
return await fn(...args);
} finally {
next();
}
};
const enqueue = (fn, ...args) => {
return new Promise((resolve, reject) => {
const runTask = () => {
run(fn, ...args).then(resolve).catch(reject);
};
if (activeCount < concurrency) {
runTask();
} else {
queue.push(runTask);
}
});
};
return enqueue;
}
// 使用示例
const limit = pLimit(3);
async function processFiles(files) {
return Promise.all(
files.map(file => limit(() => processFile(file)))
);
}串行执行
javascript
// 串行执行异步任务
async function series(tasks) {
const results = [];
for (const task of tasks) {
results.push(await task());
}
return results;
}
// 使用示例
const results = await series([
() => fs.readFile('file1.txt'),
() => fs.readFile('file2.txt'),
() => fs.readFile('file3.txt'),
]);流量控制与背压
背压处理
javascript
const fs = require('fs');
// 处理可写流的背压
async function writeWithBackpressure() {
const readable = fs.createReadStream('source.txt');
const writable = fs.createWriteStream('dest.txt');
for await (const chunk of readable) {
// 检查是否可以继续写入
if (!writable.write(chunk)) {
// 等待 drain 事件
await new Promise(resolve => writable.once('drain', resolve));
}
}
writable.end();
}限流器实现
javascript
// 令牌桶限流器
class RateLimiter {
constructor(rate, capacity) {
this.rate = rate; // 每秒添加的令牌数
this.capacity = capacity; // 桶容量
this.tokens = capacity;
this.lastTime = Date.now();
this.queue = [];
}
async acquire() {
this.refill();
if (this.tokens >= 1) {
this.tokens -= 1;
return true;
}
// 需要等待
return new Promise(resolve => {
this.queue.push(resolve);
});
}
refill() {
const now = Date.now();
const elapsed = (now - this.lastTime) / 1000;
this.tokens = Math.min(
this.capacity,
this.tokens + elapsed * this.rate
);
this.lastTime = now;
// 处理队列
while (this.queue.length > 0 && this.tokens >= 1) {
this.tokens -= 1;
const resolve = this.queue.shift();
resolve(true);
}
}
}
// 使用示例
const limiter = new RateLimiter(10, 100); // 每秒 10 个请求,最大 100 个
async function processRequest() {
await limiter.acquire();
// 执行请求
}缓存策略
内存缓存
javascript
class AsyncCache {
constructor(ttl = 60000) {
this.cache = new Map();
this.ttl = ttl;
}
async get(key, fetcher) {
const cached = this.cache.get(key);
if (cached && Date.now() - cached.timestamp < this.ttl) {
return cached.value;
}
// 防止缓存击穿:多个请求同时获取
if (cached?.promise) {
return cached.promise;
}
const promise = fetcher().then(value => {
this.cache.set(key, { value, timestamp: Date.now() });
return value;
});
this.cache.set(key, { promise, timestamp: Date.now() });
return promise;
}
delete(key) {
this.cache.delete(key);
}
clear() {
this.cache.clear();
}
}
// 使用示例
const cache = new AsyncCache(30000);
async function getUser(id) {
return cache.get(`user:${id}`, async () => {
const response = await fetch(`/api/users/${id}`);
return response.json();
});
}文件缓存
javascript
const fs = require('fs/promises');
const path = require('path');
class FileCache {
constructor(cacheDir) {
this.cacheDir = cacheDir;
}
async get(key) {
try {
const filePath = path.join(this.cacheDir, `${key}.json`);
const data = await fs.readFile(filePath, 'utf8');
const cached = JSON.parse(data);
// 检查是否过期
if (cached.expiry && Date.now() > cached.expiry) {
await fs.unlink(filePath);
return null;
}
return cached.data;
} catch {
return null;
}
}
async set(key, data, ttl = 3600000) {
const filePath = path.join(this.cacheDir, `${key}.json`);
const cached = {
data,
expiry: Date.now() + ttl
};
await fs.mkdir(this.cacheDir, { recursive: true });
await fs.writeFile(filePath, JSON.stringify(cached));
}
}小结
- 使用
util.promisify或fs/promises统一异步 API 风格 - 使用 Stream 处理大文件,避免内存溢出
- CPU 密集型任务使用 Worker Threads 并行处理
- 实现全局错误处理和重试机制
- 限制并发数量,避免资源耗尽
- 使用背压处理防止内存问题
- 合理使用缓存减少重复计算和网络请求
Node.js 22+ 异步编程实践新特性
node:test 内置测试框架
Node.js 22+ 内置 node:test,无需安装 Jest/Mocha:
javascript
import { describe, it, beforeEach, afterEach, mock } from 'node:test'
import assert from 'node:assert/strict'
describe('用户服务', () => {
beforeEach(() => {
mock.fn(console.log)
})
afterEach(() => {
mock.restoreAll()
})
it('应该获取用户列表', async () => {
const users = await getUsers()
assert.ok(Array.isArray(users))
})
it('应该正确处理并发请求', async () => {
const results = await Promise.allSettled([
fetchUser(1),
fetchUser(2),
fetchUser(3)
])
const succeeded = results.filter(r => r.status === 'fulfilled')
assert.ok(succeeded.length >= 2)
})
})bash
# 运行测试
node --test
# 运行指定测试文件
node --test test/user-service.test.js
# 监视模式
node --test --watchPromise.withResolvers
Node.js 22+ 支持 Promise.withResolvers(),简化 Promise 创建:
javascript
// 传统方式
let resolve, reject
const promise = new Promise((res, rej) => {
resolve = res
reject = rej
})
// Promise.withResolvers(Node.js 22+)
const { promise, resolve, reject } = Promise.withResolvers()
// 实际应用:可取消的异步操作
function createCancellableTask(timeout = 5000) {
const { promise, resolve, reject } = Promise.withResolvers()
const timer = setTimeout(() => reject(new Error('超时')), timeout)
return {
promise: promise.finally(() => clearTimeout(timer)),
cancel: () => reject(new Error('已取消')),
complete: (value) => resolve(value)
}
}AsyncIterator helpers
Node.js 22+ 支持异步迭代器辅助方法:
javascript
// 异步生成器 + for await...of
async function* fetchPages(url) {
let page = 1
while (true) {
const response = await fetch(`${url}?page=${page}`)
const data = await response.json()
if (data.length === 0) break
yield data
page++
}
}
// 消费异步迭代器
for await (const page of fetchPages('https://api.example.com/users')) {
for (const user of page) {
console.log(user.name)
}
}