Flink SQL多表Join的状态存储与内存占用疑问
Flink三表Upsert-Kafka关联的状态存储与优化问题解答
一、关联执行时的状态存储内容
结合你看到的两节点执行计划,各节点的状态存储细节如下:
- 第一个节点(source1 ⇨ INNER JOIN ⇨ source2):
由于是基于col1主键的Upsert流关联,节点会以col1为键,分别存储:- source1中每个
col1对应的最新col2值(仅保留关联所需字段) - source2中每个
col1对应的最新col2值
只有当两边存在相同col1的记录时,才会输出<s1.col1, s1.col2, s2.col2>的中间Upsert流(当source1或source2的某个col1更新时,会同步输出更新后的中间结果)。
- source1中每个
- 第二个节点(中间结果 ⇨ INNER JOIN ⇨ source3):
同样基于col1关联,节点会存储:- 上游中间结果中每个
col1对应的最新<s1.col2, s2.col2>组合 - source3中每个
col1对应的最新col2值
当中间结果与source3存在相同col1时,输出最终的<s1.col2, s2.col2, s3.col2>结果。
- 上游中间结果中每个
二、你的假设是否正确?
你的假设不完全正确:
嵌套关联确实会存储中间结果的状态,但因为所有流都是Upsert类型(基于主键更新),每个主键只会保留最新的状态条目,不会存储全量历史数据。比如source1有1000个不同col1,状态中只会存1000条记录,而非所有历史更新。
如果三个表的col1主键集合高度重叠,中间结果的状态条目数不会超过单个表的主键数,此时总状态大小不会显著超过三个源表的状态总和;即便主键集合差异较大,中间结果带来的额外存储开销也远达不到“大幅超出总和”的程度。
三、Flink的优化机制
Flink SQL针对这类关联场景提供了多种优化手段,可有效减少状态存储:
- 星型关联重写:当多个JOIN都基于同一个主键时,Flink优化器会自动将嵌套JOIN计划重写为星型关联——让source1分别与source2、source3做JOIN后合并结果。这样无需存储
<source1,source2>的中间结果,仅需维护三个源表的状态,直接降低状态总量。 - Upsert流增量状态维护:Upsert-Kafka数据流本身是主键驱动的更新流,Flink处理时会自动用新记录替换对应主键的旧状态,不会累积历史数据,从根源上控制状态规模。
- 状态TTL配置:可为关联状态设置生存时间(TTL),对长时间无更新的主键自动清理对应状态条目,释放内存/存储资源。需注意的是,TTL会清理旧数据,要确保业务场景允许该行为。
- 状态后端选型:使用RocksDB状态后端时,部分状态会持久化到磁盘,仅在内存中保留热点数据,大幅降低JVM堆内存占用,适配大状态关联场景。
内容的提问来源于stack exchange,提问作者Bogey
相关产品推荐
相关产品推荐

