KOA技术分享

专注 Koa.js 框架的编程知识分享

Koa.js 响应式编程与高性能并发处理

引言

Koa.js 基于 ES6 的 async/await 语法,天然支持异步编程。在高并发场景下,如何高效地处理并发请求、优化资源利用、提升系统吞吐量,是每个 Node.js 开发者需要掌握的核心技能。

Promise 并发控制

Promise.all 的正确使用姿势:

// 并发限制器
class ConcurrencyLimiter {
  constructor(limit) {
    this.limit = limit;
    this.running = 0;
    this.queue = [];
  }

  // 执行带并发限制的任务
  async run(task) {
    return new Promise((resolve, reject) => {
      this.queue.push({ task, resolve, reject });
      this.process();
    });
  }

  // 处理队列
  process() {
    while (this.running < this.limit && this.queue.length > 0) {
      const { task, resolve, reject } = this.queue.shift();

      this.running++;
      task()
        .then(resolve)
        .catch(reject)
        .finally(() => {
          this.running--;
          this.process();
        });
    }
  }
}

// 使用示例:限制并发数为 5
const limiter = new ConcurrencyLimiter(5);

const tasks = Array.from({ length: 20 }, (_, i) => () =>
  fetch(`/api/item/${i}`).then(r => r.json())
);

// 并发执行,但最多同时 5 个
const results = await Promise.all(
  tasks.map(task => limiter.run(task))
);

// 并发映射(带限制)
async function concurrentMap(array, fn, limit = 10) {
  const limiter = new ConcurrencyLimiter(limit);
  return Promise.all(array.map(item => limiter.run(() => fn(item))));
}

顺序与并发混合策略

根据业务场景选择合适的执行策略:

// 执行策略管理器
class ExecutionStrategy {
  // 纯并发:所有任务同时执行
  static async parallel(tasks) {
    return Promise.all(tasks.map(task => task()));
  }

  // 纯顺序:任务依次执行
  static async sequential(tasks) {
    const results = [];
    for (const task of tasks) {
      results.push(await task());
    }
    return results;
  }

  // 分组并发:分成 N 批执行
  static async batched(tasks, batchSize = 10) {
    const results = [];
    for (let i = 0; i < tasks.length; i += batchSize) {
      const batch = tasks.slice(i, i + batchSize);
      const batchResults = await Promise.all(batch.map(task => task()));
      results.push(...batchResults);
    }
    return results;
  }

  // 管道流:前一个任务结果作为下一个输入
  static async pipeline(tasks) {
    let result;
    for (const task of tasks) {
      result = await task(result);
    }
    return result;
  }

  // 哀兵策略:失败时重试
  static async withRetry(task, options = {}) {
    const {
      maxRetries = 3,
      delay = 1000,
      backoff = 2,
      shouldRetry = () => true
    } = options;

    let lastError;
    let currentDelay = delay;

    for (let attempt = 0; attempt <= maxRetries; attempt++) {
      try {
        return await task();
      } catch (error) {
        lastError = error;

        if (attempt < maxRetries && shouldRetry(error)) {
          await this.sleep(currentDelay);
          currentDelay *= backoff;
        }
      }
    }

    throw lastError;
  }

  // 超时控制
  static async withTimeout(promise, ms) {
    return Promise.race([
      promise,
      new Promise((_, reject) =>
        setTimeout(() => reject(new Error('Timeout')), ms)
      )
    ]);
  }

  static sleep(ms) {
    return new Promise(resolve => setTimeout(resolve, ms));
  }
}

// 实际应用示例
async function getUserOrdersWithDetails(userId) {
  // 第一步:获取用户信息(必须先执行)
  const userInfo = await ExecutionStrategy.withRetry(
    () => fetchUser(userId),
    { maxRetries: 3 }
  );

  // 第二步:并发获取订单和推荐(可并行)
  const [orders, recommendations] = await ExecutionStrategy.parallel([
    () => fetchUserOrders(userId),
    () => fetchRecommendations(userId)
  ]);

  // 第三步:顺序获取订单详情(依赖订单列表)
  const orderDetails = await ExecutionStrategy.sequential(
    orders.map(order => () => fetchOrderDetail(order.id))
  );

  return { userInfo, orders: orderDetails, recommendations };
}

流式处理大文件

使用流处理避免内存溢出:

// 流式文件处理
const fs = require('fs');
const { pipeline } = require('stream/promises');
const { Transform } = require('stream');

// 流式 JSON 处理
class JSONLineParser extends Transform {
  constructor(options = {}) {
    super({ ...options, objectMode: true });
    this.buffer = '';
  }

  _transform(chunk, encoding, callback) {
    this.buffer += chunk.toString();

    const lines = this.buffer.split('\n');
    this.buffer = lines.pop(); // 最后一行可能不完整

    for (const line of lines) {
      if (line.trim()) {
        try {
          this.push(JSON.parse(line));
        } catch (e) {
          this.emit('error', new Error(`Invalid JSON: ${line}`));
        }
      }
    }

    callback();
  }

  _flush(callback) {
    if (this.buffer.trim()) {
      try {
        this.push(JSON.parse(this.buffer));
      } catch (e) {
        // 忽略最后不完整的行
      }
    }
    callback();
  }
}

// 流式 CSV 处理
async function processLargeCSV(filePath, handler, batchSize = 1000) {
  const readStream = fs.createReadStream(filePath, { encoding: 'utf8' });
  const parser = new JSONLineParser();

  let batch = [];
  let processed = 0;

  for await (const record of readStream.pipe(parser)) {
    batch.push(record);

    if (batch.length >= batchSize) {
      await handler(batch);
      processed += batch.length;
      batch = [];
    }
  }

  // 处理剩余数据
  if (batch.length > 0) {
    await handler(batch);
    processed += batch.length;
  }

  return processed;
}

// 文件上传流处理
async function handleUpload(ctx) {
  const stream = ctx.req;
  const uploadDir = './uploads';
  const filename = `${Date.now()}-${Math.random()}.tmp`;
  const out = fs.createWriteStream(path.join(uploadDir, filename));

  let uploaded = 0;

  await pipeline(
    stream,
    new Transform({
      transform(chunk, encoding, callback) {
        // 可以在这里添加进度记录
        uploaded += chunk.length;
        ctx.set('X-Upload-Progress', `${uploaded}`);
        callback(null, chunk);
      }
    }),
    out
  );

  return { filename, size: uploaded };
}

事件驱动架构

使用 EventEmitter 实现解耦:

const { EventEmitter } = require('events');

// 事件总线
class EventBus extends EventEmitter {
  constructor() {
    super();
    this.setMaxListeners(100);
  }

  // 同步发布
  publish(event, data) {
    this.emit(event, data);
  }

  // 异步发布(带确认)
  async publishAsync(event, data, timeout = 5000) {
    const handlers = this.listeners(event);

    if (handlers.length === 0) return [];

    const promises = handlers.map(handler =>
      Promise.race([
        handler(data),
        new Promise((_, reject) =>
          setTimeout(() => reject(new Error('Handler timeout')), timeout)
        )
      ])
    );

    return Promise.allSettled(promises);
  }

  // 订阅(支持一次性订阅)
  subscribe(event, handler, once = false) {
    if (once) {
      this.once(event, handler);
    } else {
      this.on(event, handler);
    }

    // 返回取消订阅函数
    return () => this.off(event, handler);
  }
}

// 订单领域事件
class OrderEvents extends EventBus {
  constructor() {
    super();

    // 订单创建
    this.on('order:created', this.handleOrderCreated.bind(this));

    // 订单支付
    this.on('order:paid', this.handleOrderPaid.bind(this));

    // 订单发货
    this.on('order:shipped', this.handleOrderShipped.bind(this));

    // 订单完成
    this.on('order:completed', this.handleOrderCompleted.bind(this));
  }

  async handleOrderCreated(order) {
    console.log(`Order created: ${order.id}`);

    // 发送确认邮件
    await this.sendEmail(order.userId, 'order_created', order);

    // 更新统计
    await this.updateStatistics('orders_created');
  }

  async handleOrderPaid(order) {
    console.log(`Order paid: ${order.id}`);

    // 通知仓库备货
    await this.notifyWarehouse(order);

    // 发送支付成功通知
    await this.sendNotification(order.userId, 'Payment successful');
  }

  async handleOrderShipped(order) {
    const tracking = await this.generateTrackingNumber(order);

    // 发送物流通知
    await this.sendNotification(order.userId, 'Your order has been shipped', {
      tracking
    });

    // 更新物流信息
    await this.updateLogistics(order.id, tracking);
  }

  async handleOrderCompleted(order) {
    // 计算用户积分
    await this.calculatePoints(order.userId, order.amount);

    // 邀请评价
    await this.scheduleReviewRequest(order.id);

    // 更新销售统计
    await this.updateSalesStats(order);
  }
}

// 在 Koa 中间件使用
const orderEvents = new OrderEvents();

router.post('/orders', async (ctx) => {
  const order = await createOrder(ctx.request.body);

  // 触发领域事件
  orderEvents.publish('order:created', order);

  ctx.body = order;
});

// 异步处理事件(不阻塞响应)
orderEvents.publishAsync('order:created', order).catch(console.error);

并发安全与锁

处理共享资源的并发访问:

// 分布式锁
class DistributedLock {
  constructor(redis) {
    this.redis = redis;
    this.defaultTTL = 30000; // 30秒
  }

  // 获取锁
  async acquire(key, ttl = this.defaultTTL) {
    const lockKey = `lock:${key}`;
    const lockValue = `${Date.now()}-${Math.random()}`;

    const result = await this.redis.set(lockKey, lockValue, 'PX', ttl, 'NX');

    if (result === 'OK') {
      return {
        release: () => this.release(key, lockValue),
        extend: (newTTL) => this.extend(key, lockValue, newTTL)
      };
    }

    return null;
  }

  // 释放锁(Lua 脚本确保原子性)
  async release(key, value) {
    const script = `
      if redis.call("get", KEYS[1]) == ARGV[1] then
        return redis.call("del", KEYS[1])
      else
        return 0
      end
    `;

    return await this.redis.eval(script, 1, `lock:${key}`, value);
  }

  // 延长锁时间
  async extend(key, value, ttl) {
    const script = `
      if redis.call("get", KEYS[1]) == ARGV[1] then
        return redis.call("pexpire", KEYS[1], ARGV[2])
      else
        return 0
      end
    `;

    return await this.redis.eval(script, 1, `lock:${key}`, value, ttl);
  }

  // 使用锁的装饰器
  static withLock(lock, key, fn) {
    return async function(...args) {
      const lockInstance = await lock.acquire(key);

      if (!lockInstance) {
        throw new Error('Failed to acquire lock');
      }

      try {
        return await fn.apply(this, args);
      } finally {
        await lockInstance.release();
      }
    };
  }
}

// 信号量(内存版)
class Semaphore {
  constructor(count) {
    this.count = count;
    this.waitQueue = [];
  }

  async acquire() {
    if (this.count > 0) {
      this.count--;
      return true;
    }

    return new Promise(resolve => {
      this.waitQueue.push(resolve);
    });
  }

  release() {
    if (this.waitQueue.length > 0) {
      const resolve = this.waitQueue.shift();
      resolve(true);
    } else {
      this.count++;
    }
  }

  // 使用 async with 模式
  async use(fn) {
    await this.acquire();
    try {
      return await fn();
    } finally {
      this.release();
    }
  }
}

// API 限流中间件
const rateLimiter = new Semaphore(100); // 最多 100 并发

async function rateLimitMiddleware(ctx, next) {
  try {
    await rateLimiter.acquire();
    await next();
  } finally {
    rateLimiter.release();
  }
}

背压处理

当消费者慢于生产者时的应对策略:

// 背压处理机制
class BackPressureHandler {
  constructor(options = {}) {
    this.highWatermark = options.highWatermark || 1000;
    this.lowWatermark = options.lowWatermark || 100;
    this.paused = false;
    this.queue = [];
    this.pauseCallbacks = [];
    this.resumeCallbacks = [];
  }

  // 添加数据
  push(data) {
    if (this.paused) {
      // 等待或拒绝
      return Promise.reject(new Error('Stream paused'));
    }

    this.queue.push(data);

    // 触发背压
    if (this.queue.length >= this.highWatermark) {
      this.pause();
    }

    return Promise.resolve();
  }

  // 消费数据
  async pop() {
    if (this.queue.length === 0) {
      // 等待新数据
      await new Promise(resolve => {
        const checkData = () => {
          if (this.queue.length > 0) {
            resolve();
          } else {
            setTimeout(checkData, 10);
          }
        };
        checkData();
      });
    }

    const data = this.queue.shift();

    // 恢复背压
    if (this.paused && this.queue.length <= this.lowWatermark) {
      this.resume();
    }

    return data;
  }

  // 暂停(通知生产者)
  pause() {
    if (!this.paused) {
      this.paused = true;
      this.pauseCallbacks.forEach(cb => cb());
    }
  }

  // 恢复
  resume() {
    if (this.paused) {
      this.paused = false;
      this.resumeCallbacks.forEach(cb => cb());
    }
  }

  on(event, callback) {
    if (event === 'pause') this.pauseCallbacks.push(callback);
    if (event === 'resume') this.resumeCallbacks.push(callback);
  }
}

// 在 Koa 中使用背压
async function streamingResponse(ctx, dataSource) {
  const handler = new BackPressureHandler({
    highWatermark: 100,
    lowWatermark: 20
  });

  // 背压时暂停数据生产
  handler.on('pause', () => {
    dataSource.pause();
  });

  handler.on('resume', () => {
    dataSource.resume();
  });

  ctx.set('Content-Type', 'application/x-ndjson');
  ctx.set('Transfer-Encoding', 'chunked');

  // 消费数据并发送
  while (true) {
    const data = await handler.pop();
    ctx.res.write(JSON.stringify(data) + '\n');

    // 检查连接是否断开
    if (ctx.res.writableEnded) break;
  }
}

最佳实践建议

总结

响应式编程与并发处理是 Node.js 性能优化的关键:

掌握这些技术,可以构建高性能、高可用的 Koa.js 应用。

← 下一篇:Koa.js WebAssembly 集成与高性能计算