Flink 学习笔记
一、Flink 基础概念
1.1 什么是 Flink
Apache Flink 是一个开源的分布式流处理框架,专门用于高效处理大规模数据流。它支持流处理和批处理两种模式,但核心是流处理。
主要特点:
- 高吞吐、低延迟的流处理引擎
- 支持事件时间(Event Time)处理,解决数据乱序问题
- 提供精确一次(Exactly-Once)的语义保证
- 支持有状态的计算
- 灵活的窗口操作
- 完善的容错机制
1.2 Flink 与其他框架的对比
- Flink vs Spark Streaming:Flink 是真正的流处理(微批处理),Spark Streaming 是微批处理框架
- Flink vs Kafka Streams:Flink 更强大灵活,支持更复杂的计算逻辑和更多状态管理能力
- Flink vs Storm:Flink 提供更高的吞吐量和更强大的 API
1.3 Flink 架构中的核心概念
Source:数据源,负责读取外部数据(Kafka、文件、Socket 等)
Operator:算子,执行具体的计算逻辑(map、filter、reduce 等)
Sink:数据汇,负责输出结果到外部系统(Kafka、数据库、文件等)
Stream:数据流,无限的数据序列
Transformation:转换操作,将一个或多个 DataStream 转换为新的 DataStream
二、Flink 部署架构
2.1 系统架构
Flink 采用典型的 Master-Slave 架构:
JobManager(主节点)
- 接收提交的作业
- 负责作业的调度和协调
- 管理 TaskManager 的生命周期
- 存储和维护作业状态
TaskManager(从节点)
- 执行任务的实际工作节点
- 管理自己的资源(内存、CPU、网络)
- 与 JobManager 定期心跳通信
Client
- 用户提交作业的客户端
- 将作业转换成 JobGraph 并提交给 JobManager
2.2 部署模式
Standalone 模式:独立部署,适合开发测试
Yarn 模式:运行在 Hadoop Yarn 之上,适合与其他大数据系统集成
Kubernetes 模式:容器化部署,适合云原生场景
Flink on K8s:现在推荐的部署方式
三、DataStream API
3.1 基本编程模型
// 1. 获取执行环境
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 2. 创建数据源(Source)
DataStream<String> stream = env.addSource(new FlinkKafkaConsumer<>(...));
// 3. 数据转换(Transformation)
stream.map(x -> x * 2)
.filter(x -> x > 10)
.print();
// 4. 数据输出(Sink)
stream.addSink(new FlinkKafkaProducer<>(...));
// 5. 提交执行
env.execute("Job Name");
3.2 常用转换算子
map:一对一转换,逐个元素映射
filter:过滤元素,保留满足条件的元素
flatMap:一对多转换,每个元素可以产生 0 个或多个结果
keyBy:按键分组,返回 KeyedStream,后续操作会分组执行
reduce:聚合操作,同一分组内的元素进行规约
aggregate:更灵活的聚合方式,可以使用自定义聚合逻辑
union:合并多个 DataStream
connect:连接两个 DataStream,可以进行联合处理
split / select:分流操作(旧版本,新版本用 side output)
3.3 Source 常见实现
Kafka Source
Properties props = new Properties();
props.setProperty("bootstrap.servers", "localhost:9092");
props.setProperty("group.id", "flink-group");
DataStream<String> stream = env.addSource(
new FlinkKafkaConsumer<>("topic", new SimpleStringSchema(), props)
);
Socket Source(用于测试)
DataStream<String> stream = env.socketTextStream("localhost", 9999);
集合 Source
DataStream<Integer> stream = env.fromCollection(Arrays.asList(1, 2, 3, 4, 5));
自定义 Source:实现 SourceFunction 接口
3.4 Sink 常见实现
Print Sink:直接打印到控制台
stream.print();
Kafka Sink
stream.addSink(new FlinkKafkaProducer<>(
"topic",
new SimpleStringSchema(),
props
));
自定义 Sink:实现 SinkFunction 接口
四、时间概念(Time Semantics)
4.1 三种时间
Processing Time(处理时间):事件被处理的时间,即系统时钟时间。简单但无法处理乱序数据。
Event Time(事件时间):事件发生的真实时间,即数据本身携带的时间戳。能准确处理乱序数据,是生产环境推荐方式。
Ingestion Time(摄入时间):事件进入 Flink 系统的时间。介于两者之间,使用较少。
4.2 配置时间语义
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 配置使用事件时间
env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);
// 配置最大允许延迟
env.getConfig().setAutoWatermarkInterval(1000);
4.3 Watermark(水位线)
Watermark 用于衡量事件时间的进度,标记"在这个时间点之前的数据都已经到达"。
作用:
- 触发窗口计算
- 处理迟到数据
- 确保计算的完整性
生成 Watermark
DataStream<Event> stream = env.addSource(...)
.assignTimestampsAndWatermarks(
WatermarkStrategy
.forBoundedOutOfOrderness(Duration.ofSeconds(20))
.withTimestampAssigner((event, recordTimestamp) -> event.getTimestamp())
);
allowedLateness:指定允许迟到数据的最大时间
stream.keyBy(...)
.window(TumblingEventTimeWindows.of(Time.seconds(10)))
.allowedLateness(Time.minutes(1))
.process(...)
五、窗口(Window)
5.1 窗口的概念
窗口是处理无限数据流的关键机制,将无限数据流分割成有限的数据集合。
5.2 窗口类型
时间窗口
-
滚动窗口(Tumbling Window):固定大小,无重叠
.window(TumblingEventTimeWindows.of(Time.seconds(10))) -
滑动窗口(Sliding Window):固定大小,有重叠
.window(SlidingEventTimeWindows.of(Time.seconds(10), Time.seconds(5))) -
会话窗口(Session Window):根据数据活跃度动态生成
.window(EventTimeSessionWindows.withGap(Time.seconds(5)))
计数窗口
-
滚动计数窗口
.countWindow(100) -
滑动计数窗口
.countWindow(100, 50)
5.3 窗口计算
stream.keyBy(x -> x.userId)
.window(TumblingEventTimeWindows.of(Time.seconds(10)))
.reduce((x, y) -> new Event(x.userId, x.count + y.count))
.print();
或使用 process() 获得更多控制权:
stream.keyBy(...)
.window(...)
.process(new ProcessWindowFunction<Event, Result, String, TimeWindow>() {
@Override
public void process(String key, Context context, Iterable<Event> elements, Collector<Result> out) {
// 自定义窗口计算逻辑
}
})
六、状态管理(State)
6.1 状态概念
状态是指在流处理过程中维护的信息,用于存储计算的中间结果。
6.2 状态类型
Operator State(算子状态)
- 绑定到单个算子实例
- 分布式场景下不推荐使用
- 应用场景有限
Keyed State(键值状态)
- 按键分组存储
- 每个键对应一份状态
- 生产环境推荐使用
6.3 Keyed State 类型
ValueState:存储单个值
ValueState<Integer> countState = getRuntimeContext()
.getState(new ValueStateDescriptor<>("count", Integer.class));
ListState:存储列表
ListState<String> listState = getRuntimeContext()
.getListState(new ListStateDescriptor<>("list", String.class));
MapState:存储键值对
MapState<String, Integer> mapState = getRuntimeContext()
.getMapState(new MapStateDescriptor<>("map", String.class, Integer.class));
ReducingState 和 AggregatingState:专门用于聚合
6.4 有状态计算示例
public class StatefulCounter extends KeyedProcessFunction<String, Event, Result> {
private transient ValueState<Integer> countState;
@Override
public void open(Configuration parameters) throws Exception {
ValueStateDescriptor<Integer> descriptor =
new ValueStateDescriptor<>("count", Integer.class);
countState = getRuntimeContext().getState(descriptor);
}
@Override
public void processElement(Event event, Context ctx, Collector<Result> out)
throws Exception {
Integer current = countState.value();
current = current == null ? 1 : current + 1;
countState.update(current);
out.collect(new Result(event.getKey(), current));
}
}
七、容错机制
7.1 Checkpoint
Checkpoint 是 Flink 故障恢复的核心机制,定期保存整个应用的状态快照。
启用 Checkpoint
env.enableCheckpointing(60000); // 每 60 秒进行一次 checkpoint
env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
env.getCheckpointConfig().setMinPauseBetweenCheckpoints(30000);
env.getCheckpointConfig().setCheckpointTimeout(600000);
env.getCheckpointConfig().setMaxConcurrentCheckpoints(1);
Checkpoint 模式
- EXACTLY_ONCE:精确一次语义,性能略低但最安全
- AT_LEAST_ONCE:至少一次语义,性能更好但可能重复处理
7.2 State Backend
状态后端用于存储状态数据。
MemoryStateBackend:存储在内存中,适合开发测试
env.setStateBackend(new MemoryStateBackend());
FsStateBackend:存储在文件系统中,支持大状态
env.setStateBackend(new FsStateBackend("hdfs:///flink/checkpoints"));
RocksDBStateBackend:基于 RocksDB,支持最大的状态量
env.setStateBackend(new RocksDBStateBackend("hdfs:///flink/checkpoints"));
7.3 Savepoint
Savepoint 是人工触发的状态快照,用于应用版本升级、迁移或维护。
# 触发 savepoint
flink savepoint <job-id> [target-directory]
# 从 savepoint 恢复
flink run -s <savepoint-path> <jar-path>
八、CEP(复杂事件处理)
8.1 什么是 CEP
CEP 用于在事件流中检测和提取复杂的事件模式。
8.2 基本用法
DataStream<Event> input = env.addSource(...);
Pattern<Event, ?> pattern = Pattern
.<Event>begin("start").where(e -> e.getId() == 1)
.next("middle").where(e -> e.getId() == 2)
.followedBy("end").where(e -> e.getId() == 3)
.within(Time.seconds(10));
PatternStream<Event> patternStream = CEP.pattern(input, pattern);
DataStream<String> result = patternStream.process(
new PatternProcessFunction<Event, String>() {
@Override
public void processMatch(Map<String, List<Event>> pattern,
Context ctx,
Collector<String> out) {
out.collect("Pattern matched!");
}
}
);
九、高级特性
9.1 Side Output(旁路输出)
用于将数据分为多个输出流,例如将异常数据和正常数据分开处理。
OutputTag<String> outputTag = new OutputTag<String>("side-output"){};
SingleOutputStreamOperator<Integer> mainStream = stream
.process(new ProcessFunction<String, Integer>() {
@Override
public void processElement(String value, Context ctx, Collector<Integer> out) {
if (value.startsWith("ERROR")) {
ctx.output(outputTag, value);
} else {
out.collect(Integer.parseInt(value));
}
}
});
DataStream<String> sideStream = mainStream.getSideOutput(outputTag);
9.2 Broadcast State
广播状态用于将某个流的数据广播到所有并行任务。
MapStateDescriptor<String, String> descriptor =
new MapStateDescriptor<>("broadcast", String.class, String.class);
BroadcastStream<String> broadcastStream = configStream.broadcast(descriptor);
stream.connect(broadcastStream)
.process(new BroadcastProcessFunction<Event, String, Result>() {
@Override
public void processElement(Event event, ReadOnlyContext ctx, Collector<Result> out) {
ReadOnlyBroadcastState<String, String> state = ctx.getBroadcastState(descriptor);
// 使用广播状态
}
@Override
public void processBroadcastElement(String element, Context ctx, Collector<Result> out) {
BroadcastState<String, String> state = ctx.getBroadcastState(descriptor);
state.put("key", element);
}
});
9.3 异步函数(Async I/O)
用于异步调用外部服务,提高吞吐量。
AsyncDataStream.unorderedWait(
stream,
new AsyncFunction<Event, Result>() {
@Override
public void asyncInvoke(Event input, ResultFuture<Result> resultFuture) {
// 异步调用外部服务
CompletableFuture.supplyAsync(() -> {
return callExternalService(input);
}).thenAccept(resultFuture::complete);
}
},
1000,
TimeUnit.MILLISECONDS
);
十、调优和最佳实践
10.1 性能调优
并行度设置:与数据源分区数相匹配
env.setParallelism(8);
减少序列化开销:使用 Kryo 序列化器
env.getConfig().enableForceKryo();
状态大小控制:定期清理过期状态,使用 TTL
StateTtlConfig ttlConfig = StateTtlConfig
.newBuilder(Time.days(1))
.setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
.build();
descriptor.enableTimeToLive(ttlConfig);
10.2 最佳实践
- 使用事件时间而非处理时间,保证计算的准确性
- 正确配置 Watermark,平衡延迟和数据完整性
- 启用 Checkpoint,确保故障恢复能力
- 监控作业状态,及时发现和解决问题
- 合理使用状态,避免状态过大导致性能下降
- 充分测试,特别是复杂的业务逻辑
十一、常见问题排查
数据处理延迟高
- 检查并行度设置是否合理
- 检查状态大小是否过大
- 优化 Checkpoint 配置
数据丢失
- 确保启用了 Checkpoint
- 检查 Source 和 Sink 的容错能力
- 验证数据源的可靠性
内存溢出
- 减少并行度或增加 TaskManager 内存
- 检查是否存在状态泄漏
- 优化数据结构
更多推荐
所有评论(0)