{T}

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.promisifyfs/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 --watch

Promise.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)
  }
}