Flink DataStream API
概述
DataStream API 是 Flink 流处理的核心编程接口:数据从 Source 流入,经过 Transformation 算子处理,最终写入 Sink。本文讲清三大环节、KeyedStream/WindowedStream 分层模型,以及算子链(Operator Chain)与执行优化。
一、API 全景
Source ──► Transformation ──► Sink
│ │ │
读数据 处理逻辑 输出结果| 环节 | 作用 | 示例 |
|---|---|---|
| Source | 数据入口 | Kafka、Socket、集合 |
| Transformation | 计算转换 | map、filter、keyBy、window |
| Sink | 数据出口 | Kafka、HDFS、MySQL、打印 |
java
DataStream<String> stream = env.addSource(kafkaSource);
stream
.flatMap(new SplitWords())
.map(word -> new Tuple2<>(word, 1))
.keyBy(t -> t.f0)
.sum(1)
.addSink(consoleSink);
env.execute("word-count");二、Source
2.1 常用 Source
| Source | 说明 |
|---|---|
env.addSource(...) | 自定义/连接器源 |
| Kafka Source | FlinkKafkaConsumer |
| Socket | 测试用 |
| 集合/文件 | 批式/测试 |
2.2 并行与限速
| 参数 | 说明 |
|---|---|
| Source 并行度 | 决定读取并发 |
| 背压 | 下游慢自动限速 |
| 时间戳/水位线 | 在 Source 后 assignTimestampsAndWatermarks |
三、Transformation
3.1 基础算子
| 算子 | 作用 |
|---|---|
| map | 1 进 1 出 |
| flatMap | 1 进 0~N 出 |
| filter | 条件过滤 |
| connect/union | 合并流 |
3.2 键控算子
| 算子 | 说明 |
|---|---|
| keyBy | 按键分区(类似 groupBy) |
| reduce | 按键归约 |
| aggregate | 聚合 |
| process | 底层处理(可访问状态/时间) |
3.3 KeyedStream 模型
keyBy 之后:
DataStream → KeyedStream → WindowedStream → DataStream
│
数据按键分区,状态按 key 隔离| 层级 | 特征 |
|---|---|
| DataStream | 普通流 |
| KeyedStream | 按键分区,状态按 key |
| WindowedStream | 窗口化流,窗口计算 |
| ConnectedStreams | 双流连接 |
四、Sink
4.1 常用 Sink
| Sink | 说明 |
|---|---|
| Kafka Sink | 写 Kafka |
| HDFS Sink | 落文件系统 |
| JDBC Sink | 写数据库 |
| 自定义 SinkFunction | 任意输出 |
4.2 可靠性
| 语义 | 实现 |
|---|---|
| At-Most-Once | 直接写 |
| At-Least-Once | 写后提交 |
| Exactly-Once | 两阶段提交(事务 Sink) |
java
// Kafka 精确一次 Sink
kafkaSink.withTransaction() // 语义配置五、UDF 算子链
5.1 什么是算子链
相邻算子满足条件时合并为同一 Task(同一线程):
map → filter → map (合并成一个 Task,避免线程切换)
↓
window(需要 Shuffle) (断开)| 合并条件 | 说明 |
|---|---|
| 并行度相同 | 上下游一致 |
| 无 Shuffle | 无 keyBy 等重分区 |
| 无状态冲突 | 单输入 |
5.2 链的收益
| 收益 | 说明 |
|---|---|
| 减少线程切换 | 合并算子共享线程 |
| 减少网络传输 | 本地直接传递 |
| 降低开销 | 序列化/缓冲减少 |
5.3 控制链
| 方法 | 作用 |
|---|---|
startNewChain() | 从此处断开 |
disableChaining() | 禁止链 |
env.disableOperatorChaining() | 全局禁用 |
六、执行环境与模式
6.1 环境
java
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(4);
env.enableCheckpointing(60000);6.2 批流模式
| 模式 | 说明 |
|---|---|
| STREAMING | 流处理(默认) |
| BATCH | 批处理(DataStream 也支持) |
6.3 延迟与吞吐权衡
| 参数 | 说明 |
|---|---|
execution.buffer-timeout | 缓冲超时(默认 100ms) |
| 缓冲越小 | 延迟低、吞吐低 |
| 缓冲越大 | 吞吐高、延迟高 |
七、实践要点
| 要点 | 说明 |
|---|---|
| 键控优先 | 能用 keyBy 状态不用全局状态 |
| 避免大对象 | 减少序列化开销 |
| 合理并行 | 与 Source/Sink 吞吐匹配 |
| 水位线 | 乱序场景必须设置 |
| Checkpoint | 生产必开 |
常见问题速查
| 问题 | 原因与处理 |
|---|---|
| 结果乱序 | 未设置事件时间与水位线 |
| 并行度不一致 | 各算子并行度统一或调整链 |
| 背压高 | 优化算子或扩容 |
| Key 数据倾斜 | 检查 key 分布,预聚合或加盐 |
| Sink 慢 | 批量写入、控制并行 |