最近被RabbitMQ队列堆积折腾过的人,应该都懂那种“半夜被电话叫醒,打开管理后台一看,队列里躺着几十万条消息”的绝望。业务方催着处理,消费者疯狂重启,消息越积越多,最后只能靠临时扩容和手动清队列救火。更尴尬的是,这种事往往不是第一次发生——团队没有一套轻量的监控告警体系,每次只能等用户投诉或者业务指标异常了才后知后觉。所以今天这篇就聊点实战的:怎么用最短的时间,给RabbitMQ补一套队列堆积监控,一旦指标异常,直接往钉钉群和企业微信群推告警。写给自己用、写给小团队用都可以,整体思路不复杂,也不需要引入特别重的监控平台,适合那些“RabbitMQ在用,但监控基本靠肉眼”的朋友直接抄作业。
先说结论:标题里的“一分钟”不是夸张,但也不是说从零搭完所有东西只要一分钟。它指的是,如果脚本骨架你已经有了、机器上Python环境是现成的,替换掉配置里的队列名、Webhook地址、阈值,一分钟内能让告警跑起来。而完整的搭建流程,包括开插件、建监控账号、写脚本、配定时任务、验证告警通路,熟练的话十几分钟也能搞定。这篇会把这些步骤全部拆开讲,包括我踩过的坑和一些细节取舍。
1. 为什么必须盯队列数据堆积?
1.1 堆积不是小事,是连环事故的起点
很多人觉得消息堆积只是“处理慢了点”,顶多延迟几分钟。但等你真正经历过一次,就会发现它是个连锁反应的开端。我用一个生活化的类比解释下:消息队列就像一个快递中转站,生产者是各地发来的包裹,消费者是仓库里的分拣员。平时包裹到一件分拣一件,中转站一直空荡荡的。突然有一天分拣员集体生病(消费者挂了),包裹还在源源不断涌进来,中转站的货架很快塞满,新包裹没地方放,整个系统就堵死了。
映射到技术指标上,堆积的后果很直接:
- 实时业务变离线业务:订单状态、支付回调、积分变动这类消息一旦堆积,用户侧的反馈就变成了“页面一直转圈”“操作半天没反应”,本质上系统还在跑,但业务已经失真了。
- 消息过期被丢弃:如果队列设置了TTL,消息在队列里待太久会自动消失。堆积期间流失的数据是永远找不回来的,对账时就会发现少了一大批记录。
- 内存和磁盘双双告急:RabbitMQ的队列消息默认会驻留内存,堆积量一大,内存先扛不住,接着触发流控,连正常的小流量消息都会被限速。磁盘也一样,持久化消息全写到磁盘上,磁盘满了节点直接拒绝服务。
- 殃及池鱼:这是最阴的。同一台节点上往往跑着多套业务队列,堆积的队列吃光了节点的内存和磁盘资源,其他队列的正常消息也被连累,跟着消费变慢甚至阻塞。你本来只想处理一个队列的问题,结果整个节点的业务都被拖下水。
所以,堆积监控不是“锦上添花”,而是RabbitMQ生产环境的底线配置。没监控的时候,你是等用户投诉了才知道出事了;有监控之后,你可以在堆积刚冒头的时候就把问题摁死。
1.2 光看总量不够,要拆开看Ready和Unacked
RabbitMQ管理后台里每个队列都有两个关键数字:Ready(待消费消息数)和Unacknowledged(已投递但未确认的消息数)。很多人只看消息总数,这是不够的。这两个数字背后对应的是完全不同的故障场景:
- Ready持续上涨:说明消费者没有拉取消息,或者拉取的速度远低于生产速度。原因大概率是消费者进程挂了、阻塞了、或者消费逻辑卡在某个外部调用上。
- Unacked持续上涨:说明消费者一直在接收消息,但处理完后没有确认(ack)。原因可能是消费逻辑抛了异常没捕获、处理超时、或者手动ack模式漏写了确认逻辑。
排查思路上,看到Ready涨,先去看消费者的存活和日志;看到Unacked涨,得去找消费代码里的bug。如果你只盯着总数,两边混在一起,排查效率会低很多。监控指标时最好把这两项分开设阈值,告警信息里也分别标明,方便值班的人第一时间判断方向。
1.3 阈值怎么定?给的参考公式
阈值设得太小容易告警风暴,三分钟响一次,群里全是噪音,最后大家直接把群屏蔽了;设得太大又起不到预警作用。我的经验是参考业务常态的3到5倍:
- 先让脚本跑一周左右,记录各个队列的Ready正常水位(非高峰期均值和高峰期峰值)。
- 堆积阈值取“平时峰值再乘3”,低于这个值一般只是瞬时抖动,不用打扰人。
- 如果队列本身允许短暂堆积,可以再叠加“持续时间”条件:连续N次巡检(比如连续3分钟)都超过阈值才告警,避免偶发毛刺刷屏。
Unacked的阈值更严格一些,这个数字正常情况下应该趋近于0。只要超过几百条持续一分多钟,基本可以断定消费者在“假死”,可以直接把告警阈值定低一点。
2. 监控方案怎么选?为什么我推荐“管理API + 定时脚本”
2.1 先盘点市面上的主流方案
把RabbitMQ监控方案放一起对比,大概有这几种:Prometheus + RabbitMQ Exporter + Grafana的组合、RabbitMQ自带的管理API轮询、商业监控平台。它们各有侧重,我做了个简单对比:
| 方案 | 落地成本 | 实时性 | 灵活度 | 适合场景 |
|---|---|---|---|---|
| Prometheus + Exporter + Grafana | 偏高,需要部署维护一套监控体系 | 高,可精确到秒级抓取 | 高,指标全面 | 已有Prometheus体系的团队,想长期完善监控大盘 |
| 管理API + 脚本轮询 | 极低,一个脚本 + 一个定时任务 | 中,分钟级足够 | 中,告警逻辑完全可控 | 小团队、个人项目、不想引入重组件的场景 |
| 商业监控平台 | 高,需要采购和接入 | 高 | 受平台限制 | 企业级统一运维平台,有合规要求 |
如果你团队里已经跑着Prometheus和Grafana,那直接用Exporter会更正规,告警规则写在Prometheus里,图表也好看。但如果你只是“想快速给RabbitMQ加一道告警护城河”,为这一件事去搭一套Prometheus,成本其实有点划不来——光是维护Exporter本身、处理指标采集异常、配Alertmanager路由,就够喝一壶了。
2.2 管理API轮询的优势,恰恰是“轻”
RabbitMQ从3.x开始内置了Management插件,开启后默认会暴露一套RESTful API,能查队列、连接、消费者、节点状态等几乎所有指标。用Python或Shell定时去拉这个接口,把数值和阈值一比,触发条件就推Webhook,整个过程不侵入业务、不依赖额外中间件。
我选这个方案的核心理由有四个:
- 零新增组件:管理API是RabbitMQ自带的,开启插件就行,不用装Exporter、不用改RabbitMQ配置、不用引入时序数据库。
- 告警逻辑完全在自己手里:阈值、冷却时间、恢复通知、告警文案,全是你说了算,想改就改。Prometheus那条路虽然也行,但规则要写成PromQL,调试起来不如直接对着Python字典改来得直观。
- 内网安全好把控:管理API走内网地址,脚本部署在同一台机器或内网机器上,不额外暴露端口。
- 一分钟能跑起来:管理API返回的JSON结构非常规整,队列的
messages_ready和messages_unacknowledged两个字段直接拿过来比较就行。
当然它也有短板:没有历史趋势图表,看不了“过去一小时堆积量的走势”。但我的观点是,监控告警要解决的核心问题是“出事了第一时间知道”,趋势分析是第二步的事。先把告警通了,后面有余力再考虑接Grafana补图表。
2.3 告警通道选型:钉钉和企业微信双通道
告警推送通道上,我最常用的是钉钉自定义机器人和企业微信机器人。这俩都是官方开放的Webhook能力,配置方式大同小异:在群里添加一个机器人,拿到Webhook地址,往这个地址POST一段JSON,消息就到群里了。
钉钉自定义机器人有个“加签”安全设置,配置后会给你一个密钥,推送时需要把时间戳和密钥拼起来做HMAC-SHA256签名,拼到Webhook地址里。企业微信机器人则靠“关键词”过滤,webhook地址里带上关键词,消息内容里必须包含这个关键词才能发出去。这些安全设置建议都开上,防止Webhook地址泄露后被人乱刷。
之所以推荐“双通道”,是因为群里不一定所有人同时看钉钉或企业微信,有的公司主流办公是钉钉,有的团队习惯挂企业微信。脚本里两个通道都实现,配置里填了哪个就推哪个,兼容性最好,也方便切换。
3. 实操实现:用Python脚本一分钟搭建堆积监控
3.1 前置准备:开启Management插件并创建监控账号
先确认RabbitMQ的管理插件已经开启。如果你是通过官方安装包或Docker部署的RabbitMQ,一般默认没开Management插件,需要手动执行:
rabbitmq-plugins enable rabbitmq_management然后创建一个只读权限的监控账号,别直接用admin账号跑监控脚本。管理API有专门的正则表达式权限控制,监控账号只需要能读队列信息就够了:
rabbitmqctl add_user monitor YourStrongPassword rabbitmqctl set_permissions -p / monitor "^$" "^$" "^.*"这里的set_permissions三个字段分别对应配置、写、读权限。给监控账号配读权限就行,配置留空、写留空、读放行全部的^.*。这样万一脚本账号泄露,别人也删不了队列、发不了消息。
最后验证API是否通:
curl -u monitor:YourStrongPassword http://127.0.0.1:15672/api/overview能看到JSON返回就说明管理API正常。
3.2 完整Python脚本:监控、判断、推钉钉和企业微信
脚本本身不长,我直接给出完整代码,然后逐块拆解。你可以复制下来按需改配置:
#!/usr/bin/env python3 # -*- coding: utf-8 -*- """ RabbitMQ 队列堆积监控与告警脚本 支持:钉钉、企业微信群机器人推送 """ import base64 import hashlib import hmac import json import time import urllib.parse from datetime import datetime import requests # ====== 配置区 ====== RABBITMQ_HOST = "127.0.0.1" RABBITMQ_PORT = 15672 RABBITMQ_USER = "monitor" RABBITMQ_PASSWORD = "YourStrongPassword" RABBITMQ_VHOST = "/" # 默认vhost # 每个队列独立配置阈值:max_ready 表示Ready堆积上限;max_unacked 表示Unacked上限 QUEUE_RULES = [ {"name": "order.delay.queue", "max_ready": 10000, "max_unacked": 1000}, {"name": "order.paid.queue", "max_ready": 5000, "max_unacked": 500}, {"name": "sms.send.queue", "max_ready": 2000, "max_unacked": 200}, ] # 告警通道配置,不需要的通道置空字符串即可 DINGTALK_WEBHOOK = "https://oapi.dingtalk.com/robot/send?access_token=your_token" DINGTALK_SECRET = "your_secret" # 如果没开加签,留空 WECOM_WEBHOOK = "https://qyapi.weixin.qq.com/cgi-bin/webhook/send?key=your_key" # 全局开关 CHECK_INTERVAL = 30 # 巡检间隔(秒),配合cron或循环使用 SCAN_ALL_QUEUES = True # True表示巡检所有队列;False表示只看QUEUE_RULES里列的 # ====== 配置区结束 ====== # 告警状态记录:queue_name -> "ok" / "alerting" ALERT_STATE = {} def rabbitmq_api(path): """调用RabbitMQ Management API并返回解析后的JSON""" url = f"http://{RABBITMQ_HOST}:{RABBITMQ_PORT}{path}" resp = requests.get(url, auth=(RABBITMQ_USER, RABBITMQ_PASSWORD), timeout=5) resp.raise_for_status() return resp.json() def dingtalk_sign(timestamp): """钉钉加签:把时间戳和密钥拼接后做HMAC-SHA256签名""" secret_enc = DINGTALK_SECRET.encode("utf-8") string_to_sign = f"{timestamp}\n{DINGTALK_SECRET}".encode("utf-8") hmac_code = hmac.new(secret_enc, string_to_sign, digestmod=hashlib.sha256).digest() sign = urllib.parse.quote_plus(base64.b64encode(hmac_code)) return sign def send_dingtalk(text): """推送钉钉机器人消息,支持加签模式""" if not DINGTALK_WEBHOOK: return headers = {"Content-Type": "application/json"} if DINGTALK_SECRET: timestamp = str(round(time.time() * 1000)) sign = dingtalk_sign(timestamp) webhook = f"{DINGTALK_WEBHOOK}×tamp={timestamp}&sign={sign}" else: webhook = DINGTALK_WEBHOOK payload = {"msgtype": "text", "text": {"content": text}} try: resp = requests.post(webhook, headers=headers, json=payload, timeout=5) data = resp.json() if data.get("errcode") != 0: print(f"[dingtalk] push failed: {data}") except Exception as exc: print(f"[dingtalk] push error: {exc}") def send_wecom(text): """推送企业微信机器人消息""" if not WECOM_WEBHOOK: return headers = {"Content-Type": "application/json"} payload = {"msgtype": "text", "text": {"content": text}} try: resp = requests.post(WECOM_WEBHOOK, headers=headers, json=payload, timeout=5) data = resp.json() if data.get("errcode") != 0: print(f"[wecom] push failed: {data}") except Exception as exc: print(f"[wecom] push error: {exc}") def check_queue_health(queue_info): """ 判断单个队列是否异常。 返回 (状态, 需要推送的告警文本列表) """ name = queue_info.get("name", "") ready = queue_info.get("messages_ready", 0) unacked = queue_info.get("messages_unacknowledged", 0) consumers = queue_info.get("consumers", 0) rule = None for r in QUEUE_RULES: if r["name"] == name: rule = r break if rule is None and not SCAN_ALL_QUEUES: return None, [] ready_limit = rule["max_ready"] if rule else 10000 unacked_limit = rule["max_unacked"] if rule else 500 alerts = [] status = "ok" if ready > ready_limit: status = "alerting" alerts.append(f" - Ready堆积: {ready}条, 阈值{ready_limit}条") if unacked > unacked_limit: status = "alerting" alerts.append(f" - Unacked异常: {unacked}条, 阈值{unacked_limit}条") now_str = datetime.now().strftime("%Y-%m-%d %H:%M:%S") text_lines = [ f"[RabbitMQ告警] 队列: {name}", f"时间: {now_str}", f"Ready: {ready} / Unacked: {unacked} / 消费者: {consumers}", ] text_lines.extend(alerts) return status, "\n".join(text_lines) def run_check(): """巡检入口:拉取队列列表,逐个判断,按状态机推送""" # 获取所有队列 queues = rabbitmq_api("/api/queues") for queue in queues: # 如果配置里限制了队列列表,只处理需要监控的队列 if not SCAN_ALL_QUEUES: names = [r["name"] for r in QUEUE_RULES] if queue.get("name") not in names: continue status, text = check_queue_health(queue) if status is None: continue queue_name = queue.get("name") prev_state = ALERT_STATE.get(queue_name, "ok") if status == "alerting" and prev_state != "alerting": # 从正常切到告警,推送告警文本 final_text = text + "\n状态: 首次告警,请尽快处理" send_dingtalk(final_text) send_wecom(final_text) ALERT_STATE[queue_name] = "alerting" print(f"[alert] {queue_name} -> alerting") elif status == "ok" and prev_state == "alerting": # 从告警恢复,推送恢复通知 now_str = datetime.now().strftime("%Y-%m-%d %H:%M:%S") recover_text = f"[RabbitMQ恢复] 队列: {queue_name} 已恢复正常。时间: {now_str}" send_dingtalk(recover_text) send_wecom(recover_text) ALERT_STATE[queue_name] = "ok" print(f"[recover] {queue_name} -> ok") if __name__ == "__main__": # 配合crontab使用时跑一次即可;也可以直接while True循环 # 下面两种模式二选一: while True: try: run_check() except Exception as exc: print(f"[run_check] error: {exc}") time.sleep(CHECK_INTERVAL)3.3 代码拆解:每个环节为什么这么写
这个脚本看着不长,但有几个细节值得单独说:
指标采集部分用的是/api/queues接口,一次请求返回所有队列的完整指标,比逐个队列调/api/queues/{vhost}/{name}接口效率高得多。返回的JSON里每个元素包含name、messages_ready、messages_unacknowledged、consumers等字段,直接取值就行。注意如果队列名里带斜杠或特殊字符,要在单独查询时做URL编码,但用/api/queues批量接口就不用操心这个。
状态机设计是整个脚本的灵魂。我用了ALERT_STATE字典记录每个队列上次的状态,只有“正常到告警”和“告警到恢复”这两次状态切换才推消息。这个设计解决的是“告警风暴”问题——没有状态机的话,脚本每隔30秒跑一次,堆积期间群消息能把你手机震到没电。有了状态切换控制,一个队列在故障期间只会收一条告警和一条恢复通知,安静又及时。
恢复通知经常被人忽略,但它比告警还重要。故障恢复之后,值班的人需要在群里看到一个明确的“已恢复”信号,不然会一直悬着心,重复确认好几遍状态。所以脚本里在状态从alerting切回ok时,会再推一条恢复消息。
钉钉加签实现里有个小坑:签名用的Webhook地址需要把timestamp和sign两个参数拼到原Webhook后面,而且签名用的字符串是timestamp + "\n" + secret,不是直接把secret拿去加密。还有,base64.b64encode出来的字节串要decode("utf-8")再去quote_plus,很多人第一次写都卡在这。
超时处理用的是timeout=5,避免RabbitMQ管理API卡住导致脚本长时间阻塞。推消息的请求也套了try-except,Webhook地址万一失效了不能影响主流程继续巡检。
3.4 部署:用systemd或crontab让它定时跑
脚本里写了while True循环加sleep,但这种常驻模式适合用systemd管理。如果你更习惯传统的crontab,也可以把主流程改成单次执行,然后交给cron每分钟调度:
# crontab 方式,每分钟执行一次巡检 * * * * * cd /opt/rabbitmq-monitor && /usr/bin/python3 rabbitmq_monitor.py >> /opt/rabbitmq-monitor/monitor.log 2>&1但这里有个矛盾:脚本里是while True,cron又每分钟调一次,会起一堆重复进程。我建议二选一:
- 用cron调度的话,把
while True循环删掉,只保留run_check()的单次执行逻辑。 - 用systemd常驻的话,配一个简单的service文件,反而更稳,进程崩溃了能自动拉起。
systemd配置大概长这样:
[Unit] Description=RabbitMQ Monitor After=network.target [Service] ExecStart=/usr/bin/python3 /opt/rabbitmq-monitor/rabbitmq_monitor.py Restart=always RestartSec=10 [Install] WantedBy=multi-user.target写完systemctl daemon-reload && systemctl enable --now rabbitmq-monitor就能跑起来。
3.5 验证:模拟堆积看告警能不能收到
部署完别急着走,务必做一次端到端验证。最简单的方式是手动往队列里塞消息模拟堆积:
# 循环往监控的队列里发消息,快速撑过阈值 for i in $(seq 1 20000); do rabbitmqadmin publish exchange=amq.default routing_key=order.delay.queue payload="test-$i" done发完消息等一个巡检周期,看钉钉群或企业微信群有没有收到告警。然后再把消费者恢复(或者直接清队列),确认能不能收到恢复通知。整个链路跑通了,这套监控才算真正落地。
4. 常见问题与排查实录
4.1 管理API返回401或403
这是新手最容易踩的坑。两个原因:一是创建监控账号时权限写错了,记住监控账号的set_permissions第三段(读权限)一定要放行,否则/api/queues接口会拒绝访问。二是API的基础认证试试用IP访问管理后台能通,但脚本怎么都401,大概率是密码里有@、:之类的特殊字符,requests的auth=(user, password)虽然会自动编码,但如果你用的是拼接URL的方式,就要用urllib.parse.quote处理。
4.2 队列名为中文或含特殊字符时查询报错
批量接口/api/queues不受影响,但如果你改用单队列接口/api/queues/{vhost}/{name},vhost和队列名里的/会破坏URL路径。这时候要用urllib.parse.quote对每一段单独编码,比如vhost是/就得编码成%2F,不能直接拼。这也是我为什么推荐直接拉全量队列列表再过滤的原因,省掉一层编码的麻烦。
4.3 钉钉机器人推送报“关键字不匹配”
钉钉自定义机器人有两种安全设置:加签和自定义关键字。如果你没开加签而是设置了“自定义关键字”,比如“告警”或“监控”,而推出去的消息文本里恰好没有这个关键词,钉钉会直接拒绝。解决方法是把告警文案里的固定词设计好,比如统一带上“Rainfall告警”或“RabbitMQ监控”这类包含关键字的词。企业微信机器人同理,关键词不匹配会返回errmsg错误,日志里仔细看就能定位。
4.4 机器人每分钟20条限制导致丢消息
钉钉和企业微信的Webhook机器人都有频率限制,都是“每个机器人每分钟最多20条”。如果多队列同时告警,或者告警状态频繁抖动,很容易触发这个限制。我的处理办法有两个:一是靠状态机去重,让每个队列在故障期间只推一条;二是把多个队列状态汇总成一条消息推出去,比如巡检时把所有异常队列整理成一条“XX个队列异常”的文本,避免逐条推送。
4.5 RabbitMQ从3.x升到4.x后字段对不上
RabbitMQ 4.x的Management API大部分字段和3.x保持一致,messages_ready、messages_unacknowledged还在。但4.x新增了分页参数,返回结构上有些统计数据字段(比如message_stats里的细分计数)有小调整,对只取基础字段的监控脚本影响不大。如果你升级后脚本报错,先去/api/queues接口返回的JSON里人工确认字段名,基本都能排查到。
4.6 推送机器人通知超时或失败
Webhook请求失败最常见的原因是内网机器访问不了外网,或者防火墙把oapi.dingtalk.com、qyapi.weixin.qq.com挡了。两种解决思路:一是给请求挂代理(环境变量HTTP_PROXY或requests的proxies参数);二是如果公司有内网网关,想办法把告警推到内网IM或短信平台。小团队的话,最省事的还是让脚本部署在能访问外网的机器上,RabbitMQ的Python脚本通过内网API拉数据,出网只走Webhook这一条链路。
5. 把监控从“能跑”做成“好用”
脚本跑通只是第一步,实际上我在运维中发现,真正让监控体系“好用”的是下面这些细节。先说说告警文案的设计。一条告警消息如果只有“order.delay.queue堆积了”,值班的人还要去查这个队列是干嘛的、属于哪个业务线、平时水位多少,效率很低。我习惯在配置里给每个队列加一个“业务说明”字段,告警文案里直接带出来,比如“订单延迟队列(订单服务-延迟关单),堆积30000条,超过阈值10000条”。收到告警的人不用再翻资料,能直接判断影响面。
阈值这块也不要一劳永逸。业务有周期波动,比如电商大促期间链路水位明显高于日常。我给每个队列配置的阈值是写死的,但每到活动前我会手动调大一档,活动结束后调回来。如果要做得更自动,可以把阈值抽到独立的配置文件里,按天或按环境加载,不过小团队用脚本改改也够了。
还有两个加分项:一是除了堆积,顺手监控一下消费者数量。如果某个队列的消费者数掉到0了,那就算当前没堆积也要告警,因为这是“即将堆积”的前兆,早处理比晚处理舒服。二是有条件的话把告警接一个“分级”概念,比如超过阈值3倍是P1告警,1倍是P2告警,不同级别的消息用不同文案、推不同群,避免重要事故被淹没在海量普通告警里。
我在实际使用中最深的体会是:监控脚本的代码其实不值钱,真正值钱的是对业务水位的心中有数。你有多少队列、每个队列正常水位是多少、峰值能冲到多高,这些数据比技术方案更稀缺。先让脚本跑起来收集基线,再根据基线调整阈值和告警策略,这套体系才真正贴合自己的业务,而不是套用网上模板写完就完事。最后再分享一个实用小技巧:告警推送到群里之后,可以在文案里带上“恢复请在后台确认”之类的提示,减少群里“这个处理了吗”的追问。监控是手段,目标是让收到告警的人能快速、正确地行动。