启用ObjectReuse后Flink Table API仍因RowData序列化高CPU占用
Flink Table API流水线序列化CPU开销优化问题
我有一个基于Flink Table API的流水线,针对15列执行1分钟滚动窗口的计数聚合操作。尽管已启用ObjectReuse、MiniBatch和本地两阶段聚合优化,但火焰图显示约40%的CPU消耗在每行的序列化操作上。
执行计划
SecurityEvents[1] -> MiniBatchAssigner[2] -> Calc[3] -> WindowTableFunction[4] -> Correlate[5] -> Calc[6] -> LocalWindowAggregate[7]
当前配置
- 从列向量化的
DataStream<RowData>转换为Table - 每行包含一个
authorizations<Row>数组,会被执行展开(unnested)操作 - 火焰图显示每个
Calc节点(对应生成的字节码StreamExecCalc$33.processElement_split14:-1)会复制投影行 - 大部分CPU(约40%)消耗在
StringDataSerializer.copy方法上
代码片段
tableEnv .sqlQuery( """ SELECT * FROM TABLE( TUMBLE( TABLE SecurityEvents, DESCRIPTOR(rowtime), INTERVAL '1' MINUTE ) ) LEFT JOIN UNNEST(authorizations) ON TRUE """.trimIndent() ) // 必须同时按窗口起始和结束时间分组,才能让该表成为可转换回DataStream的仅追加表 .groupBy( *topLevelGroupings.toCol(), *authorizationsGroupings.toCol(), *attributesGroupings, col("window_start"), col("window_end"), col("window_time") ) .select( *topLevelGroupings.toCol(), *authorizationsGroupings.toCol { "authorization_$it" }, *authorizationAttributesGroupings.toCol { "authorization_attribute_$it" }, col("window_start"), col("window_time").`as`("rowtime"), col("dataset").count().`as`("count") )
核心问题
火焰图显示每个Calc节点会复制15个分组列对应的投影行,大部分CPU消耗在StringDataSerializer.copy(约40%)。即使启用了ObjectReuse,且LocalWindowAggregate已通过RowDataSerilizer.toBinaryRow复制键值,该问题仍存在。
疑问
- 为何启用ObjectReuse后仍会复制每行数据?
- 是否有办法在Table API中避免此序列化开销?
- 除切换到DataStream API(需重新实现本地聚合和滚动窗口)外,还有其他替代方案吗?
环境信息
- Apache Flink版本:1.20.0
- 检查点间隔:5分钟
内容的提问来源于stack exchange,提问作者Tomás Cerdá
相关产品推荐
相关产品推荐

