Flink

Flink 是一个以流为本的分布式计算引擎。它和 Spark 都能做批也能做流,但出发点相反:Spark 把流当成「一小批一小批的批」(micro-batch),Flink 则把批当成「有界的流」——流才是它的原生模型。这个区别决定了 Flink 在低延迟、事件时间处理、状态管理上的天然优势,是实时数仓的主力引擎。

流处理的三个硬骨头

实时计算难,难在这三点,Flink 的设计几乎都围绕它们:

1. 时间语义:事件时间 vs 处理时间
  • 处理时间(Processing Time):数据被处理的那一刻的机器时间。简单,但结果不可复现——同一份数据重跑结果可能不同。
  • 事件时间(Event Time):数据实际发生的时间(埋在消息里的时间戳)。用事件时间才能算出正确、可复现的结果,哪怕数据迟到、乱序。

生产上几乎都用事件时间,因为网络抖动会让数据乱序到达,只有按「事情真正发生的时间」聚合才对得上账。

2. Watermark(水位线)

事件时间的难点是:我怎么知道某个时间窗口的数据到齐了? Watermark 就是答案——它是一个不断推进的时间戳,含义是「比这个时间更早的数据我认为基本都到了」。当 Watermark 越过窗口末尾,就触发该窗口计算。Watermark 设得越宽容忍乱序越强、延迟越高;设得越紧延迟越低、越容易丢迟到数据。这是延迟和准确性的核心权衡。

3. State(状态)与容错

流是无穷的,很多计算(去重、累计、Join)都要记住「之前发生过什么」,这就是状态。Flink 把状态作为一等公民管理:

  • Checkpoint(检查点):周期性给所有算子的状态打一份分布式快照。作业失败后从最近一次 checkpoint 恢复,配合可重放的数据源(如 Kafka 回拨位点),实现 exactly-once
  • Savepoint(保存点):手动触发的快照,用于升级作业、改并行度、迁移集群。
  • 状态后端:小状态放内存/堆,大状态放 RocksDB(落本地磁盘 + 增量 checkpoint)。
Window(窗口)

把无界流切成有界片段来聚合:

窗口类型含义
滚动窗口 Tumbling固定长度、首尾相接、不重叠(每 5 分钟一段)
滑动窗口 Sliding固定长度但按步长滑动、会重叠(每 1 分钟统计近 5 分钟)
会话窗口 Session按活动间隙切分,一段时间没数据就关窗
  • 模型:Flink 是逐条处理的真流;Spark(Structured Streaming)本质是微批,延迟下限受批间隔约束。
  • 延迟:Flink 可做到毫秒级,Spark 微批通常在秒级。
  • 状态与事件时间:两者都支持,但 Flink 的 Watermark、状态管理更成熟细腻。
  • 选型:对延迟极敏感、状态复杂的场景选 Flink;已有 Spark 生态、延迟要求不高、想批流一套代码,Spark 也够用。
背压(Backpressure)

下游处理不过来时,Flink 会自然地把压力沿着算子链反向传导到源头,让上游放慢读取速度,而不是把数据堆爆内存。排查性能瓶颈时,看哪个算子是背压的起点,就找到了瓶颈所在。

在数仓里的位置

Flink 是实时数仓的计算核心:消费 Kafka 里的 ODS 流,做实时清洗、维表关联(Lookup Join)、多流 Join、实时指标聚合,再写入 OLAP 引擎(如 Doris/ClickHouse)或回写 Kafka。它也是流批一体、实时湖仓(配合 Paimon/Iceberg)方案里的主力。

评论 (0)

登录后参与评论。

还没有评论,来做第一个。

登录后可以选中正文添加批注(仅自己可见)。

Flink