Skip to content
On this page

kue

简介

kue 是一个基于 Redis 的分布式任务队列工具包,专为 Node.js 设计。它提供了优先级、重试机制、事件系统、进度追踪等功能,适用于处理后台任务和异步操作。

主要特性

  • 基于 Redis 的持久化存储
  • 任务优先级管理
  • 失败任务自动重试
  • 任务进度跟踪
  • 并发控制
  • 延迟任务处理
  • Web 界面监控
  • 事件驱动架构

安装

bash
npm install kue

基本使用

1. 创建任务队列

javascript
const kue = require('kue')
const queue = kue.createQueue({
  redis: {
    port: 6379,
    host: 'localhost',
    auth: 'password', // 可选
    db: 0, // 可选
  },
})

2. 创建任务

javascript
const job = queue
  .create('email', {
    title: '发送欢迎邮件',
    to: 'user@example.com',
    template: 'welcome-email',
  })
  .priority('high')
  .attempts(5)
  .backoff({ delay: 60000, type: 'exponential' })
  .save()

3. 处理任务

javascript
queue.process('email', 3, function (job, done) {
  sendEmail(job.data.to, job.data.template, function (err) {
    if (err) return done(err)
    done()
  })
})

高级特性

1. 任务优先级

javascript
queue.create('email', data).priority('high').save()

支持的优先级:

  • low (10)
  • normal (0)
  • medium (-5)
  • high (-10)
  • critical (-15)

优先级数值越低,任务处理顺序越靠前。可以通过数字直接设置优先级:

javascript
job.priority(-10) // 等同于 high

2. 任务进度

javascript
queue.process('video', function (job, done) {
  job.progress(0, 100)
  convertVideo(
    job.data.file,
    function (progress) {
      job.progress(progress, 100)
    },
    done,
  )
})

3. 任务事件

javascript
job
  .on('complete', function () {
    console.log('任务完成')
  })
  .on('failed', function (err) {
    console.log('任务失败', err)
  })
  .on('progress', function (progress, total) {
    console.log('任务进度:', progress)
  })

4. 并发处理

javascript
queue.process('email', 20, function (job, done) {
  // 同时处理最多 20 个任务
})

最佳实践

  1. 错误处理

    javascript
    queue.on('error', function (err) {
      console.error('队列错误:', err)
    })
  2. 优雅关闭

    javascript
    process.once('SIGTERM', function () {
      queue.shutdown(5000, function (err) {
        console.log('Kue 已关闭', err || '')
        process.exit(0)
      })
    })
  3. 任务超时

    javascript
    queue.watchStuckJobs(1000)
  4. 任务清理

    javascript
    kue.Job.rangeByState('complete', 0, 1000, 'asc', function (err, jobs) {
      jobs.forEach(function (job) {
        job.remove(function () {
          console.log('已删除完成的任务:', job.id)
        })
      })
    })

应用场景

1. 邮件发送系统

javascript
// 生产者:创建邮件发送任务
app.post('/register', function (req, res) {
  const user = saveUser(req.body)

  queue
    .create('send-welcome-email', {
      userId: user.id,
      email: user.email,
      name: user.name,
    })
    .priority('high')
    .attempts(3)
    .save()

  res.send('注册成功!')
})

// 消费者:处理邮件发送
queue.process('send-welcome-email', 5, function (job, done) {
  const { userId, email, name } = job.data

  emailService
    .sendWelcomeEmail(email, name)
    .then(() => {
      console.log(`欢迎邮件已发送至 ${email}`)
      done()
    })
    .catch((err) => {
      console.error(`发送邮件失败: ${err.message}`)
      done(err)
    })
})

2. 图片处理服务

javascript
// 生产者:创建图片处理任务
app.post('/upload', upload.single('image'), function (req, res) {
  const originalPath = req.file.path

  queue
    .create('image-processing', {
      originalPath,
      formats: ['thumbnail', 'medium', 'large'],
      userId: req.user.id,
    })
    .attempts(2)
    .backoff({ delay: 60000, type: 'fixed' })
    .save()

  res.send('图片上传成功,正在处理...')
})

// 消费者:处理图片转换
queue.process('image-processing', 2, function (job, done) {
  const { originalPath, formats, userId } = job.data
  let completed = 0

  // 更新进度
  job.progress(0, formats.length)

  formats.forEach((format) => {
    imageProcessor
      .convert(originalPath, format)
      .then(() => {
        completed++
        job.progress(completed, formats.length)
        if (completed === formats.length) {
          done()
        }
      })
      .catch((err) => done(err))
  })
})

常见问题与解决方案

1. 任务卡住不执行

问题: 有时任务会卡在队列中不执行。

解决方案:

javascript
// 监控卡住的任务
queue.watchStuckJobs(5000)

// 设置任务超时
job.ttl(60000) // 60秒后任务超时

2. Redis 连接丢失

问题: Redis 连接断开导致任务处理中断。

解决方案:

javascript
// 设置重连选项
const queue = kue.createQueue({
  redis: {
    port: 6379,
    host: 'localhost',
    auth: 'password',
    options: {
      // 启用自动重连
      retry_strategy: function (options) {
        if (options.error && options.error.code === 'ECONNREFUSED') {
          return new Error('Redis服务器拒绝连接')
        }
        if (options.total_retry_time > 1000 * 60 * 60) {
          return new Error('重试时间已用尽')
        }
        if (options.attempt > 10) {
          return undefined // 停止重试
        }
        return Math.min(options.attempt * 100, 3000) // 重连延迟
      },
    },
  },
})

3. 内存占用过高

问题: 大量任务积压导致内存占用过高。

解决方案:

javascript
// 定期清理已完成任务
function cleanupJobs() {
  kue.Job.rangeByState('complete', 0, 1000, 'asc', function (err, jobs) {
    if (err) return console.error(err)

    jobs.forEach((job) => {
      job.remove((err) => {
        if (err) console.error(`无法删除任务 ${job.id}:`, err)
      })
    })
  })

  // 每小时运行一次
  setTimeout(cleanupJobs, 3600000)
}

cleanupJobs()

注意事项

  1. 确保 Redis 服务器正在运行
  2. 合理设置任务优先级和重试策略
  3. 监控任务队列的长度和处理速度
  4. 定期清理已完成的任务
  5. 在生产环境中使用可靠的 Redis 配置
  6. 任务处理函数应该是幂等的,以防重复执行
  7. 对于关键任务,实现完整的日志记录和监控
  8. 考虑使用集群模式提高可靠性和性能

要保持清醒 永远不抱有意外的幻想 凭空的期待最要命