Flink快速入门:大数据流处理框架核心概念

前言

如果你是一个后端开发者,可能经常听到"实时计算"、"流处理"这些词。比如双十一大屏上那个不断跳动的成交额数字,又比如你刷短视频时实时推荐的内容,这些背后都有流处理框架的身影。

今天我们要聊的 Apache Flink,就是当下最火的流处理框架之一。这篇文章不会上来就怼代码,而是先把 Flink 的核心概念讲清楚,让你知道它是什么、能干什么、为什么选它。

🏠个人主页:你的主页


目录


一、什么是流处理

在聊 Flink 之前,我们得先搞清楚一个问题:什么是流处理?

批处理 vs 流处理

想象一下你在餐厅吃饭:

  • 批处理:就像自助餐,厨师先把所有菜做好摆出来,你一次性打完所有想吃的菜再去吃。数据是"攒一波再处理"。
  • 流处理:就像火锅,食材一片片下锅,熟了就捞起来吃,边吃边涮。数据是"来一条处理一条"。

用技术语言来说:

特点批处理流处理
数据特征有限数据集(Bounded)无限数据流(Unbounded)
处理时机数据收集完毕后统一处理数据到达即处理
延迟高(分钟/小时级)低(毫秒/秒级)
典型场景报表统计、数据仓库实时监控、实时推荐

举个例子

  • 统计"昨天的网站PV"——用批处理,跑个定时任务就行
  • 统计"当前这一秒的在线人数"——必须用流处理,不然数据早过时了

二、Flink是什么

Apache Flink 是一个开源的分布式流处理框架,由德国柏林工业大学的几个博士生在 2010 年左右开始研发,2014 年成为 Apache 顶级项目。

一句话概括 Flink:

Flink 是一个支持有状态计算的、高吞吐低延迟的、流批一体的分布式计算引擎。

别被这句话吓到,我们拆开来看:

  • 有状态计算:能记住之前处理过的数据,比如统计"过去5分钟的订单总金额"
  • 高吞吐低延迟:处理速度快,每秒能处理百万级事件,延迟在毫秒级
  • 流批一体:一套代码既能跑流处理,也能跑批处理
  • 分布式:能把任务分散到多台机器上并行执行

Flink的发展历程

时间事件
2010柏林工业大学启动 Stratosphere 项目
2014项目更名为 Flink,成为 Apache 顶级项目
2016阿里巴巴开始大规模使用,并贡献 Blink 分支
2019阿里收购 Flink 商业公司 Ververica
2020+Flink 成为流处理事实标准,版本迭代至 1.18+

可以说,Flink 是被阿里带火的。双十一那个实时成交额大屏,背后就是 Flink 在撑着。


三、Flink的核心特性

1. 真正的流处理(Event-Driven)

很多框架说自己是流处理,其实是"微批处理"(把数据攒成小批次再处理)。Flink 不一样,它是真正的逐条处理,数据来一条就处理一条,延迟更低。

┌─────────────────────────────────────────────────────────────┐
│  微批处理(如Spark Streaming)                               │
│                                                             │
│  数据流:● ● ● ● ● ● ● ● ●                                  │
│           ↓     ↓     ↓                                     │
│        [批次1] [批次2] [批次3]  ← 攒够一批再处理              │
│                                                             │
├─────────────────────────────────────────────────────────────┤
│  真正的流处理(Flink)                                       │
│                                                             │
│  数据流:● → ● → ● → ● → ●                                  │
│          ↓   ↓   ↓   ↓   ↓   ← 来一条处理一条               │
└─────────────────────────────────────────────────────────────┘

2. 精确一次语义(Exactly-Once)

在分布式系统中,保证数据"不丢不重"是很难的。Flink 通过 Checkpoint机制 实现了精确一次语义,意思是:

  • 即使程序挂了重启,也能保证每条数据恰好被处理一次
  • 不会丢数据,也不会重复处理

这对金融、电商等场景非常重要——你肯定不希望用户付了一次钱,系统却扣了两次吧?

3. 高吞吐低延迟

Flink 单节点每秒可以处理数百万条消息,延迟可以控制在毫秒级。这得益于:

  • 高效的内存管理(自己管内存,不完全依赖JVM GC)
  • 异步快照(Checkpoint不会阻塞数据处理)
  • 网络传输优化

4. 流批一体

Flink 把批处理看作是流处理的特例(有限流)。也就是说:

  • 你用 DataStream API 写的代码,既能处理实时流,也能处理历史数据
  • 不需要像以前那样,流处理用 Storm,批处理用 Spark,维护两套代码

5. 丰富的时间语义

Flink 支持三种时间概念:

时间类型含义场景
Event Time事件真正发生的时间大多数业务场景
Processing Time数据被处理的时间对时间精度要求不高
Ingestion Time数据进入Flink的时间折中方案

比如用户在 10:00:00 下单,但数据 10:00:05 才到 Flink。如果用 Event Time,订单会被归到 10:00:00 那个窗口;如果用 Processing Time,会被归到 10:00:05 那个窗口。


四、Flink能用来做什么

典型应用场景

1. 实时数据分析
  • 电商实时大屏:GMV、订单量、UV等指标实时展示
  • 实时报表:不用等到T+1,数据实时可见
2. 实时ETL
  • 数据清洗:去重、过滤、格式转换
  • 数据同步:MySQL → Kafka → Flink → Elasticsearch
3. 实时风控
  • 信用卡欺诈检测:发现异常交易立即告警
  • 薅羊毛识别:检测刷单、恶意注册等行为
4. 实时推荐
  • 短视频推荐:根据实时行为调整推荐结果
  • 广告投放:实时计算用户特征,精准投放
5. IoT实时监控
  • 工业设备监控:传感器数据实时分析,异常预警
  • 车联网:车辆轨迹实时追踪

哪些公司在用Flink

公司使用场景
阿里巴巴双十一大屏、搜索推荐、风控
字节跳动推荐系统、广告、数据中台
美团实时营销、配送调度
腾讯游戏数据分析、广告
Netflix实时数据管道
Uber实时定价、ETA预估

五、Flink vs Spark Streaming

很多同学会问:Spark Streaming 也能做流处理,为什么要用 Flink?

这是个好问题。我们来对比一下:

维度FlinkSpark Streaming
处理模型真正的流处理微批处理(Mini-Batch)
延迟毫秒级秒级(取决于批次间隔)
状态管理原生支持,功能强大需要借助外部存储
容错语义精确一次(原生支持)精确一次(需要配置)
窗口机制灵活,支持会话窗口相对简单
SQL支持Flink SQL(成熟)Spark SQL(更成熟)
生态流处理更强批处理生态更丰富
学习曲线稍陡稍平缓

一句话总结

  • 如果你的场景对延迟要求高(毫秒级)、需要复杂的状态计算、需要精确一次语义,选 Flink
  • 如果你的场景延迟要求不高(秒级够用)、已有 Spark 技术栈、以批处理为主,选 Spark

当然,现在很多公司是 Flink + Spark 混用的,各取所长。


六、Flink核心概念速览

在后续的文章中,我们会详细讲每个概念。这里先混个脸熟:

1. JobManager 和 TaskManager

┌─────────────────────────────────────────────────────────────┐
│                        Flink 集群                           │
│  ┌─────────────┐      ┌─────────────┐                      │
│  │ JobManager  │      │ JobManager  │  ← 主节点(负责调度)  │
│  │  (Active)   │      │  (Standby)  │    支持高可用         │
│  └─────────────┘      └─────────────┘                      │
│         │                                                   │
│         ↓                                                   │
│  ┌─────────────┐  ┌─────────────┐  ┌─────────────┐         │
│  │ TaskManager │  │ TaskManager │  │ TaskManager │         │
│  │   (Worker)  │  │   (Worker)  │  │   (Worker)  │         │
│  └─────────────┘  └─────────────┘  └─────────────┘         │
│        ↑                ↑                ↑                  │
│        └────────────────┴────────────────┘                  │
│                  从节点(负责执行任务)                        │
└─────────────────────────────────────────────────────────────┘
  • JobManager:老板,负责接收任务、分配任务、协调Checkpoint
  • TaskManager:员工,负责干活,执行具体的计算任务

2. DataStream 和 DataSet

  • DataStream:处理无界流(实时数据)
  • DataSet:处理有界数据集(批处理)
  • Flink 1.12+ 推荐统一使用 DataStream API,流批一体

3. Source、Transformation、Sink

这是 Flink 程序的三大组成部分:

数据源 → 数据处理 → 数据输出

Source → Transformation → Sink
 │            │            │
 │            │            └── 输出到 Kafka/MySQL/ES/文件...
 │            └── map/filter/keyBy/window 等算子
 └── 从 Kafka/文件/Socket/数据库 读取数据

4. Checkpoint 和 Savepoint

  • Checkpoint:自动的、轻量级的状态快照,用于故障恢复
  • Savepoint:手动触发的、完整的状态快照,用于版本升级、迁移

5. 窗口(Window)

流数据是无限的,但我们的计算往往要划分边界。窗口就是用来切分数据的:

时间轴:─────────────────────────────────────────────►

滚动窗口(Tumbling):
[  窗口1  ][  窗口2  ][  窗口3  ][  窗口4  ]
← 窗口之间不重叠

滑动窗口(Sliding):
[    窗口1    ]
      [    窗口2    ]
            [    窗口3    ]
← 窗口之间有重叠

会话窗口(Session):
[  会话1  ]          [  会话2  ]      [ 会话3 ]
← 根据活动间隙动态划分

七、总结

这篇文章我们从宏观角度认识了 Flink:

  1. 流处理 vs 批处理:流处理是"来一条处理一条",适合实时场景
  2. Flink是什么:一个高吞吐、低延迟、流批一体的分布式流处理框架
  3. 核心特性:真正的流处理、精确一次语义、丰富的时间语义、流批一体
  4. 应用场景:实时大屏、实时ETL、风控、推荐、IoT监控等
  5. vs Spark:Flink 在流处理领域更强,Spark 在批处理生态更丰富
  6. 核心概念:JobManager/TaskManager、DataStream、Source/Transformation/Sink、Checkpoint、窗口

下一篇文章,我们将动手搭建 Flink 开发环境,并写出第一个 WordCount 程序。别担心,我会一步步带你把环境搞定,保证你能跑起来!


热门专栏推荐

等等等还有许多优秀的合集在主页等着大家的光顾,感谢大家的支持


文章到这里就结束了,如果有什么疑问的地方请指出,诸佬们一起来评论区一起讨论😊
希望能和诸佬们一起努力,今后我们一起观看感谢您的阅读🙏
如果帮助到您不妨3连支持一下,创造不易您们的支持是我的动力🌟

更多推荐