分布式架构下,消息队列承担着系统解耦、流量削峰和异步通信的重任,但生产环境远比理想状态复杂:下游服务偶尔超时、数据库临时不可用、消息体格式异常,这些都可能导致消费失败。如果失败的消息被直接丢弃或者无限重试,后果都很严重。死信队列(Dead Letter Queue,简称DLQ)的价值就在于此——把那些处理不了的消息单独收集起来,等待人工介入或自动补偿。这篇文章就用Node.js实际写一遍完整的死信队列处理流程。

死信队列的核心概念与触发条件
所谓死信,就是那些被消费者拒绝、无法被正常投递的消息。一般来说,消息进入死信队列有三种典型情况:第一种是消费者显式拒绝(reject或nack)且不要求重新入队;第二种是消息超过设定的重试次数仍然失败;第三种是消息在队列中存活时间超过TTL(Time To Live)限制。理解这三种触发条件,是设计重试策略的前提。
很多团队在实践中容易走两个极端:要么完全不设重试,一次失败就进死信,导致网络抖动这类瞬时故障也被放大成事故;要么无限重试,消费端反复抛错,把CPU和日志系统全部打满。比较合理的做法是采用指数退避重试,比如第一次失败后延迟2秒,之后依次是4秒、8秒、16秒,一般重试3到5次仍失败,再进入死信队列。这样既能容忍短时故障,又不会让毒消息(poison message)无限占用资源。
另外需要明确一点:死信队列不是垃圾桶,而是待处理的暂存区。消息进入DLQ之后,必须有配套的消费、告警和补偿机制,否则只是把问题从一个地方挪到另一个地方。
基于BullMQ实现死信队列
BullMQ是Node.js生态中基于Redis的高性能任务队列,它内置了failed队列机制,天然适合做死信处理。下面是一个完整示例,包含生产者、消费者和死信处理三个角色。
const { Queue, Worker } = require('bullmq');
// 创建队列,配置指数退避重试策略
const emailQueue = new Queue('email-queue', {
connection: { host: '127.0.0.1', port: 6379 },
defaultJobOptions: {
attempts: 5, // 最多尝试5次
backoff: {
type: 'exponential', // 指数退避
delay: 2000 // 初始延迟2秒
},
removeOnComplete: 100, // 完成后保留最近100条
removeOnFail: false // 失败任务保留,充当死信
}
});
// 添加任务
async function produce() {
await emailQueue.add('send-email', {
to: 'user@ipipp.com',
subject: '订单确认',
body: '您的订单已支付成功'
});
}
// 消费者:处理业务逻辑
const worker = new Worker('email-queue', async (job) => {
const res = await fetch('http://127.0.0.1:3000/api/send', {
method: 'POST',
body: JSON.stringify(job.data),
headers: { 'Content-Type': 'application/json' }
});
if (!res.ok) {
throw new Error(`发送失败,状态码: ${res.status}`);
}
}, { connection: { host: '127.0.0.1', port: 6379 } });
// 监听失败事件,超过重试上限即视为死信
worker.on('failed', async (job, err) => {
console.error(`任务 ${job.id} 第 ${job.attemptsMade} 次失败: ${err.message}`);
if (job.attemptsMade >= job.opts.attempts) {
// 这里触发告警并写入死信记录表
await notifyOps(`死信产生: ${job.id}`, job.data);
}
});
上面的代码有几个关键点值得展开。首先是attempts: 5配合backoff,构成了完整的重试语义,BullMQ会自动把失败任务按延迟重新调度。其次是removeOnFail: false,确保失败任务不会从Redis中删除,这些任务可以通过queue.getFailed()方法批量取出,形成事实上的死信集合。
死信的回放补偿也很简单,BullMQ提供了重试指定任务的API,配合一个定时巡检的调度器即可:
const { Queue } = require('bullmq');
const queue = new Queue('email-queue', {
connection: { host: '127.0.0.1', port: 6379 }
});
// 定时巡检死信,人工确认后批量回放
async function replayDeadLetters() {
const failedJobs = await queue.getFailed(0, 50); // 取50条
for (const job of failedJobs) {
await job.retry(); // 重新入队
console.log(`已回放任务 ${job.id}`);
}
}
setInterval(replayDeadLetters, 60 * 1000);
回放之前一定要做业务校验,比如先检查下游服务健康状态、校验消息体是否合法,否则回放只是重复失败。建议给死信加上人工审批环节,通过管理后台逐条或批量确认后再执行job.retry()。
基于RabbitMQ的经典死信队列方案
如果你的消息需要跨语言消费,或者对消息可靠性有更高要求,RabbitMQ是更常见的选择。它的死信机制是通过队列声明中的dead-letter-exchange参数实现的:业务队列把死信转发到一个指定的交换机,再由该交换机路由到死信队列。
const amqp = require('amqplib');
async function main() {
const conn = await amqp.connect('amqp://127.0.0.1:5672');
const ch = await conn.createChannel();
// 声明死信交换机和死信队列
await ch.assertExchange('dlx.exchange', 'direct', { durable: true });
await ch.assertQueue('dead-letter-queue', { durable: true });
await ch.bindQueue('dead-letter-queue', 'dlx.exchange', 'order.dead');
// 声明业务队列,绑定死信交换机
await ch.assertQueue('order-queue', {
durable: true,
arguments: {
'x-dead-letter-exchange': 'dlx.exchange',
'x-dead-letter-routing-key': 'order.dead',
'x-message-ttl': 60000 // 消息60秒未被消费则成为死信
}
});
// 业务消费者
ch.consume('order-queue', async (msg) => {
try {
await processOrder(JSON.parse(msg.content.toString()));
ch.ack(msg);
} catch (err) {
// requeue为false,消息直接进入死信队列
ch.nack(msg, false, false);
}
});
}
main();
这段配置里,x-dead-letter-exchange指定了死信去向,x-message-ttl定义了消息最大存活时间。当消费者调用ch.nack(msg, false, false)且第三个参数为false时,消息不会重新回到业务队列,而是被投递到死信交换机。需要注意的是,这种拒绝方式本身不带重试语义,如果想要先重试几次再进死信,通常的做法是声明一组带递增TTL的延迟重试队列,或者安装rabbitmq_delayed_message_exchange插件。
死信队列的消费者与普通消费者写法一致,区别在于处理逻辑:这里应该做的是持久化落库、发送告警,而不是简单打日志。例如可以把死信内容连同原始的headers(RabbitMQ会在死信上附加x-death头信息,记录死信原因和来源队列)写入MySQL或MongoDB,供管理后台查询和批量重发。
监控告警与方案选型建议
无论选哪种方案,死信队列都必须配上监控。核心指标包括:死信队列深度(积压量)、死信产生速率、TOP失败原因分布。BullMQ可以直接用queue.getFailedCount()定时上报,RabbitMQ则可以通过Management插件的HTTP API获取队列消息数。一旦死信速率超过阈值,比如每分钟超过10条,就应该立刻触发电话或即时通讯告警,因为这往往意味着下游服务已经出了严重问题。
两种方案怎么选?BullMQ方案的优势是部署轻量、API友好、重试与延迟控制粒度细,适合以Node.js为主的技术栈、任务型场景(发邮件、生成报表、同步数据)。RabbitMQ的优势在于协议级别的死信路由、完善的TTL支持以及跨语言能力,适合多团队、多语言混合、对消息可靠性要求苛刻的场景。中小型Node.js项目直接用BullMQ通常性价比更高,而大型微服务体系里RabbitMQ或Kafka加自建补偿平台会更稳妥。说到底,死信队列的难点不在于代码,而在于流程:谁接收告警、谁判断能否回放、回放失败后如何升级处理,把这些环节定义清楚,才算真正落地了一套可靠的异常消息处理体系。