什么是 Apache Flink?
Apache Flink 是一个开源的分布式流处理框架,以低延迟、高吞吐、Exactly-Once 语义著称,广泛应用于实时数据分析、事件驱动应用和 ETL 管道等场景。
与 Spark Streaming 的微批处理不同,Flink 采用真正的流处理模型,能够达到毫秒级的处理延迟。
核心概念
1. 数据流(DataStream)
Flink 中一切都是数据流。有界流对应批处理场景,无界流对应实时流处理场景。
2. 时间语义
- Event Time:数据本身携带的时间戳,最准确。
- Processing Time:数据到达算子时的系统时间,性能最好。
- Ingestion Time:数据进入 Flink 系统的时间。
3. 窗口(Window)
- 滚动窗口:固定大小、不重叠。
- 滑动窗口:固定大小、有重叠。
- 会话窗口:按活动间隔动态划分。
4. 状态与容错
Flink 支持有状态计算,结合 Checkpoint 机制实现 Exactly-Once 语义。
架构概览
- JobManager:负责作业调度、Checkpoint 协调和故障恢复。
- TaskManager:负责实际数据处理,每个 TaskManager 包含若干 Task Slot。
快速上手:DataStream API
以下是一个最简单的 WordCount 示例(Java):
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
DataStream<String> text = env.fromElements("hello world", "hello flink");
text.flatMap(...).keyBy(t -> t.f0).sum(1).print();
env.execute("WordCount");Table API 与 SQL
Flink 提供 Table API 和 Flink SQL,让你用 SQL 语法处理流数据,大大降低开发门槛。
常见使用场景
- 实时数仓:实现分钟级数据入仓。
- 实时监控告警:对指标实时计算并触发告警。
- 用户行为分析:实时计算 PV/UV、漏斗转化等。
- 流式 ETL:实时清洗、转换、加载数据。
下一步学习建议
- 本地搭建 Flink 单机环境,运行官方示例。
- 学习 DataStream API 核心算子:map、filter、keyBy、window、process 等。
- 学习 Flink 与 Kafka 集成,构建实时数据管道。
- 深入学习 Checkpoint 机制与 RocksDB State Backend。
- 尝试用 Flink SQL 重写业务逻辑。
Apache Flink 生态正在快速成熟,现在就开始你的 Flink 之旅吧!