Kafka / Pulsar 数据越来越多,磁盘却越来越小?聊聊实时流平台的容量与回溯
大家好,我是 Echo_Wish。
做运维这些年,我发现一个特别有意思的现象:
Kafka、Pulsar 平时看着特别稳定,一到业务高峰,大家突然开始研究磁盘。
“昨天磁盘怎么 80% 了?”
“这个 Topic 能不能多保留 7 天?”
“老板说要支持历史数据回放,Kafka 消息是不是不能删?”
“磁盘快满了怎么办?”
然后运维开始疯狂扩容。
扩完以后呢?
过几个月继续扩。
最后 Kafka 变成了一个很贵的“数据仓库”。
其实这事儿真正难的,从来不是 Kafka 磁盘不够,而是一个问题:
实时数据到底应该保留多久?出了问题,我们到底需要回溯多久?
这两个问题如果没想清楚,容量规划基本就是拍脑袋。
一、Kafka 运维最容易踩的坑:把“保留时间”当成了“容量策略”
很多公司的 Kafka 配置是这样的:
log.retention.hours=168
7 天。
看起来挺合理。
但问题来了:
7 天到底需要多少磁盘?
假设某个业务 Topic:
- 每秒写入 20 MB
- 每天数据量约:
20 MB × 60 × 60 × 24
≈ 1.65 TB/天
保留 7 天:
1.65 TB × 7
≈ 11.55 TB
这还只是业务数据本身。
再考虑:
- 副本
- 峰值流量
- Segment
- Broker 磁盘预留
- 数据压缩比例
- Consumer Lag
- Rebalance
- 突发流量
容量很容易继续往上蹿。
如果副本因子是 3:
11.55 TB × 3
≈ 34.65 TB
所以你会发现一个很现实的问题:
Kafka 的容量不是“消息有多少”这么简单,而是“消息量 × 保留时间 × 副本数”。
二、真正应该关注的是“数据生命周期”
我个人非常反对一种 Kafka 运维方式:
所有 Topic 一律保留 7 天。
这其实就是懒。
不同数据,生命周期完全不一样。
比如:
| 数据类型 | 建议保留 |
|---|---|
| 实时监控指标 | 1~3 天 |
| 普通业务事件 | 3~7 天 |
| 订单事件 | 7~30 天 |
| 关键交易事件 | 30 天以上 |
| CDC 数据 | 根据下游恢复能力 |
| 审计日志 | 单独进入日志/对象存储 |
| 临时计算 Topic | 几小时~1天 |
为什么?
因为:
Kafka 不是数据仓库。
Kafka 最擅长的是:
让数据在不同系统之间可靠、实时地流动。
而不是:
“把公司过去三年的所有数据都存起来。”
这两个定位一定要分清楚。
三、容量规划不要问“磁盘还有多少”,要问“还能撑多久”
这是我比较推荐的一个运维指标:
Disk Time To Full
也就是:
按照当前写入速度,磁盘还能撑多久?
例如:
Kafka Broker:
磁盘容量:10 TB
当前使用:7 TB
剩余:3 TB
过去 6 小时平均增长:
200 GB / 小时
那么:
3 TB / 200 GB
≈ 15 小时
这时候你不能说:
“磁盘还有 30% 呢,不急。”
错了。
应该直接告警:
预计 15 小时后磁盘打满。
甚至可以进一步做 Prometheus 告警。
例如:
(
node_filesystem_avail_bytes
/
predict_linear(
node_filesystem_avail_bytes[6h],
86400
)
) < 0
当然,实际生产环境还需要根据磁盘挂载点、采集指标和业务特点调整。
我的观点是:
容量告警一定要从“百分比告警”升级到“时间告警”。
四、别等磁盘 90% 才扩容
这是 Kafka 运维里面特别容易犯的错误。
很多人的告警:
磁盘使用率 > 80%:Warning
磁盘使用率 > 90%:Critical
然后 90% 的时候开始扩容。
问题是:
Kafka 这种系统最怕的就是:
你以为还有 10%,其实根本没有 10%。
因为高峰期流量可能突然翻倍。
例如:
平时:
100 MB/s
大促:
500 MB/s
你按照平时容量计算:
10 TB / 100 MB/s
觉得还能撑很久。
结果业务一上量:
10 TB / 500 MB/s
直接原地起飞。
所以我更建议:
70%:开始容量评估
75%:制定扩容计划
80%:必须完成扩容
85%:限制非核心 Topic
90%:应急处理
不要把 90% 当成“扩容提醒”,而应该把它当成“事故现场”。
五、Kafka 真正厉害的地方:回溯
说到实时流平台,容量只是第一半。
另一半其实是:
回溯。
假设今天上午 10 点:
订单系统正常
↓
Kafka 正常
↓
实时计算正常
↓
10:20 开始出现数据异常
业务发现的时候已经 12 点了。
现在问题来了:
能不能重新计算 10:20~12:00 的数据?
这就是回溯能力。
Kafka 为什么比传统 MQ 更适合做这种事情?
因为消息消费以后,并不是马上从 Kafka 消失。
Consumer 只是移动 Offset。
例如:
Topic: order_event
Partition 0
0 1 2 3 4 5 6 7 8 9
↑
Offset
消费者现在:
offset = 8
并不代表:
0~7 的消息不存在了
只要消息还在 retention 范围内,你完全可以重新消费。
六、回溯千万别直接“改生产 Consumer”
这是很多新手运维特别容易干的事情。
比如生产消费者:
order-consumer
发现:
昨天 14:00~15:00 数据处理有问题
然后直接:
auto.offset.reset=earliest
然后重启。
这时候……
祝你好运。
因为生产 Consumer Group 可能开始重新消费大量历史数据。
结果:
Kafka
↓
Consumer
↓
数据库
↓
重复写入
↓
数据库压力暴增
严重一点:
Kafka Lag ↑
数据库 CPU ↑
接口延迟 ↑
业务报警 ↑
运维人员 ↑
最后从一个数据问题变成整个系统的问题。
七、正确的回溯方式:创建独立 Consumer Group
例如:
生产:
order-consumer-prod
回溯:
order-consumer-replay-20260909
两个 Group 完全隔离。
例如:
kafka-console-consumer.sh \
--bootstrap-server kafka01:9092 \
--topic order_event \
--group order-consumer-replay-20260909 \
--from-beginning
这样做最大的好处:
生产消费不受影响。
回溯任务自己慢慢跑。
甚至可以限制消费速度。
八、但是“从头开始”其实也不是好策略
如果 Topic 有:
10 TB
你只是想回溯:
昨天 14:00~15:00
结果:
--from-beginning
直接把 10 TB 全扫一遍。
这不是回溯。
这是:
给 Kafka 做压力测试。
真正合理的做法是根据时间定位 Offset。
例如:
10:00 → offset 1200000
14:00 → offset 1800000
15:00 → offset 1950000
然后:
从 offset 1800000
消费到 1950000
这才叫精准回溯。
九、Pulsar 的思路也类似,但更灵活
Pulsar 在消息保留和回溯方面,有自己的优势。
它的核心概念是:
Topic
↓
Subscription
↓
Cursor
不同 Subscription 可以拥有不同的消费位置。
例如:
order-sub-prod
负责生产。
order-sub-replay
负责回溯。
所以可以做到:
生产消费
┌──→ prod
Topic ──┤
└──→ replay
这也是我比较喜欢实时流平台的一点:
数据生产和数据消费天然可以解耦。
生产者不用关心:
消费者今天要不要重新算?
消费者也不用要求:
生产者重新发一次消息。
这就是流平台真正的价值。
十、但是回溯不是免费的
这一点很多人容易忽略。
假设:
正常消费速度:
100 MB/s
现在需要回溯:
2 TB
理论上:
2 TB / 100 MB/s
≈ 5.5 小时
如果你把回溯速度提高到:
500 MB/s
虽然只需要:
≈ 1.1 小时
但下游可能直接扛不住。
所以回溯其实是一种:
资源再分配。
你从 Kafka 把数据读出来,只是第一步。
后面还有:
Kafka
↓
Consumer
↓
反序列化
↓
计算
↓
RPC
↓
数据库
↓
缓存
任何一个环节都可能成为瓶颈。
因此回溯系统最好具备:
限速
暂停
恢复
失败重试
进度记录
幂等
监控
十一、我更推荐:热数据 + 冷数据两级架构
如果业务真的要求:
“我要支持半年甚至一年以前的数据重新计算。”
那我不建议:
Kafka 保留 365 天
因为这太贵了。
更合理的是:
实时数据
↓
Kafka / Pulsar
↓
实时计算
↓
对象存储 / 数据湖
例如:
0~7 天
Kafka
7~180 天
对象存储
180 天以后
归档
这样:
Kafka 负责实时。
对象存储负责长期。
出现问题:
最近 7 天
→ Kafka 快速回溯
7 天以前
→ 从对象存储重新构建流
这才是比较健康的架构。
十二、最后聊一个特别容易被忽略的问题:回溯必须考虑幂等
假设:
订单 ID:
ORDER001
第一次处理:
INSERT ORDER001
回溯的时候又来了:
INSERT ORDER001
数据库:
Duplicate Key
如果没有唯一约束,甚至可能:
ORDER001
ORDER001
变成两条。
所以实时系统设计的时候,我一直比较强调:
任何需要回溯的系统,都必须从第一天开始考虑幂等。
例如:
INSERT INTO order_result
(
order_id,
amount,
update_time
)
VALUES
(
@orderId,
@amount,
GETDATE()
);
不要简单认为:
消费一次 = 处理一次
应该设计成:
消息可能重复
↓
业务必须幂等
↓
允许重新消费
↓
允许回溯
这才是真正可靠的实时系统。
十三、我心目中的 Kafka / Pulsar 运维模型
如果让我重新设计一套实时数据平台,我不会只盯着:
CPU
Memory
Disk
我会重点看下面这些指标:
① 写入速率
② 消费速率
③ Consumer Lag
④ 磁盘增长速度
⑤ 磁盘预计耗尽时间
⑥ Topic 数据保留时间
⑦ 峰值流量
⑧ 副本同步状态
⑨ 回溯任务进度
⑩ 回溯对下游的影响
尤其是:
Lag
+
Retention
+
Replay
+
Capacity
这四个东西其实是一套完整的闭环。
最后:实时平台最重要的不是“能存多少”,而是“出问题能不能回来”
很多人做 Kafka 运维,第一反应是:
“磁盘不够了,扩磁盘。”
做久一点以后会发现:
扩磁盘只是最简单的解决方案。
真正应该思考的是:
数据为什么要留?
留多久?
谁会消费?
出了问题能回溯多久?
回溯需要多少资源?
回溯会不会影响线上?
历史数据放 Kafka 还是对象存储?
如果这些问题没有答案,那么:
Kafka 留 30 天和留 300 天,本质上都只是“先存着再说”。
而真正成熟的实时数据平台,应该做到:
实时数据
↓
短期保留
↓
快速消费
↓
异常可回溯
↓
历史数据沉淀
↓
必要时重新计算
所以我一直觉得:
Kafka/Pulsar 运维的核心,从来不是“把消息存住”,而是“让数据在可控成本下,随时有机会重新跑一遍”。
因为线上系统真正可怕的,从来不是今天的数据处理失败。
而是:
失败了以后,你发现根本没有办法回到昨天。
这才是实时数据平台最应该解决的问题。
我是 Echo_Wish,关注我,聊点真正落地的 AI、Python、云原生、DevOps 和数据平台实战。