)
Flink 2.3.0 从理论到实践 —— 第 4 章 状态管理State课程定位状态是 Flink 区别于普通流处理引擎的核心能力。本章深入 Keyed State 与 Operator State 的区别、三种状态后端HashMap / ForSt / RocksDB、状态 TTL 与过期策略、状态大小监控与 Schema 演进——这些都是后续 Checkpoint 容错、性能调优、故障恢复的理论基础。版本基线Flink 2.3.0 ForSt 状态后端章节导读4.1 有状态计算的概念4.2 Keyed State vs Operator State4.3 状态后端HashMap / ForSt / RocksDB4.4 状态类型详解4.5 状态 TTL 与过期策略4.6 状态大小估算与监控4.7 状态迁移与 Schema 演进4.8 本章小结与下章预告4.1 有状态计算的概念4.1.1 什么是状态状态State是 Flink 算子在处理数据时记住的历史信息。例如聚合作业每辆车当天的累计里程 之前所有报文里程之和累加状态Join 作业左侧流的当前数据要与右侧流的某些历史数据匹配缓存状态CEP 作业检测5 分钟内连续 3 次告警需要记住已发生的告警模式状态无状态计算: 输入 → 算子 → 输出 每条数据独立处理,与历史无关 有状态计算: 输入 [历史状态] → 算子 → 输出 [更新后的状态] 算子记住历史,基于历史做决策4.1.2 为什么状态如此重要能力没有状态有状态聚合每条数据独立无法累加维护累加器逐条更新Join无法跨数据匹配缓存一侧数据等另一侧到达窗口无法累积一段时间数据缓冲窗口内所有数据CEP无法检测跨数据模式维护模式匹配进度容错数据丢失无法恢复Checkpoint 快照后可恢复结论Flink 的实时数仓 / 风控 / CEP / 机器学习等场景几乎全是有状态计算。理解状态是理解 Flink 的关键。4.2 Keyed State vs Operator StateFlink 有两类状态Keyed State按 key 分区和Operator State算子级。4.2.1 对比维度Keyed StateOperator State粒度按 key 分区每个 key 一份算子级每个 SubTask 一份前提必须先keyBy不需要keyBy典型算子reduce/aggregate/window/processKafka Sourceoffset/ListCheckpointed恢复方式按 key 重新分配按 SubTask 重新分配可自定义APIValueState/ListState/MapState等ListState/ 自定义CheckpointedFunction4.2.2 Keyed State 示例// 每辆车当天的累计里程publicclassDailyMileageFunctionextendsKeyedProcessFunctionString,VehicleEvent,Tuple2String,Double{// Keyed State: 按车辆 VIN 分区,每辆车一份privateValueStateDoublemileageState;Overridepublicvoidopen(Configurationparameters){ValueStateDescriptorDoubledescriptornewValueStateDescriptor(dailyMileage,Double.class);// 配置 TTL (见 4.5)descriptor.enableTimeToLive(StateTtlConfig.newBuilder(Time.hours(48))// 48 小时过期.setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite).build());mileageStategetRuntimeContext().getState(descriptor);}OverridepublicvoidprocessElement(VehicleEventevent,Contextctx,CollectorTuple2String,Doubleout)throwsException{DoublecurrentmileageState.value();if(currentnull)current0.0;currentevent.getMileage();mileageState.update(current);out.collect(Tuple2.of(event.getVin(),current));}}4.2.3 Operator State 示例// Kafka Source 内部用 Operator State 存 offsetpublicclassKafkaSourceFunctionextendsRichSourceFunctionEventimplementsCheckpointedFunction{// Operator State: 每个 SubTask 一份,不按 key 分区privatetransientListStateLongoffsetState;OverridepublicvoidinitializeState(FunctionInitializationContextcontext)throwsException{ListStateDescriptorLongdescriptornewListStateDescriptor(kafkaOffsets,Long.class);offsetStatecontext.getOperatorStateStore().getListState(descriptor);}OverridepublicvoidsnapshotState(FunctionSnapshotContextcontext)throwsException{offsetState.clear();offsetState.add(currentOffset);// 保存当前 offset}}4.2.4 选型决策作业里有 keyBy 吗? ├─ 是 → Keyed State │ 可用 ValueState/ListState/MapState/ReducingState/AggregatingState └─ 否 → Operator State 用于 Source 的 offset / 自定义算子的非分区状态4.3 状态后端HashMap / ForSt / RocksDB状态后端State Backend决定了状态如何存储。Flink 2.x 提供三种选择。4.3.1 三种后端对比维度HashMapForSt2.x 推荐RocksDB1.x 主流存储位置TM 堆内存TM 本地磁盘RocksDB 改进版TM 本地磁盘状态大小受 TM 堆限制GB 级受磁盘限制TB 级受磁盘限制TB 级访问延迟微秒级最快毫秒级毫秒级Checkpoint直接写文件快增量快照高效增量快照内存管理受 GC 影响托管内存不受 GC托管内存2.x 状态默认小状态推荐大状态兼容保留4.3.2 ForStFlink 2.x 的新选择Flink 2.x 变更ForSt 是 Flink 2.0 引入的新的状态后端作为 RocksDB 的替代实现。它基于 RocksDB 但做了深度优化更好的 Checkpoint 性能本地快照 增量传输更精细的托管内存控制与 Flink 2.x 内存模型深度集成# flink-conf.yaml 配置 ForStstate.backend:forststate.backend.local-recovery:true# 本地恢复,加速重启state.backend.forst.memory.managed:true# 托管内存模式state.backend.incremental:true# 增量 Checkpoint4.3.3 选型决策树状态总大小 1 GB 且要求最低延迟? ├─ 是 → HashMap │ 注意: 状态膨胀会 OOM,建议严格 TTL └─ 否 状态总大小 100 GB? ├─ 是 → ForSt 增量 Checkpoint │ 推荐: Flink 2.x 大多数场景 └─ 否 → ForSt 增量 本地恢复 极大状态场景,需配 SSD 大磁盘4.3.4 生产建议场景后端配置维表 Lookup 缓存HashMapstate.backend: hashmap TTL 1h聚合作业小状态HashMap TTL Mini-Batch聚合作业大状态ForSt 增量 Checkpoint 本地恢复Join 作业大状态ForSt RocksDB Options 调优详见第 17 章4.4 状态类型详解4.4.1 Keyed State 类型类型说明典型用途ValueState单值状态累加器、最新值、计数器ListState列表状态窗口数据缓存、CEP 模式匹配MapStateKV 映射状态维表缓存、按子 key 聚合ReducingState自动 Reduce 的状态SUM/MIN/MAX 聚合AggregatingState复杂聚合状态IN/OUT 类型不同AVG 加权平均4.4.2 完整示例四种状态并用publicclassVehicleStatsFunctionextendsKeyedProcessFunctionString,VehicleEvent,Stats{privateValueStateDoubletotalMileage;// 累计里程privateListStateVehicleEventeventBuffer;// 事件缓存(窗口)privateMapStateString,IntegereventTypeCount;// 各事件类型计数privateReducingStateDoublemaxSpeed;// 最大速度Overridepublicvoidopen(Configurationparameters){// ValueStatetotalMileagegetRuntimeContext().getState(newValueStateDescriptor(totalMileage,Double.class));// ListStateeventBuffergetRuntimeContext().getListState(newListStateDescriptor(eventBuffer,VehicleEvent.class));// MapStateeventTypeCountgetRuntimeContext().getMapState(newMapStateDescriptor(eventTypeCount,String.class,Integer.class));// ReducingStatemaxSpeedgetRuntimeContext().getReducingState(newReducingStateDescriptor(maxSpeed,(a,b)-Math.max(a,b),Double.class));}OverridepublicvoidprocessElement(VehicleEventevent,Contextctx,CollectorStatsout)throwsException{// ValueState: 累加里程DoublemileagetotalMileage.value();if(mileagenull)mileage0.0;totalMileage.update(mileageevent.getMileage());// ListState: 缓存事件eventBuffer.add(event);// MapState: 按 event_type 计数IntegercnteventTypeCount.get(event.getType());if(cntnull)cnt0;eventTypeCount.put(event.getType(),cnt1);// ReducingState: 自动取最大值maxSpeed.add(event.getSpeed());// 输出统计StatsstatsnewStats();stats.setVin(event.getVin());stats.setTotalMileage(totalMileage.value());stats.setMaxSpeed(maxSpeed.get());out.collect(stats);}}4.5 状态 TTL 与过期策略4.5.1 为什么需要 TTL流作业长期运行状态会持续膨胀车辆VIN 注册时间 最新事件时间 状态大小 VIN001 2026-01-01 2026-03-01 累计 60 天数据 VIN002 2026-01-15 2026-03-01 累计 45 天数据 ... → 不加 TTL,状态会无限增长,最终 OOM4.5.2 TTL 配置StateTtlConfigttlConfigStateTtlConfig.newBuilder(Time.hours(48))// TTL 48 小时.setTtlTimeCharacteristic(...)// 时间语义.setUpdateType(UpdateType.OnCreateAndWrite)// 更新策略.setStateVisibility(StateVisibility.NeverReturnExpired)// 可见性.setCleanupStrategies(...)// 清理策略.build();descriptor.enableTimeToLive(ttlConfig);4.5.3 关键参数参数选项说明时间语义ProcessingTime/EventTime建议用 ProcessingTime更稳定更新策略OnCreateAndWrite/OnReadAndWrite写时更新 vs 读写都更新可见性NeverReturnExpired/ReturnExpiredIfNotCleanedUp过期立即不可见 vs 清理前仍可读清理策略FULL_STATE_SCAN_SNAPSHOT/INCREMENTAL_CLEANUP/IN_HEAP全量扫描 / 增量清理 / 堆内即时清理生产建议来自项目经验聚合作业状态 TTL 设 48 小时覆盖 1 天 容错窗口维表缓存 TTL 设 1 小时避免缓存脏数据用ProcessingTime语义不依赖 Watermark 推进NeverReturnExpired保证业务正确性4.6 状态大小估算与监控4.6.1 状态大小估算状态类型单 key 大小估算公式ValueStatevalue 类型大小N_keys × value_sizeListState元素数 × 元素大小N_keys × N_elements × element_sizeMapState条目数 × 条目大小N_keys × N_entries × entry_sizeReducingState单值N_keys × value_size示例10000 辆车 × 平均 100 个事件类型 × 每条目 50 字节 50 MB4.6.2 监控指标通过 Web UI 或 REST API 查看状态大小# 作业状态大小curlhttp://flink-master:8081/jobs/job-id/checkpoints# {# counts: {...},# summary: {# state_size: 5368709120 # 5 GB# }# }# 各算子状态curlhttp://flink-master:8081/jobs/job-id/vertices关键监控指标指标告警阈值处理状态总大小 TM 内存 80%收紧 TTL / 切 ForSt状态增长率每日 10%检查 TTL 是否生效Checkpoint 时长 30s优化状态 / 切增量Checkpoint 失败率 5%检查存储 / 网络4.6.3 状态膨胀排查状态持续增长? ├─ TTL 是否生效? │ └─ 检查 StateTtlConfig 是否正确配置 ├─ TTL 时间语义? │ └─ EventTime 但 Watermark 不推进 → 用 ProcessingTime ├─ 清理策略? │ └─ 增量清理: cleanupInBackground 设为 true └─ 业务逻辑漏更新? └─ 某些 key 的状态从未被读取,需要主动清理4.7 状态迁移与 Schema 演进4.7.1 Schema 演进场景业务需求变化导致状态类型变化变更类型是否兼容说明加字段nullable✅旧状态读到新字段为 null加字段非空❌旧状态没有默认值需无状态重启删字段✅旧状态读到字段被忽略改字段类型❌序列化不兼容改字段名❌视为删除加新字段4.7.2 演进 SOP1. 评估变更类型(参考上表) 2. 若不兼容 → 停作业(savepoint) → 改代码 → 无状态重启 3. 若兼容 → 停作业(savepoint) → 改代码 → 从 savepoint 恢复 4. 验证数据正确性项目硬约束改表结构的铁律是先停受影响作业及其全部下游 → DDL → 全部无状态重启。有状态恢复会因源快照不连续报OutOfRangeException崩溃循环。4.7.3 状态迁移工具# 1. 生成 Savepointflink savepointjob-idhdfs:///savepoints/# 2. 升级代码 / 配置# 3. 从 Savepoint 恢复(无状态重启需加 -n)flink run-shdfs:///savepoints/savepoint-xxx-dmy-job.jar# 无状态重启: flink run -n -d my-job.jar (丢弃状态)4.8 本章小结与下章预告本章小结┌────────────────────────────────────────────────────────────────┐ │ 第 4 章 要点回顾 │ └────────────────────────────────────────────────────────────────┘ ✓ 状态 算子记住的历史,是 Flink 区别普通流引擎的核心 聚合/Join/窗口/CEP 都依赖状态 ✓ Keyed State vs Operator State: Keyed: 按 key 分区,需要 keyBy Operator: 算子级,用于 Source offset 等 ✓ 三种状态后端: HashMap: 堆内存,最快,小状态( 1GB) ForSt: 2.x 推荐,大状态,TB 级,增量 Checkpoint RocksDB: 1.x 主流,2.x 兼容保留 ✓ 五种 Keyed State 类型: ValueState / ListState / MapState / ReducingState / AggregatingState ✓ TTL: 防状态膨胀 推荐: ProcessingTime OnCreateAndWrite NeverReturnExpired 聚合 48h,维表缓存 1h ✓ 监控: 状态总大小 增长率 Checkpoint 时长 告警: TM 内存 80% / 日增 10% / Checkpoint 30s ✓ Schema 演进: 兼容(加 nullable 字段/删字段) → savepoint 恢复 不兼容(改类型/加非空) → 无状态重启 铁律: 停作业 → DDL → 无状态重启下章预告第 5 章 时间语义与 Watermark讲解 Event Time / Processing Time / Ingestion Time 三种时间语义、Watermark 的生成与传递机制、Flink 2.3 的 Watermark 对齐增强、迟到数据处理Allowed Lateness / Side Output。时间是窗口与 Checkpoint 的基础。官方参考资料Flink State 官方文档https://nightlies.apache.org/flink/flink-docs-stable/docs/concepts/state/State Backendhttps://nightlies.apache.org/flink/flink-docs-stable/docs/ops/state/state_backends/State TTLhttps://nightlies.apache.org/flink/flink-docs-stable/docs/dev/datastream/fault_tolerance/state_ttl/ForSt 状态后端https://nightlies.apache.org/flink/flink-docs-stable/docs/ops/state/forst_state_backend/