上一篇把事务消息和最终一致讲清楚后,消息系统已经从"能跑"进到了"业务依赖它"。一旦业务依赖它,它就变成了需要持续观测和提前规划的基建设施。本篇回答两个问题:监控该盯哪些指标、告警阈值怎么定;以及容量规划怎么算,避免上线半年后磁盘写满才后悔没留余量。这两个能力是消息系统从"能用"走向"敢依赖"的最后门槛。
一、监控指标与告警:先盯 Lag,再看资源
消息系统要监控的指标可以分成两类:业务侧看 Lag(消费积压)和吞吐速率,资源侧看磁盘、CPU、网络。其中 Lag 是最重要的业务健康度指标——它直接反映"消费有没有掉队"。但告警不能只设一个笼统的"Lag 大就报警",要分级:轻微滞后提示关注,严重滞后要立即处理。同时还要看"趋势":入流量持续大于出流量,即使当前 Lag 还不大,也预示着积压即将到来。下面实现一套分级告警规则。
class Metrics:
def __init__(self, lag=0, in_rate=0, out_rate=0, disk_pct=40):
self.lag = lag
self.in_rate = in_rate
self.out_rate = out_rate
self.disk_pct = disk_pct
def check_alerts(m):
alerts = []
if m.lag > 10_000:
alerts.append("consumer lag > 10000, 消费严重落后")
elif m.lag > 2_000:
alerts.append("consumer lag > 2000, 需关注")
if m.disk_pct > 85:
alerts.append("磁盘使用率 > 85%, 接近写满")
if m.in_rate > m.out_rate * 1.5:
alerts.append("入流量 > 出流量 1.5 倍, 积压在增长")
return alerts
samples = [
Metrics(lag=500, in_rate=200, out_rate=100, disk_pct=40),
Metrics(lag=12000, in_rate=800, out_rate=100, disk_pct=90),
]
for i, m in enumerate(samples, 1):
print(f"样本{i} lag={m.lag} in={m.in_rate} out={m.out_rate} disk={m.disk_pct}%")
for a in check_alerts(m):
print(" ALERT:", a)
运行输出:
样本1 lag=500 in=200 out=100 disk=40%
ALERT: 入流量 > 出流量 1.5 倍, 积压在增长
样本2 lag=12000 in=800 out=100 disk=90%
ALERT: consumer lag > 10000, 消费严重落后
ALERT: 磁盘使用率 > 85%, 接近写满
ALERT: 入流量 > 出流量 1.5 倍, 积压在增长
样本 1 的 Lag 只有 500,但入流量是出流量的 2 倍,趋势告警已经亮起——这正是"趋势比绝对值更早报警"的价值。样本 2 则是典型的全面恶化:Lag 严重、磁盘告急、趋势恶化同时出现。这里有一个运维常识:磁盘告警阈值要设得比业务直觉更保守。因为消息系统一旦磁盘写满,Broker 会拒绝写入,生产端直接报错,而且清理积压需要时间,从"85% 报警"到"写满宕机"往往比想象中快。建议磁盘在 70% 就开始预警,85% 进入处理流程。
二、容量规划:吞吐量、副本与保留时间的乘法
容量规划的公式不复杂,难的是把所有变量都考虑进去。磁盘占用 = 吞吐量 × 平均消息大小 × 副本因子 × 保留时间,再乘一个安全余量。很多人只算了"每秒发多少条",却漏了副本因子(3 副本就是 3 倍存储)和保留时间(7 天就是 604800 秒)。下面实现一个容量估算器,看三个规模档位的存储、网络和分区需求。
def capacity_plan(throughput_per_sec, avg_msg_kb, replica_factor,
retention_hours, buffer_ratio=1.5):
bytes_per_sec = throughput_per_sec * avg_msg_kb * 1024
storage_gb = (bytes_per_sec * retention_hours * 3600 * replica_factor) / (1024 ** 3)
storage_gb *= buffer_ratio
network_mbps = (bytes_per_sec * replica_factor) / (1024 ** 2) * 8
per_partition_throughput = 50
partitions = max(1, round(throughput_per_sec / per_partition_throughput))
return storage_gb, network_mbps, partitions
cases = [
("小流量", 100, 1, 3, 24),
("中流量", 1000, 2, 3, 72),
("大流量", 10000, 5, 3, 168),
]
for name, tps, kb, rf, rh in cases:
storage, net, parts = capacity_plan(tps, kb, rf, rh)
print(f"{name}: tps={tps}, msg={kb}KB, 副本={rf}, 保留={rh}h")
print(f" -> 存储≈{storage:.1f}GB, 网络≈{net:.1f}Mbps, 分区≈{parts}")
运行输出:
小流量: tps=100, msg=1KB, 副本=3, 保留=24h
-> 存储≈37.1GB, 网络≈2.3Mbps, 分区≈2
中流量: tps=1000, msg=2KB, 副本=3, 保留=72h
-> 存储≈2224.7GB, 网络≈46.9Mbps, 分区≈20
大流量: tps=10000, msg=5KB, 副本=3, 保留=168h
-> 存储≈129776.0GB, 网络≈1171.9Mbps, 分区≈200
三个数字暴露了容量规划里最容易被低估的两点。一是副本因子的乘法效应:同样 3 副本,小流量只需 37GB,大流量却要 126TB,其中三分之二是副本开销。二是保留时间是静默的存储杀手:中流量从 24h 提到 72h,存储直接翻 3 倍。分区数则要反过来看——它决定了消费并行的上限,规划时宁可多留,因为 Kafka 增加分区容易、减少分区困难,而分区过少会卡死消费组的扩容空间。
容量规划还要和监控联动:规划出的数字要落成告警阈值(磁盘、网络、Lag),否则规划只是纸面数字。下一篇作为系列收官,把前面九篇的所有知识点压缩成一张生产环境落地清单,给出上线前必须逐项确认的检查表和故障演练方法。
三、更多指标与扩容策略
Lag 和吞吐是最核心的指标,但还不够。消费延迟(一条消息从生产到被消费的时间)反映的是端到端体验,Lag 相同但消息大小不同,延迟可能差很多;Broker 连接数和网络连接复用情况反映客户端有没有正确复用连接,连接泄漏会悄悄耗尽 Broker 的文件描述符;死信队列的积压数则反映"有多少消息在反复失败",是业务健康度的先行指标。这些指标要集中到统一的监控看板(Prometheus + Grafana 是主流组合),并和告警规则绑定。
告警的难点不是太少而是太多。每个指标都设一个阈值,晚上就会收到一堆没人看的告警,最后真正重要的告警被淹没。降噪的原则是"告警要可行动":每条告警都应该对应一个明确的处理动作,没有处理动作的告警就不要设,或者只记日志。可以用"分级 + 聚合":严重告警立即通知值班,轻微告警聚合到日报,趋势告警只在持续恶化时触发。
扩容策略要和分区模型匹配。Kafka 增加分区容易(但要重算哈希,可能破坏顺序)、减少分区困难,所以分区数要按"未来一两年峰值"提前规划,而不是按当前量设。Broker 扩容要关注数据再均衡——新节点不会自动分担老节点的存储,需要手动触发分区迁移。容量规划里还要留出"峰值倍数":平时 1000 tps,大促可能冲到 5000 tps,容量要按峰值算,而不是按平均值算,否则一到活动就全线告警。
落地工具上,Kafka 生态有 Kafka Exporter、Burrow 这类现成的指标采集器,把 Lag、吞吐等暴露成 Prometheus 指标;RabbitMQ 内置管理插件能导出指标;RocketMQ 和 Pulsar 也有各自的监控组件。不要从零自己造监控,先接现成的 exporter,再补业务自定义指标。看板建议分两层:一层给运维看资源(磁盘、网络、CPU),一层给业务看健康(Lag、延迟、死信),避免两拨人挤在同一张图里互相干扰。
监控告警最后要落到值班响应上:告警发出去没人处理,等于没有告警。每类告警都要有对应 runbook(处理手册),写清楚"这条告警意味着什么、第一步查什么、怎么恢复",否则告警只会变成噪音。
参考来源
- Kafka:运维与监控
- RabbitMQ:监控与度量
- AWS:Kafka 容量规划最佳实践
👍 觉得有用就点个 赞 + 收藏,方便回头查阅;有疑问直接在评论区留言,我看到都会回。
🚀 本文属于 《消息队列实战》 系列,持续更新,关注不迷路。
📌 文章里的代码都能直接跑。想要可直接 clone 的完整工程 + 配套部署脚本 / 踩坑清单?评论一声或发邮件到 cj2664@qq.com,我免费发你。 如果你正好在做类似系统、或有工程化难题想找人做,也欢迎邮件聊一句——我按实际情况评估,能落地的就接单或出方案。评论和邮件都能直接找到我,不用跳别的平台。

