一、为什么状态是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个算子状态:
image.png
算子状态有三种结构类型: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键来维护每一个状态
image.png
键控状态有三种结构类型:ValueState、ListState、MapState、ReducingState、AggregatingState

2.2.1 ValueState(值状态)

特点:存储单一可更新值,支持 update()value()clear()
场景:记录最新用户属性、会话计数器。
APIgetRuntimeContext().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()
典型场景:收集历史事件、缓冲队列。
APIgetRuntimeContext().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 操作。
典型场景:用户画像标签、动态配置映射。
APIgetRuntimeContext().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()
典型场景:实时求和、求最大值(无需手动维护中间结果)。
APIgetRuntimeContext().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)、带中间状态的聚合。
APIgetRuntimeContext().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 StreamkeyBy() 之后)的 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(默认):创建或写入时刷新 TTL
    • OnReadAndWrite:读写都刷新(适合“最后访问时间”场景)
  • 过期可见性(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 StateKeyed 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 性能。配置方式主要有代码设置(灵活,适合开发/测试)和配置文件设置(集群统一管理,适合生产)。生产环境中,有些任务可能需要单独配置,可以使用代码设置

  1. 推荐配置决策(快速选择)
    • 小状态(< 几百 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和状态的生产级调优与故障排查。

Logo

腾讯云面向开发者汇聚海量精品云计算使用和开发经验,营造开放的云计算技术生态圈。

更多推荐