RabbitMQ消费者遇406确认超时错误,如何处理避免Node重启?
解决RabbitMQ消费者ACK超时导致通道关闭的问题
错误原因分析
你遇到的PRECONDITION-FAILED - delivery acknowledgement on channel ... timed out错误,本质是RabbitMQ在配置的超时时间内(这里是30分钟)没有收到消费者的消息确认(ACK),因此主动关闭了通道。常见触发场景:
- 消息处理逻辑(
appController.exec(data))执行时间过长,超过了RabbitMQ的消费者超时阈值 - 代码异步处理存在漏洞,导致ACK无法及时发送
- 消费者同时处理的消息过多,负载过高导致部分任务卡住
具体解决方案
1. 优化消息处理耗时
- 排查
appController.exec(data)的业务逻辑:检查是否有慢数据库查询、阻塞式IO、外部API调用超时等情况,针对性优化(比如加缓存、异步化非核心操作、拆分大任务),确保单条消息处理时间控制在超时阈值内。 - 如果业务确实需要超长处理时间,可以调整RabbitMQ的超时配置:
- 全局修改:在RabbitMQ配置文件中设置
consumer_timeout参数(比如改成1小时:consumer_timeout = 3600000) - 单消费者设置:在声明消费者时通过
arguments指定超时,示例:channel.consume(service.channel, async (message) => { // 处理逻辑 }, { noAck: false, arguments: { 'x-consumer-timeout': 3600000 } // 1小时超时 })
- 全局修改:在RabbitMQ配置文件中设置
2. 修复代码异步逻辑问题
- 替换
forEach为for...of:forEach不会等待异步函数执行完成,会导致多个通道初始化逻辑混乱,换成for...of确保顺序执行:// 原代码的services.forEach(async (service) => { ... }) 替换为: for (const service of services) { // 通道创建和消费逻辑 } - 给消息处理加超时限制:用
Promise.race确保即使exec卡住,也能及时处理,避免ACK超时:const timeoutPromise = new Promise((_, reject) => { setTimeout(() => reject(new Error('Message processing timeout')), 1700000); // 比RabbitMQ超时少100秒 }); try { await Promise.race([appController.exec(data), timeoutPromise]); channel.ack(message); } catch (error) { logger.error(`Process message failed: ${error.message}`); channel.nack(message, false, true); // 重新入队,避免丢失 }
3. 添加通道重连机制
当前代码没有处理通道关闭的情况,一旦通道被RabbitMQ关闭,消费者就会失效。给通道添加close事件监听,自动重新初始化消费:
// 把通道初始化逻辑抽成单独函数 async function setupConsumer(connection, service, app) { const channel = await connection.createChannel(); await channel.assertQueue(service.channel, { durable: false }); const appController = app.get(service.controller); channel.prefetch(10); // 调小prefetch值,降低并发负载 channel.consume(service.channel, async (message) => { // 消息处理逻辑(带超时的版本) }, { noAck: false }); // 监听通道关闭事件,自动重连 channel.on('close', async () => { logger.warn(`Channel ${service.channel} closed, reconnecting...`); await setupConsumer(connection, service, app); }); channel.on('error', (err) => { logger.error(`Channel ${service.channel} error: ${err.message}`); }); logger.log(`Watching channel ${service.channel} ...`); } // 在bootstrap中调用 amqp.connect(process.env.RABBITMQ_URL) .then(async (connection) => { for (const service of services) { await setupConsumer(connection, service, app); } })
4. 调整prefetch参数
当前channel.prefetch(30)意味着RabbitMQ会一次性推送30条消息给消费者,如果这些消息都耗时较长,很容易导致多个任务同时卡住,触发ACK超时。建议调小这个值(比如5-10),减少并发处理的消息数量,降低消费者负载。
5. 完善错误处理逻辑
- 替换错误时的
channel.ack为channel.nack:错误的消息直接确认会导致丢失,用nack并设置重新入队(channel.nack(message, false, true)),让消息可以重新被消费。但要注意添加重试次数限制,避免不可恢复错误导致死循环(比如给消息加x-death头判断重试次数)。 - 给
JSON.parse单独加错误处理:避免消息格式错误导致整个消费者逻辑崩溃:let data; try { data = JSON.parse(message.content.toString()); } catch (parseErr) { logger.error(`Invalid message format: ${parseErr.message}`); channel.nack(message, false, false); // 不重新入队,直接丢弃或转死信队列 return; }
修改后的完整示例代码
const services: Service[] = [ { channel: "post-now", controller: PostNowController, }, { channel: "post-queue", controller: PostQueueController, }, ]; async function setupConsumer(connection, service, app) { try { const channel = await connection.createChannel(); await channel.assertQueue(service.channel, { durable: false }); const appController = app.get(service.controller); channel.prefetch(10); channel.consume(service.channel, async (message) => { if (!message) return; let data; try { data = JSON.parse(message.content.toString()); } catch (parseErr) { logger.error(`[${service.channel}] Invalid message format: ${parseErr.message}`); channel.nack(message, false, false); return; } const timeoutPromise = new Promise((_, reject) => { setTimeout(() => reject(new Error('Processing timeout')), 1700000); }); try { await Promise.race([appController.exec(data), timeoutPromise]); channel.ack(message); } catch (processErr) { logger.error(`[${service.channel}] Process message failed: ${processErr.message}`); // 可以在这里判断是否是超时错误,决定是否重新入队 channel.nack(message, false, processErr.message !== 'Processing timeout'); } }, { noAck: false }); channel.on('close', async () => { logger.warn(`[${service.channel}] Channel closed, reconnecting...`); await setupConsumer(connection, service, app); }); channel.on('error', (err) => { logger.error(`[${service.channel}] Channel error: ${err.message}`); }); logger.log(`[${service.channel}] Started watching...`); } catch (err) { logger.error(`[${service.channel}] Setup consumer failed: ${err.message}`); // 失败后延迟重试 setTimeout(() => setupConsumer(connection, service, app), 5000); } } async function boostrap(services: Service[]) { try { const app = await NestFactory.createApplicationContext(AppModule); const connection = await amqp.connect(process.env.RABBITMQ_URL); // 监听连接断开事件 connection.on('close', () => { logger.error('RabbitMQ connection closed, reconnecting...'); // 重新初始化所有消费者 boostrap(services); }); connection.on('error', (err) => { logger.error(`RabbitMQ connection error: ${err.message}`); }); for (const service of services) { await setupConsumer(connection, service, app); } } catch (error) { logger.error(`Bootstrap failed: ${error.message}`); setTimeout(() => boostrap(services), 5000); } } boostrap(services);
内容的提问来源于stack exchange,提问作者Azickri
相关产品推荐
相关产品推荐

