Appearance
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) // 等同于 high2. 任务进度
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 个任务
})最佳实践
错误处理
javascriptqueue.on('error', function (err) { console.error('队列错误:', err) })优雅关闭
javascriptprocess.once('SIGTERM', function () { queue.shutdown(5000, function (err) { console.log('Kue 已关闭', err || '') process.exit(0) }) })任务超时
javascriptqueue.watchStuckJobs(1000)任务清理
javascriptkue.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()注意事项
- 确保 Redis 服务器正在运行
- 合理设置任务优先级和重试策略
- 监控任务队列的长度和处理速度
- 定期清理已完成的任务
- 在生产环境中使用可靠的 Redis 配置
- 任务处理函数应该是幂等的,以防重复执行
- 对于关键任务,实现完整的日志记录和监控
- 考虑使用集群模式提高可靠性和性能