Kafka / Pulsar 数据越来越多,磁盘却越来越小?聊聊实时流平台的容量与回溯

简介: Kafka / Pulsar 数据越来越多,磁盘却越来越小?聊聊实时流平台的容量与回溯

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 和数据平台实战。

目录
相关文章
|
5天前
|
人工智能 运维 BI
阿里云千问办公QwenWork深度解析:基于Qwen3.8,六大核心能力重构企业全自动化工作流与计费选型指南
传统AI办公工具大多停留在对话问答、文档摘要、简单文案生成层面,只能完成单点碎片化任务,无法自主拆解复杂业务流程,很难串联多工具、多文档、外部业务系统完成端到端完整工作交付。很多企业在落地AI办公的时候,需要组合多款不同工具,来回切换界面,手动复制粘贴中间结果,智能化改造落地门槛居高不下。千问办公QwenWork是整合多款智能体产品能力打造的一体化企业办公智能体平台,底层基座依托Qwen3.8大模型,打通桌面端Agent、云端Agent、企业协同Agent三种运行形态,不再局限简单问答,接收业务目标之后自主拆解任务步骤,调用各类工具,处理文档、表格、浏览器自动化、数据查询,直接输出可交付的办公
1472 0
|
5天前
|
人工智能 自然语言处理 安全
阿里云AI数智鉴密:AI 生成内容如何拿到一张"防篡改的身份证"
隐形水印 + C2PA签名:让AI生成内容“持证上岗”。
1127 0
|
14天前
|
人工智能 自然语言处理 安全
阿里云千问办公、Qoder Teams、Qoder CN区别与选择指南:模型能力、适用场景与最新活动参考
本文聚焦阿里云2026年推出的三款自研AI办公产品,清晰拆解千问办公、Qoder Teams、Qoder CN的差异化定位与能力边界:千问办公主打职场全场景提效,支持自然语言指令一键完成PPT生成、数据分析等高频办公任务;Qoder Teams面向程序员团队,深度整合AI代码生成、团队协同与企业知识库能力;Qoder CN则专为金融、政务等强合规场景打造,实现数据不出境与VPC私有化部署。文章同步给出分场景选型指南与最新活动定价,帮助不同类型的企业按需组合产品,实现业务岗、研发岗与强合规场景的AI能力全覆盖。
3767 4
阿里云千问办公、Qoder Teams、Qoder CN区别与选择指南:模型能力、适用场景与最新活动参考
|
5天前
|
人工智能 安全 前端开发
刚刚 GPT-6 Astra 发布,全球最强,AGI 时代到来!
OpenAI 正式推出 GPT-6 Astra 模型,带大家看看这次 GPT 有哪些提升,跟 Claude Fable 5.1 有什么差距?AI 编程能力如何?AGI 真的来了么?
630 0
|
2天前
|
SQL 人工智能 前端开发
QoderWake 1.0 正式发布:从桌面里的 Agent,到工作现场的数字员工
QoderWake v1.0正式发布:企业级数字员工团队平台。支持“一句话建岗”,预置10类特训岗位;Waker常驻钉钉/飞书群,@即响应、自动协作、跨任务记忆;具备定时/事件/API多触发方式与统一任务看板;已沉淀27.6万条记忆、12.3万项技能,助力组织实现人机协同增效。
594 0
|
6天前
|
网络协议 Linux iOS开发
【2026实测】Wireshark下载+安装+汉化+使用教程(图文版,巨详细)
Wireshark 是一款免费开源的网络协议分析工具,可实时捕获、解析并可视化数据包,助你诊断网络故障、分析通信协议(如HTTP、DNS、TCP等)。支持Windows/macOS/Linux,含中文界面,新手入门便捷。(239字)

热门文章

最新文章