Flink 1.20中ListView比ArrayList性能更优的原因及转换逻辑问询
Flink 1.20自定义中位数UDAF:ArrayList vs ListView性能差异解析
一、性能差距的核心原因
- 内存管理与持久化策略差异
- 使用
ArrayList作为状态时,整个集合会被当作单一状态对象存储。每次更新状态(如添加元素),Flink需要序列化整个ArrayList并写入状态后端,数据量越大,序列化/反序列化开销、IO成本越高,极端情况还会触发OOM(全量加载到内存)。 ListView是Flink专为批量数据设计的DataView实现,采用分段存储+惰性加载机制:数据被分成多个小块(segment)存到状态后端,仅当访问特定元素时才加载对应分段,避免全量序列化/反序列化开销;同时利用状态后端的增量更新能力,只同步修改过的分段,大幅降低IO和序列化成本。
- 使用
- 状态访问模式适配
中位数计算需要频繁添加元素和排序/中间值查找。ArrayList内存连续存储,查找快但状态更新成本极高;ListView针对批量追加场景优化,状态更新开销远低于ArrayList,即使排序也可按需加载分段到内存处理,整体吞吐量更高。
二、1.20版本中ListView的状态转换逻辑位置与调试方法
转换逻辑的新位置
Flink 1.14之后,原AggregationCodeGenerator的相关逻辑被重构,ListView的状态转换逻辑现在分布在:
org.apache.flink.table.planner.codegen.agg.AggsHandlerCodeGenerator的generateAccumulatorCode方法中,会调用DataViewCodeGenerator处理ListView这类DataView的代码生成。- 深层状态后端适配逻辑在
org.apache.flink.table.runtime.dataview.ListView实现中,当元素数量超过阈值(默认由state.backend.listview.state.threshold配置,默认值100)时,自动切换为状态后端存储(而非内存)。
观察与调试方法
- 配置调试日志
开启指定包的DEBUG级别日志,可查看ListView切换状态后端的时机和代码生成细节:<Logger name="org.apache.flink.table.runtime.dataview" level="DEBUG"/> <Logger name="org.apache.flink.table.planner.codegen" level="DEBUG"/> - 断点调试
- 在
ListView.add(Object value)方法设断点,观察内部state变量类型:数据量小时用ArrayList内存存储;超过阈值后切换为StateListView(基于状态后端的实现)。 - 在
AggsHandlerCodeGenerator.generateAccumulatorCode方法中,跟踪DataViewCodeGenerator.generateDataViewSerializer和generateDataViewStateAccess的调用,查看ListView对应的状态访问代码生成过程。
- 在
- 状态后端验证
通过Flink Web UI的状态页面查看UDAF状态结构:使用ListView会看到多个分段命名的状态条目;ArrayList则是单一状态对象。
内容的提问来源于stack exchange,提问作者Apollo Elon
相关产品推荐
相关产品推荐

