Flin知识点(六)|Flink状态管理
一、为什么状态是Flink的基石?
流计算本质是有状态的计算。
状态指算子在处理数据流过程中需要“记住”的中间结果,例如窗口聚合值、用户会话信息或Kafka Offset。
当状态规模从MB级增长到GB甚至TB级,且作业面临节点故障时,如何高效可靠地管理这些“记忆”成为核心挑战。
本文目标是通过状态类型全景、状态后端机制以及容错体系,帮助读者掌握让Flink作业在生产环境中稳定运行的关键能力。
二、Flink状态类型全景图
Flink提供两种主要状态形式,分别适用于不同场景。
-
Operator State(算子状态)
与算子并行实例绑定,而非具体Key。
典型应用:Kafka Source的Offset管理、自定义算子的缓冲池。
支持API:ListState、BroadcastState。
扩展时需手动处理状态重分区,复杂度较高。 -
Keyed State(键控状态)
通过keyBy()后的Key进行分区,是分布式状态的主流形式。
支持类型:ValueState、ListState、MapState、AggregatingState。
访问方式:在RichFunction中通过RuntimeContext获取。
2.1 Operator State(算子状态)
状体和算子并行实例绑定,一个算子的状态不能被其他算子所访问。比如一个算子的并行度为2,那个该算子实际有2个算子状态:

算子状态有三种结构类型:ListState、UnionListState、BroadcastState
2.1.1 ListState(列表状态)
特点:状态以列表形式存储,与算子并行实例绑定。
重缩容行为:按轮询(round-robin)或均匀分发,每个实例只拿到部分元素。
场景:大部分用于source、sink。比如Kafka Source 的 Offset 管理(部分实现)、自定义缓冲列表。
示例:自定义缓冲池、部分 Offset 管理。
import org.apache.flink.api.common.state.ListState;
import org.apache.flink.api.common.state.ListStateDescriptor;
import org.apache.flink.api.common.functions.RichFlatMapFunction;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.util.Collector;
import org.apache.flink.api.java.tuple.Tuple2;
public class ListStateDemo extends RichFlatMapFunction<Tuple2<String, Integer>, String> {
// Operator State:ListState
private transient ListState<Integer> bufferState;
@Override
public void open(Configuration parameters) {
ListStateDescriptor<Integer> descriptor =
new ListStateDescriptor<>("buffer-list", Integer.class);
// 获取 Operator State
bufferState = getRuntimeContext().getListState(descriptor);
}
@Override
public void flatMap(Tuple2<String, Integer> input, Collector<String> out) throws Exception {
// 添加到状态列表
bufferState.add(input.f1);
// 示例业务:缓冲满5个时输出并清空
if (bufferState.get().spliterator().getExactSizeIfKnown() >= 5) {
StringBuilder sb = new StringBuilder("Buffer full: ");
for (Integer val : bufferState.get()) {
sb.append(val).append(" ");
}
out.collect(sb.toString());
bufferState.clear(); // 清空状态
}
}
}
// 引用
env.addSource(...).flatMap(new ListStateDemo()).setParallelism(2);
2.1.2 UnionListState(联合列表状态)
特点:同样是列表状态,但强调“联合”语义。
重缩容行为:所有并行实例的状态先合并成完整列表,再全量广播给每个新实例(broadcast-style recovery)。
场景:需要全局视图的元数据(如某些 connector 的全量偏移、恢复时的完整配置)。
示例:需要全局视图的元数据(如全量配置列表、某些 Connector 的偏移)。
import org.apache.flink.api.common.state.ListState;
import org.apache.flink.api.common.state.ListStateDescriptor;
import org.apache.flink.api.common.functions.RichFlatMapFunction;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.util.Collector;
import org.apache.flink.api.java.tuple.Tuple2;
public class UnionListStateDemo extends RichFlatMapFunction<Tuple2<String, String>, String> {
// Operator State:UnionListState
private transient ListState<String> unionState;
@Override
public void open(Configuration parameters) {
ListStateDescriptor<String> descriptor =
new ListStateDescriptor<>("union-list", String.class);
// 关键:使用 getUnionListState
unionState = getRuntimeContext().getUnionListState(descriptor);
}
@Override
public void flatMap(Tuple2<String, String> input, Collector<String> out) throws Exception {
// 添加元素
unionState.add(input.f1 + ":" + input.f2);
// 示例业务:每收到10个元素就输出完整联合列表(重缩容后仍保持全局一致)
long size = unionState.get().spliterator().getExactSizeIfKnown();
if (size >= 10) {
StringBuilder sb = new StringBuilder("Union full list: ");
for (String item : unionState.get()) {
sb.append(item).append(" | ");
}
out.collect(sb.toString());
// 注意:UnionListState 不建议频繁 clear(会影响全局一致性)
}
}
}
// 引用
env.addSource(...).flatMap(new UnionListStateDemo()).setParallelism(2);
2.1.3 BroadcastState(广播状态)
特点:特殊的 Map 形式状态,用于广播流与非广播流的联合处理。
重缩容行为:全量复制到每个并行实例,支持动态更新。
典型场景:动态规则/配置广播(如实时更新过滤条件、模式匹配)。
具体的详细的讲解,可以参考专栏之前的Flink知识点(三)|Flink中的广播变量(Broadcast State)
2.2 Keyed State(键控状态)
通过keyBy()后的Key进行分区,是分布式状态的主流形式。Flink会根据key键来维护每一个状态。

键控状态有三种结构类型:ValueState、ListState、MapState、ReducingState、AggregatingState
2.2.1 ValueState(值状态)
特点:存储单一可更新值,支持 update()、value()、clear()。
场景:记录最新用户属性、会话计数器。
API:getRuntimeContext().getState(ValueStateDescriptor)。
示例:累计计数
import org.apache.flink.api.common.functions.RichFlatMapFunction;
import org.apache.flink.api.common.state.ValueState;
import org.apache.flink.api.common.state.ValueStateDescriptor;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.util.Collector;
import org.apache.flink.api.java.tuple.Tuple2;
public class ValueStateDemo extends RichFlatMapFunction<Tuple2<String, Long>, String> {
private transient ValueState<Long> countState;
@Override
public void open(Configuration parameters) throws Exception {
ValueStateDescriptor<Long> desc = new ValueStateDescriptor<>("key-count", Long.class);
countState = getRuntimeContext().getState(desc);
}
@Override
public void flatMap(Tuple2<String, Long> input, Collector<String> out) throws Exception {
Long current = countState.value();
if (current == null) current = 0L;
current += input.f1;
countState.update(current);
out.collect("Key: " + input.f0 + " → count = " + current);
}
}
2.2.2 ListState(列表状态)
特点:有序列表,支持 add()、get()、update()、clear()。
典型场景:收集历史事件、缓冲队列。
API:getRuntimeContext().getListState(ListStateDescriptor)(注意:与 Operator State 的 ListState 重分布行为不同)。
示例:每个key键的最近3条记录
import org.apache.flink.api.common.state.ListState;
import org.apache.flink.api.common.state.ListStateDescriptor;
import org.apache.flink.api.common.functions.RichFlatMapFunction;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.util.Collector;
import org.apache.flink.api.java.tuple.Tuple2;
public class ListStateDemo extends RichFlatMapFunction<Tuple2<String, String>, String> {
private transient ListState<String> events;
@Override
public void open(Configuration parameters) throws Exception {
ListStateDescriptor<String> desc = new ListStateDescriptor<>("events", String.class);
events = getRuntimeContext().getListState(desc);
}
@Override
public void flatMap(Tuple2<String, String> input, Collector<String> out) throws Exception {
events.add(input.f1);
// 示例:每累积 3 条就输出一次并清空(模拟小窗口)
long size = 0;
for (String ignored : events.get()) size++;
if (size >= 3) {
StringBuilder sb = new StringBuilder("Key " + input.f0 + " events: ");
for (String e : events.get()) {
sb.append(e).append(", ");
}
out.collect(sb.toString());
events.clear();
}
}
}
2.2.3 MapState(Map状态)
特点:键值对映射,支持 put()、get()、contains()、entries() 等 Map 操作。
典型场景:用户画像标签、动态配置映射。
API:getRuntimeContext().getMapState(MapStateDescriptor)。
示例:用户画像
import org.apache.flink.api.common.state.MapState;
import org.apache.flink.api.common.state.MapStateDescriptor;
import org.apache.flink.api.common.functions.RichFlatMapFunction;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.util.Collector;
import org.apache.flink.api.java.tuple.Tuple3;
public class MapStateDemo extends RichFlatMapFunction<Tuple3<String, String, Integer>, String> {
private transient MapState<String, Integer> tagScores;
@Override
public void open(Configuration parameters) throws Exception {
MapStateDescriptor<String, Integer> desc =
new MapStateDescriptor<>("tag-scores", String.class, Integer.class);
tagScores = getRuntimeContext().getMapState(desc);
}
@Override
public void flatMap(Tuple3<String, String, Integer> input, Collector<String> out) throws Exception {
String userId = input.f0;
String tag = input.f1;
Integer score = input.f2;
tagScores.put(tag, score);
// 示例:输出当前用户所有标签得分
StringBuilder sb = new StringBuilder("User " + userId + " tags: ");
tagScores.entries().forEach(e -> sb.append(e.getKey()).append("=").append(e.getValue()).append(" "));
out.collect(sb.toString());
}
}
2.2.4 ReducingState(归约状态)
特点:自动应用 ReduceFunction,每次 add() 后底层自动归约成单一值,支持 get()、clear()。
典型场景:实时求和、求最大值(无需手动维护中间结果)。
API:getRuntimeContext().getReducingState(ReducingStateDescriptor)。
示例:每个 key累加
import org.apache.flink.api.common.functions.ReduceFunction;
import org.apache.flink.api.common.state.ReducingState;
import org.apache.flink.api.common.state.ReducingStateDescriptor;
import org.apache.flink.api.common.functions.RichFlatMapFunction;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.util.Collector;
import org.apache.flink.api.java.tuple.Tuple2;
public class ReducingStateDemo extends RichFlatMapFunction<Tuple2<String, Long>, String> {
private transient ReducingState<Long> sumState;
@Override
public void open(Configuration parameters) throws Exception {
ReducingStateDescriptor<Long> desc =
new ReducingStateDescriptor<>("sum", new SumReducer(), Long.class);
sumState = getRuntimeContext().getReducingState(desc);
}
@Override
public void flatMap(Tuple2<String, Long> input, Collector<String> out) throws Exception {
sumState.add(input.f1);
out.collect("Key: " + input.f0 + " → total = " + sumState.get());
}
private static class SumReducer implements ReduceFunction<Long> {
@Override
public Long reduce(Long a, Long b) {
return a + b;
}
}
}
2.2.5 AggregatingState(聚合状态)
特点:使用 Accumulator 实现自定义聚合,支持 add()、get()、clear(),比 ReducingState 更灵活。
典型场景:复杂统计(如平均值、TopN)、带中间状态的聚合。
API:getRuntimeContext().getAggregatingState(AggregatingStateDescriptor)。
示例:求平均值
import org.apache.flink.api.common.functions.AggregateFunction;
import org.apache.flink.api.common.state.AggregatingState;
import org.apache.flink.api.common.state.AggregatingStateDescriptor;
import org.apache.flink.api.common.functions.RichFlatMapFunction;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.util.Collector;
import org.apache.flink.api.java.tuple.Tuple2;
public class AggregatingStateDemo extends RichFlatMapFunction<Tuple2<String, Long>, String> {
private transient AggregatingState<Long, Double> avgState;
@Override
public void open(Configuration parameters) throws Exception {
AggregatingStateDescriptor<Long, AvgAcc, Double> desc =
new AggregatingStateDescriptor<>("avg", new AvgAggregator(), Double.class);
avgState = getRuntimeContext().getAggregatingState(desc);
}
@Override
public void flatMap(Tuple2<String, Long> input, Collector<String> out) throws Exception {
avgState.add(input.f1);
out.collect("Key: " + input.f0 + " → avg = " + avgState.get());
}
// 累加器
public static class AvgAcc {
public long sum = 0;
public long count = 0;
}
// 聚合逻辑
public static class AvgAggregator implements AggregateFunction<Long, AvgAcc, Double> {
@Override
public AvgAcc createAccumulator() {
return new AvgAcc();
}
@Override
public AvgAcc add(Long value, AvgAcc acc) {
acc.sum += value;
acc.count++;
return acc;
}
@Override
public Double getResult(AvgAcc acc) {
return acc.count == 0 ? 0.0 : (double) acc.sum / acc.count;
}
@Override
public AvgAcc merge(AvgAcc a, AvgAcc b) {
a.sum += b.sum;
a.count += b.count;
return a;
}
}
}
关键注意事项
- 以上类型仅能在 Keyed Stream(
keyBy()之后)的 RichFunction / ProcessFunction 中使用。 - ListState 是唯一同时出现在 Operator State 和 Keyed State 中的类型,但两者在 Checkpoint 重分布时的行为完全不同。
- 生产中推荐结合 State TTL(
desc.setStateTtl())避免 Key 膨胀。
2.2.6 TTL(状态生存时间)
State TTL(Time-To-Live) 是 Flink 防止 Keyed State 无限膨胀的核心机制,尤其在用户画像、会话分析等高基数场景下必不可少。它让每个 Key 的状态可以自动过期,避免内存/磁盘持续增长,同时保留 Exactly-Once 语义。
TTL 配置参数(StateTtlConfig):
使用 StateTtlConfig.newBuilder(...) 构建,关键字段如下:
- 过期时间:
Time.minutes(30)/Time.hours(24)等(必填) - 更新触发类型(UpdateType):
OnCreateAndWrite(默认):创建或写入时刷新 TTLOnReadAndWrite:读写都刷新(适合“最后访问时间”场景)
- 过期可见性(StateVisibility):
ReturnExpiredIfNotCleanedUp(默认):访问时仍返回过期值,但标记为已过期NeverReturnExpired:直接返回 null,不返回过期值(更严格)
- 清理策略:可开启
cleanupFullSnapshot()或cleanupIncrementally()
示例:以ValueState为例
import org.apache.flink.api.common.state.ValueState;
import org.apache.flink.api.common.state.ValueStateDescriptor;
import org.apache.flink.api.common.state.StateTtlConfig;
import org.apache.flink.api.common.time.Time;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.functions.KeyedProcessFunction;
import org.apache.flink.util.Collector;
import org.apache.flink.api.java.tuple.Tuple2;
public class UserProfileFunction extends KeyedProcessFunction<String, Tuple2<String, Profile>, String> {
private transient ValueState<Profile> profileState;
@Override
public void open(Configuration parameters) {
// 构建 TTL 配置:30 分钟不更新即过期,读写都刷新
StateTtlConfig ttlConfig = StateTtlConfig
.newBuilder(Time.minutes(30))
.setUpdateType(StateTtlConfig.UpdateType.OnReadAndWrite) // 最后访问时间刷新
.setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired) // 不返回过期值
.cleanupFullSnapshot() // 全量快照时彻底清理
.build();
ValueStateDescriptor<Profile> descriptor =
new ValueStateDescriptor<>("user-profile", Profile.class);
descriptor.setStateTtl(ttlConfig); // ← 关键一行
profileState = getRuntimeContext().getState(descriptor);
}
@Override
public void processElement(Tuple2<String, Profile> input, Context ctx, Collector<String> out) throws Exception {
Profile current = profileState.value();
if (current == null) {
current = new Profile();
}
// 更新业务字段...
current.lastActive = System.currentTimeMillis();
profileState.update(current);
out.collect("User " + input.f0 + " updated, TTL active");
}
}
2.3 来个总结
| 维度 | Operator State | Keyed State |
|---|---|---|
| 作用范围 | 算子并行子任务 | 当前Key |
| 横向扩展 | 状态重分区复杂 | 随Key自动分区 |
| 典型场景 | 源/汇连接器、广播变量 | 窗口聚合、用户画像、会话分析 |
三、状态后端(State Backend)
3.1 状态后端怎么选?
状态后端是Flink状态管理的核心引擎,负责在机访问与Checkpoint持久化两大职责。
3.1.1 HashMapStateBackend(内存型)
- 原理:所有状态对象直接存储在JVM堆内存中。
- 优点:访问速度极快,无序列化开销。
- 缺点:完全受限于TaskManager内存容量,易触发OOM。
- 适用场景:状态规模小(MB级)、本地开发、测试环境。
3.1.2 EmbeddedRocksDBStateBackend(硬盘+内存混合型)
- 原理:状态以
(Key, Namespace) → byte[]形式存储在本地RocksDB(高性能KV数据库)中,热数据通过LRU缓存驻留内存。 - 优点:
- 容量远超内存(受限于本地磁盘,可达TB级)。
- 支持增量Checkpoint,持久化文件体积小。
- 缺点:读写涉及序列化/反序列化与磁盘IO,吞吐量相对较低。
- 适用场景:生产环境标配,尤其适合GB~TB级大状态、高可用需求。
关键抉择建议:
小状态追求极致性能时选用HashMapStateBackend;大状态或生产环境统一采用EmbeddedRocksDBStateBackend。
口诀:小状态用HashMap,大状态用RocksDB。
3.2 状态后端怎么配置?
客观来看,状态后端(State Backend)是 Flink 作业启动时必须显式配置的核心参数,直接决定状态的存储位置、容量上限和 Checkpoint 性能。配置方式主要有代码设置(灵活,适合开发/测试)和配置文件设置(集群统一管理,适合生产)。生产环境中,有些任务可能需要单独配置,可以使用代码设置。
- 推荐配置决策(快速选择)
- 小状态(< 几百 MB):优先 HashMapStateBackend,追求极致读写速度。
- 大状态(GB~TB 级、生产环境):必须使用 EmbeddedRocksDBStateBackend,支持增量 Checkpoint 和本地磁盘扩展。
- 若不配置,Flink 会自动使用 HashMapStateBackend(不推荐生产)。
3.2.1 配置方式一:代码中设置(StreamExecutionEnvironment)
这是最常用、最直观的做法,可在作业 main 方法中直接写死或通过参数动态切换。
3.2.1.1 HashMapStateBackend(内存型)
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 基础配置(内存型)
env.setStateBackend(new HashMapStateBackend());
// 可选:关闭异步快照(仅测试用)
env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
3.2.1.2 EmbeddedRocksDBStateBackend(生产标配)
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 推荐生产配置(开启增量 Checkpoint)
EmbeddedRocksDBStateBackend rocksDB = new EmbeddedRocksDBStateBackend(true);
env.setStateBackend(rocksDB);
// 高级调优示例(强烈建议生产加上)
Configuration rocksConfig = new Configuration();
rocksConfig.setString("state.backend.rocksdb.block-cache-size", "256MB");
rocksConfig.setString("state.backend.rocksdb.write-buffer-size", "64MB");
rocksConfig.setInteger("state.backend.rocksdb.write-buffer-number", 4);
// 应用配置
env.configure(rocksConfig, getClass().getClassLoader());
关键参数说明(RocksDB 专属):
state.backend.rocksdb.memory.managed:是否让 Flink 托管内存(推荐 true)。state.backend.rocksdb.block-cache-size:块缓存大小(默认 8MB,生产建议 128MB~512MB)。state.backend.rocksdb.write-buffer-size:写缓冲大小(默认 64MB)。state.backend.rocksdb.incremental:增量 Checkpoint(默认 true,已在构造器中开启)。
3.2.2 配置方式二:flink-conf.yaml(集群级统一配置)
生产环境推荐在 flink/conf/flink-conf.yaml 中全局设置,避免每个作业重复写代码。一些定制化的任务除外。
# 基础类型选择
state.backend: rocksdb # 或 hashmap(不推荐生产)
# RocksDB 专用调优(生产必配)
state.backend.rocksdb.incremental: true
state.backend.rocksdb.block-cache-size: 256MB
state.backend.rocksdb.write-buffer-size: 64MB
state.backend.rocksdb.write-buffer-number: 4
state.backend.rocksdb.memory.managed: true
# Checkpoint 基础配置(配合状态后端使用)
state.backend.rocksdb.local-dir: /data/flink/rocksdb # 本地 SSD 路径,必配
state.checkpoints.dir: hdfs:///flink/checkpoints
作业启动时优先级:代码配置 > 配置文件配置(代码会覆盖 yaml)。
在Flink中,状态实际还和Checkpoint,以及容错机制关系密切,由于时间问题,今天就先聊到这里了,下一回咱们继续,聊Checkpoint、容错机制、以及关于checkpoint和状态的生产级调优与故障排查。
更多推荐
所有评论(0)