You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

启用ObjectReuse后Flink Table API仍因RowData序列化高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复制键值,该问题仍存在。

疑问

  1. 为何启用ObjectReuse后仍会复制每行数据?
  2. 是否有办法在Table API中避免此序列化开销?
  3. 除切换到DataStream API(需重新实现本地聚合和滚动窗口)外,还有其他替代方案吗?

环境信息

  • Apache Flink版本:1.20.0
  • 检查点间隔:5分钟

内容的提问来源于stack exchange,提问作者Tomás Cerdá

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.12 10:42:41