人人都会AI编程

分布式锁、限流、消息队列场景实现

更新时间:2026-07-10

在分布式系统中,多个 Node.js 服务实例通常需要协调对共享资源的访问、保护系统免受过载、并在服务间异步传递消息。Redis 凭借其原子操作、发布订阅和数据过期能力,成为实现这些模式的轻量级中间件。以下分别介绍在 Node.js 中如何基于 Redis 实现分布式锁、接口限流和简易消息队列。

1. 分布式锁:用 Redis 保证资源独占

分布式锁用于确保同一时间只有一个服务实例可以执行某段关键代码(如更新库存、执行定时任务)。最常用的实现是 Redis 的 SET resource value NX PX timeout 命令,它能在键不存在时设置值并附加过期时间,整个过程是原子的。

基础实现(ioredis):

const Redis = require('ioredis');
const redis = new Redis();

async function acquireLock(lockKey, lockValue, ttlMs) {
  // NX: 仅当键不存在时设置; PX: 过期时间(毫秒)
  const result = await redis.set(lockKey, lockValue, 'PX', ttlMs, 'NX');
  return result === 'OK';
}

async function releaseLock(lockKey, lockValue) {
  // 使用 Lua 脚本保证判断和删除的原子性
  const script = `
    if redis.call("get", KEYS[1]) == ARGV[1] then
      return redis.call("del", KEYS[1])
    else
      return 0
    end
  `;
  await redis.eval(script, 1, lockKey, lockValue);
}

使用示例:

async function processTask() {
  const lockKey = 'lock:task:123';
  const lockValue = `${process.pid}-${Date.now()}`; // 唯一标识
  const locked = await acquireLock(lockKey, lockValue, 5000);

  if (!locked) {
    console.log('任务已被其他实例处理');
    return;
  }

  try {
    // 执行关键操作...
    await doSomethingImportant();
  } finally {
    await releaseLock(lockKey, lockValue);
  }
}

注意事项:

  • lockValue 必须是唯一值,防止误删其他实例持有的锁。
  • 过期时间应大于业务执行时间,避免锁提前失效。可启用“看门狗”定时续期。
  • 对于严格一致性要求的场景,建议使用 Redlock 算法(Redis 官方推荐的多节点锁方案),Node.js 社区有 redlock 包可供使用。

2. 接口限流:保护服务不被过载

当接口请求量突然飙升时,限流可以丢弃超出容量的请求,防止服务雪崩。常见限流算法有固定窗口、滑动窗口、令牌桶等。Redis 的原子递增和过期机制可以轻松实现固定窗口计数器限流。

固定窗口限流:

async function rateLimiter(userId, limit, windowSeconds) {
  const key = `ratelimit:${userId}:${Math.floor(Date.now() / 1000 / windowSeconds)}`;
  const current = await redis.incr(key);
  
  if (current === 1) {
    await redis.expire(key, windowSeconds);
  }
  
  return current <= limit;
}

滑动窗口限流(基于有序集合):

async function slidingWindowRateLimiter(userId, limit, windowMs) {
  const key = `sliding:${userId}`;
  const now = Date.now();
  const windowStart = now - windowMs;

  const script = `
    redis.call('ZREMRANGEBYSCORE', KEYS[1], 0, ARGV[1]);
    local current = redis.call('ZCARD', KEYS[1]);
    if current < tonumber(ARGV[2]) then
      redis.call('ZADD', KEYS[1], ARGV[3], ARGV[3]);
      return true
    else
      return false
    end
  `;
  return await redis.eval(script, 1, key, windowStart, limit, now);
}

生产级建议:

  • 对于高性能需求,可以使用 Redis 令牌桶方案,配合 Lua 脚本保证原子性。
  • 结合 Express/Koa 中间件集成限流逻辑,例如 express-rate-limit 搭配 rate-limit-redis 存储后端。
  • 限流配置应支持动态调整,避免紧急情况需要重启服务。

3. 消息队列:异步解耦与削峰填谷

Node.js 轻量级任务队列通常基于 Redis 的列表(List)或流(Stream)实现。对于高吞吐需求,可以直接使用 Bull、BullMQ 等成熟库,它们内置了重试、延迟、优先级等功能。

基于 Bull 的任务队列:

const Bull = require('bull');
const videoQueue = new Bull('video transcoding', 'redis://localhost:6379');

// 生产者:添加任务
videoQueue.add({ videoId: 'abc123', format: 'mp4' }, { delay: 5000 }); // 延迟5秒
videoQueue.add({ videoId: 'abc124' }, { priority: 1 });                // 高优先级

// 消费者:处理任务
videoQueue.process(async (job) => {
  const { videoId } = job.data;
  await transcode(videoId);
  console.log(`完成视频 ${videoId}`);
});

基于 Redis Stream(Node.js 14+):

// 生产消息
await redis.xadd('mystream', '*', 'userId', '42', 'action', 'purchase');

// 消费组模式读取
await redis.xgroup('CREATE', 'mystream', 'mygroup', '$', 'MKSTREAM');
const results = await redis.xreadgroup(
  'GROUP', 'mygroup', 'consumer-1',
  'BLOCK', 5000, 'COUNT', 10,
  'STREAMS', 'mystream', '>'
);

实用要点:

  • 任务失败处理:Bull 会按退避策略自动重试,或可手动调用 job.retry()
  • 队列监控:Bull 提供 UI (bull-board) 可直接挂载到 Express 服务上实时查看队列状态。
  • 对于需要严格顺序的消息,应使用单个队列(FIFO),避免使用多个消费者并行处理导致乱序。

通过以上模式,Node.js + Redis 组合可以在无需引入重型消息中间件(如 RabbitMQ、Kafka)的情况下,为中小规模分布式系统提供可靠的基础设施支撑。当业务量增长到一定程度时,再逐步迁移至专用中间件即可。