Flink 2.3.0 从理论到实践 —— 第 4 章 状态管理(State)
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/

相关新闻

1.两数之和 - 力扣(Leetcode)

1.两数之和 - 力扣(Leetcode)

题目 给定一个整数数组 nums 和一个整数目标值 target,请你在该数组中找出 和为目标值 target 的那 两个 整数,并返回它们的数组下标。 你可以假设每种输入只会对应一个答案,并且你不能使用两次相同的元素。 你可以按任意顺序返回答案。 示例…

2026/10/10 2:56:07 阅读更多 →
基于PCA9422 PMIC与STM32F401RB的锂电池供电与低功耗电源管理实战

基于PCA9422 PMIC与STM32F401RB的锂电池供电与低功耗电源管理实战

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/10/10 2:55:07 阅读更多 →
2026年程序员必学大模型:小白也能掌握的AI应用开发,收藏学习路线!

2026年程序员必学大模型:小白也能掌握的AI应用开发,收藏学习路线!

文章指出,2026年AI大模型相关岗位需求暴涨,薪资远超传统后端岗位。针对程序员学习痛点,文章强调AI应用开发无需精通深度学习,只需调用API并搭建业务系统。文章提供了零基础必备工具环境、四个阶段的学习路线(Prompt工程…

2026/10/10 2:55:07 阅读更多 →

最新新闻

GPS天线设计 GNSS天线设计建议

GPS天线设计 GNSS天线设计建议

GPS天线设计 GNSS天线设计建议 天线作为导航定位设备中最重要的接收器件,它起到的作用就像是人的“耳朵”;是将卫星发送下来的电磁波能量变换成电子器件可解析的电流。因此天线的性能好坏将直接关系到GPS整机的产品性能。目前GNSS系统开放民用定位系统主要是美国GPS…

2026/10/10 5:16:30 阅读更多 →
Python实战:不规则JSON解析的容错技巧

Python实战:不规则JSON解析的容错技巧

真实项目里摸爬滚打的同学,大概率都遇到过这种场面:接口文档写得清清楚楚,联调时返回的 JSON 却一个比一个“野”。字段时有时无,价格一会儿是数字一会儿是字符串,嵌套结构深浅不一,偶尔还直接甩给你一个 J…

2026/10/10 5:16:30 阅读更多 →
vsode配置settings.json

vsode配置settings.json

一、打开方式命令面板运行“首选项:打开用户设置 (JSON)”命令 (CtrlShiftP)打开 settings.json 文件来更改默认设置二、配置文件内容{// 编辑器基本配置// 设置编辑器字体大小为 16"editor.fontSize": 16,// 控制字体样式"editor.fontFamily":…

2026/10/10 5:16:30 阅读更多 →
解密 Rust 裸指针操作的内存对齐(Alignment):未对齐访问的硬件陷阱与 read_unaligned 的真实代价

解密 Rust 裸指针操作的内存对齐(Alignment):未对齐访问的硬件陷阱与 read_unaligned 的真实代价

在编写系统底层网络协议解析、二进制序列化引擎或者直接与硬件寄存器打交道时,我们经常需要把一段原始的字节切片(&[u8])强行转译为高级结构体或整型数字。 很多从 C/C 转过来的开发者,习惯随手敲下这样的转换: //…

2026/10/10 5:16:30 阅读更多 →
vue使用el-tree 数据回显问题

vue使用el-tree 数据回显问题

treeMenus.forEach(menu > {if(!menu.hashChildren){//如果没有子节点,就勾选,这样就可以在父节点上有半选状态this.$refs.menuTree.setChecked(menu.id, true, false);} })menuTree:标签中设置的refmenu.id:节点的值找到叶子节…

2026/10/10 5:16:30 阅读更多 →
基于预训练技术的BIM与IoT数据融合及偏差预警算法实战

基于预训练技术的BIM与IoT数据融合及偏差预警算法实战

简介:这份文档面向建筑施工管理、BIM工程与智能建造方向的技术人员及研究者,围绕施工进度管控中数据维度单一、偏差预警滞后等痛点,给出基于DeepSeek预训练技术的BIM与IoT数据融合及偏差预警算法方案。全文共196页、50个大章节,从…

2026/10/10 5:15:29 阅读更多 →

日新闻

卫星轨道分类全解析:从LEO到GEO的选型逻辑与工程实践

卫星轨道分类全解析:从LEO到GEO的选型逻辑与工程实践

1. 从“卫星轨道分类”这个标题说起:为什么值得花时间搞懂第一次接触“卫星轨道分类”这个概念,很多人会觉得它离自己很远——不就是天上的星星怎么转吗?但如果你正在做航天任务规划、遥感数据接收、星座设计,甚至只是准备一场航天…

2026/10/10 0:00:39 阅读更多 →
Spring AOP 核心原理与实战:从概念到日志切面落地

Spring AOP 核心原理与实战:从概念到日志切面落地

1. 从一个真实痛点说起:为什么你的代码里到处都是重复逻辑刚入行那会儿,我写过一个用户管理模块,注册、登录、改密码、注销四个接口。每个接口里都塞了几乎一样的日志打印、参数校验、事务开启和提交。当时觉得没什么,能跑就行。直…

2026/10/10 0:00:40 阅读更多 →
Python招聘数据采集与分析可视化:从采集清洗到薪资技能城市可视化全链路

Python招聘数据采集与分析可视化:从采集清洗到薪资技能城市可视化全链路

简介:这是一套面向计算机相关专业学生与项目实战学习者的Python数据采集与分析可视化完整项目,以Boss直聘岗位数据为对象,适合用作毕业设计、课程设计或期末大作业。资源包共38个文件,约246KB,以13个py源码文件为核心&…

2026/10/10 0:00:40 阅读更多 →

周新闻

KT148A语音芯片外挂8002D功放的工程实践指南

KT148A语音芯片外挂8002D功放的工程实践指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/10/8 15:26:32 阅读更多 →
LLC谐振变换器增益公式推导:从FHA等效到完整归一化表达式

LLC谐振变换器增益公式推导:从FHA等效到完整归一化表达式

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/10/10 1:36:08 阅读更多 →
ARM架构深度解析:从RISC设计理念到交叉编译实战

ARM架构深度解析:从RISC设计理念到交叉编译实战

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/10/9 10:11:06 阅读更多 →

月新闻

我发现了一个新思路:用 Remotion + Claude Code 像写代码一样自动化生成短视频

我发现了一个新思路:用 Remotion + Claude Code 像写代码一样自动化生成短视频

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/10/8 21:13:17 阅读更多 →
Windows下 Codex 中 Chrome 和 Computer Use 插件不可用问题排查及解决参考方式:TaoToken 统一 Key 配置与验证

Windows下 Codex 中 Chrome 和 Computer Use 插件不可用问题排查及解决参考方式:TaoToken 统一 Key 配置与验证

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/10/9 21:32:20 阅读更多 →
黑夜航拍船只数据集训练YOLOV5模型全流程解析

黑夜航拍船只数据集训练YOLOV5模型全流程解析

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/10/9 6:17:20 阅读更多 →