You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

RabbitMQ未接收全部消息问题排查求助

RabbitMQ批量发送消息丢失问题排查与解决

问题描述

使用NodeJS脚本读取目录下129k个文件,通过fs.readdir获取文件名后遍历调用dispatchMessage向RabbitMQ队列发送消息,但队列仅存储3k条消息。后续测试推送15k条消息时,Broker仅接收14754条,确认问题出在Broker未接收全量消息,与消费者无关。

当前代码:

files.forEach((fileName) => {
    dispatchMessage(channel, fileName)
    dispatchedMessagesCount++
})

const dispatchMessage = (channel, targetPath) => {
    try {
        channel.sendToQueue(queue, Buffer.from(targetPath), {
            persistent: true,
            expiration: 10800000,
        })

        logEvent(`Published %s`, targetPath)
    } catch (error) {
        logEvent(`Error publishing  %s`, targetPath)
    }
}

已尝试设置3小时TTL,控制台无错误输出,怀疑RabbitMQ入站速率限制但未找到相关设置。

问题根源

NodeJS的AMQP客户端(如amqplib)默认会限制通道的未确认消息数量(prefetch count)。当未确认消息达到上限时,客户端会暂停发送新消息,直到收到Broker的确认。同步遍历发送的方式没有等待消息确认,导致部分消息被客户端缓冲但未实际发送到Broker。此外,若队列设置了长度或字节数限制,也会导致超出部分被丢弃。

解决方法

1. 启用消息确认并控制发送速率

修改为异步遍历,等待每条消息的确认后再发送下一条,确保消息被Broker接收:

// 设置通道预取数,控制未确认消息上限
channel.prefetch(100);

const dispatchMessages = async (channel, files) => {
    for (const fileName of files) {
        await new Promise((resolve, reject) => {
            channel.sendToQueue(queue, Buffer.from(fileName), {
                persistent: true,
                expiration: 10800000,
            }, (err) => {
                if (err) {
                    logEvent(`Error publishing %s`, fileName);
                    reject(err);
                } else {
                    logEvent(`Published %s`, fileName);
                    dispatchedMessagesCount++;
                    resolve();
                }
            });
        });
    }
};

// 执行批量发送
dispatchMessages(channel, files).catch(err => console.error(err));

2. 批量发送并监听通道缓冲状态

调大预取数,同时监听drain事件,当通道缓冲已满时等待恢复后继续发送,兼顾效率与可靠性:

channel.prefetch(1000);

const dispatchMessage = (channel, targetPath) => {
    return new Promise((resolve) => {
        const send = () => {
            // sendToQueue返回false表示通道缓冲已满
            const canSend = channel.sendToQueue(queue, Buffer.from(targetPath), {
                persistent: true,
                expiration: 10800000,
            });
            if (canSend) {
                logEvent(`Published %s`, targetPath);
                dispatchedMessagesCount++;
                resolve();
            } else {
                // 等待通道 drain 事件,恢复发送能力后重试
                channel.once('drain', () => send());
            }
        };
        send();
    });
};

// 分批次并行发送,控制并发量
const batchSize = 1000;
for (let i = 0; i < files.length; i += batchSize) {
    const batch = files.slice(i, i + batchSize);
    await Promise.all(batch.map(fileName => dispatchMessage(channel, fileName)));
}

3. RabbitMQ端配置检查

  • 登录RabbitMQ管理控制台,查看对应通道的Unconfirmed消息数,确认是否达到预取上限。
  • 检查队列参数,若设置了x-max-length或x-max-length-bytes限制,需调整或移除该配置。
  • 查看Broker日志,确认是否因磁盘空间不足、内存超限触发流控导致消息丢弃。

内容的提问来源于stack exchange,提问作者Félix Daniel

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.22 06:30:27