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;
}
}
最佳实践建议
- 并发控制:合理设置并发限制,避免资源耗尽
- 错误处理:使用 Promise.allSettled 而非 all 处理部分失败
- 超时机制:为所有异步操作设置超时
- 背压处理:消费者慢时暂停生产,避免内存溢出
- 监控告警:监控并发数、队列长度等指标
总结
响应式编程与并发处理是 Node.js 性能优化的关键:
- 并发限制:防止系统过载
- 执行策略:根据场景选择顺序、并发或混合
- 流式处理:大文件、大数据量场景必选
- 事件驱动:解耦复杂业务逻辑
- 背压处理:平衡生产与消费速率
掌握这些技术,可以构建高性能、高可用的 Koa.js 应用。