数据工程师别再卷实时了:我算了笔账,Flink 流的综合成本是 Spark 批的 7 倍

🔑 关键词:数据工程师,Flink,Spark Structured Streaming,实时数仓,流批选型

📖 摘要:一个把 70% 实时链路迁回批处理的数据工程师,把 Flink 与 Spark 的硬件、人力、故障恢复成本摊开算了一遍。附上流式选型 checklist 和 6 条能省一半钱的配置建议。

凌晨 2:47,我的 Flink 作业没挂,但数据全丢了

图片

去年 11 月 7 日,凌晨 2 点 47 分,PagerDuty 响了。Kafka consumer lag 冲到 400 万,还在涨。

我以为作业挂了。登 Flink UI 一看,TaskManager 全绿,checkpoint 成功率 100%,背压正常。但 watermark 卡在 6 小时前,一动不动。

排查了 40 分钟。原因很傻:上游一个订单服务当晚发版,把 create_time 从 13 位毫秒时间戳改成了 10 位秒级。数据本身没报错——反序列化照样过,因为我们用的是 long。但所有事件的时间戳变成了 1970 年 1 月。它们落进了早已关闭的窗口,进去就出不来。

那条链路丢了大概 800 万条订单事件。补数补了一天半,因为下游还有两个流式 join。

这次之后我把整条链路的成本摊开算了一遍,结论有点扎心。


我把账算了一遍:硬件差 12 倍,人力差 5 倍

图片

先摆数字。我们的实时链路大概是这样的:

  • Kafka:3 台 16C64G,2TB NVMe,三副本,保留 7 天
  • Flink 1.18 on K8s:常驻 12 个 TaskManager(4C8G),峰值并行度 50
  • RocksDB 状态后端,状态峰值约 40GB
  • Checkpoint 30 秒一次,每次落 S3 约 40GB,一天 2880 次
  • 故障恢复要走完 checkpoint 加载加状态对齐,实测 6 到 9 分钟

同样的业务逻辑,我用 Spark 3.5 写了一版小时批:单次跑 6 分钟,8 个 executor(4C16G),一天 24 次。折算下来 19.2 core-hour/天。

Flink 这边:48 core 乘 24 小时等于 1152 core-hour/天,这还不算 Kafka 常驻那 48 core。

12 倍。

人力更贵。这条链路一年出过 4 次 P2 以上的事故,每次平均 3 个人乘 3.5 小时,加上事后复盘和补数,一年大概 120 人时。批作业那边同期出过 2 次,重启一下,各 10 分钟。

而且这只是明面上的。真正让我心态崩的是版本升级那次——Flink 1.17 升 1.18,我们规规矩矩做了 savepoint。恢复的时候发现,有几个算子没显式写 uid,自动生成的 hash 变了,状态映射不上。折腾了 5 个小时,最后放弃状态,从 Kafka 最早 offset 重跑。而 Kafka 只保留 7 天。

7 天数据,直接断档。

图片

那天下班我在楼下便利店坐了半小时。不是气,是觉得自己花一年半维护的东西,可靠性还不如一个 cron。


先说清楚:不是所有流都不该上

我不想走到另一个极端说流处理没用。有三个场景我会毫不犹豫上流:

一、延迟有真金白银的代价。 不是老板想看实时,是超过 1 分钟就有钱损失。比如支付风控拦截、广告竞价、库存超卖扣减。这类需求你去问业务方晚 5 分钟行不行,他会反问你那还要它干嘛。

二、强事件时间语义。 会话窗口、乱序到达、要按事件本身发生的时间而不是到达时间聚合。这种批处理做起来会很别扭,不是做不了,是要写一堆取巧逻辑。

三、下游是机器,不是人。 事件驱动的动作——触发通知、调 API、扣额度。这种要求端到端延迟可预期,批处理做不到。

图片

反过来,下面这些我见过太多人硬上流,最后都是自找苦吃:

  • 报表和看板(T+1 的日报改成 5 分钟刷新,没人看)
  • 用户画像标签(本来就按天算,一天算一次完全够)
  • 大部分风控规则(真出事的时候,批处理跑一遍也就几分钟)
  • CDC 同步(这个用流是对的,但那是 Debezium、Flink CDC 这类工具的事,不用你自己写 DataStream)

决定上流之前,我会问的 5 个问题

这几年我给自己定了个 checklist,上流之前必须过一遍,过不了就别上:

  1. 最晚什么时候要数据? 对方回答实时的时候,追问一句晚 5 分钟会怎么样。十次里有七次,答案是那也没关系。
  2. 上游能保证时间语义吗? 有没有数据契约?timestamp 是秒还是毫秒?会不会变?我第一次栽的就是这个。现在我们要求上游 schema 变更必须走注册中心,改类型直接拒。
  3. 状态有多大? 低于 10GB 随便跑。10 到 100GB 考虑增量 checkpoint。超过 100GB,先想想这个状态能不能拆,或者干脆不用状态(改成幂等写加下游去重)。
  4. 有没有人能读 Flink 的 stack trace? 不是有人值班,是凌晨三点被叫起来,能在 20 分钟内判断出是反压、checkpoint 超时、还是 Kafka 分区 leader 切换。这个能力很难招,也很难留。
  5. 回填方案是什么? 出了事怎么补那几小时数据?如果答案是「从 Kafka 重跑」,那 Kafka 保留多久?上游能不能重发?

第 5 个问题最致命。批处理的回填是天然的——改个日期参数重跑就行。流处理没有这个概念,你得自己造一个。

图片


一个不太政治正确的观点:批的可重跑,比流的低延迟值钱

我做了 6 年数据,越来越觉得数据工程的第一性原理不是快,是可复现。

批处理有个天然优势:它是幂等的。同一份输入,跑同一份代码,出同一份输出。错了就回滚、重跑、改逻辑、再重跑。凌晨三点出事,你可以选择先睡觉,第二天早上再处理,天塌不下来。

流处理不一样。状态是活的,时间是流动的,Kafka 是有保留期的。你错过了那个窗口,就必须去补一个更难查的坑。而且流的 bug 特别隐蔽——它不报错,它静默地少算一点,你要过两周看报表才发现同比少了 3%。

我现在的默认立场是:能用批就用批,除非业务能说清楚延迟的代价。

这个立场让我跟产品经理吵过架。有次一个实时大屏,项目做了三个月,上线一周后访问量掉到每天 4 次。后来改成每小时刷一次,成本降了 90%,访问量一次没变。


图片

如果你非得做流,这 6 件事能省一半钱

不是劝退。真要做的话,有几个具体做法:

  1. 从微批开始。 Spark Structured Streaming 设 trigger(processingTime='1 minute'),先跑一个月。90% 的场景够用,代码量比 Flink DataStream 少得多,调试也简单。真不够再迁移。
  2. 状态能不用就不用。 优先用 Kafka 的 key 分区做局部聚合,而不是 Flink 的 keyed state。实在要用,state.ttl 一定要设,我们设的是 4 小时。
  3. checkpoint 间隔别太短。 30 秒是很多教程的默认值,但如果你状态有 40GB,每 30 秒往 S3 写 40GB,光存储账单就够呛。我们后来改成 3 分钟,故障恢复多花 2 分钟,S3 费用降了 80%。
  4. 算子显式写 uid。 这条是用 7 天数据换来的。所有 KeyedProcessFunctionAggregateFunction 都加上 .uid("order_agg_v1"),改名不改 uid。
  5. 开 unaligned checkpoint。 Flink 1.11 之后就有,反压时能显著缩短 checkpoint 时间。我们开了之后,checkpoint 超时从一天 2 次降到一周 0 次。
  6. 把回填做成常规功能,不是救火手段。 我们搭了个小工具,给定时间范围和 offset,能重放 Kafka 到另一个 topic 再灌进去。做这个花了两周,后面救了三次命。

最后说个数字

我们把那条实时链路里 70% 的计算挪回批处理之后,一年的云成本从大约 41 万降到 13 万,事故从 4 次降到 1 次。留下的是真正需要流的那部分:风控实时拦截和订单状态推送。

剩下的,让 Spark 每小时跑一次,挺好。