跳到主要内容

Flink 理解

· 阅读需 4 分钟

在大数据系统的发展过程中,一直存在一个核心问题:数据越来越多,但处理速度越来越慢。 这篇文章从批处理讲起,梳理 Flink 为什么出现、到底解决了什么问题。

从批处理到流计算

传统的数据处理方式,大多数是离线处理:每天晚上跑一次任务,统计用户行为、计算广告数据、生成报表、分析业务指标。这种方式叫 Batch Processing(批处理),典型工具例如:

  • Hadoop MapReduce
  • Hive
  • Spark

这些系统适合处理海量历史数据。但是随着互联网的发展,很多业务开始需要实时数据处理,例如:

  • 实时风控
  • 实时推荐
  • 实时监控
  • 实时广告竞价
  • 实时日志分析

这些场景有一个共同特点:数据必须“边产生边处理”,不能等到第二天。于是就出现了一个新的计算模式:Stream Processing(流式计算)

早期实时计算系统主要依赖 Storm 和 Spark Streaming,但这些系统都有一些问题。

Storm 虽然是流计算系统,但开发复杂、容错能力一般、状态管理困难。Spark Streaming 虽然稳定,但它的本质是微批处理(Micro Batch),也就是说,它并不是真正的流计算。例如:

每1秒处理一批数据

这种模式在很多场景下仍然不够实时。于是 Flink 出现了,它的目标很明确:构建一个真正的流计算系统。

1) 实时数据处理

Flink 可以实现真正的毫秒级实时计算。以用户点击商品为例:

点击 → Kafka → Flink → 实时推荐

整个过程可能只需要几十毫秒

2) 状态管理

在流计算中,一个非常重要的问题是状态(State)。例如统计用户点击次数:

用户A 点击 1次
用户A 点击 2次
用户A 点击 3次

系统需要记住之前的数据。Flink 内置了状态管理机制(Keyed State、Operator State),并且支持本地状态和 RocksDB 状态存储,这样可以处理超大规模状态数据

3) 容错

在分布式系统中,节点宕机是常态。Flink 通过 Checkpoint 机制来保证数据安全:系统会定期保存计算状态,如果节点宕机,可以从 Checkpoint 恢复计算。这种机制保证了 Exactly Once(精准一次处理)——每条数据只会被处理一次。

4) 事件时间

在流计算中,还有一个非常难的问题:事件时间(Event Time)。日志数据可能会延迟到达,如果只按系统时间计算,就会出现统计错误。

Flink 引入了 Watermark(水位线)机制。Watermark 可以帮助系统判断哪些数据已经“基本到齐”,这样就可以正确处理延迟数据、窗口计算和时间聚合。

1) Stream First

Flink 的核心理念是一切都是流,批处理只是有限数据流。因此,Flink 的 Batch 和 Stream 使用同一套引擎。

2) State Driven

Flink 是一个状态驱动的计算系统,所有计算都围绕状态展开。用户行为统计、实时聚合、实时监控,这些都依赖状态管理。

3) Exactly Once

Flink 非常强调数据处理的正确性。通过 Checkpoint、分布式快照和两阶段提交,Flink 可以实现 Exactly Once 语义,这是很多实时系统非常关键的能力。

Flink 在很多互联网公司都有大量应用,典型场景包括:实时推荐系统、用户行为分析、广告点击统计、实时风控系统、日志实时分析。一条常见的数据链路是:

用户行为 → Kafka → Flink → ClickHouse

Flink 在中间负责数据清洗、实时计算和实时聚合。

小结

Flink 的出现,本质上是为了做“真正的流计算”这件事。它用同一套引擎统一了批和流,用内置状态管理支撑大规模有状态计算,用 Checkpoint 机制换来 Exactly Once 的正确性,再用 Watermark 处理乱序和延迟数据。理解了这几点,也就理解了 Flink 和早期实时计算系统的根本差别。后面再去看具体的 API 和集群部署,会顺畅很多。

评论 / COMMENTS