Flink SQL流应用算子名称优化及编号含义技术咨询
背景信息
- Kafka主题关联表:输入表
table1、table2,输出表table3 - 对应Flink SQL流作业代码:
CREATE TEMPORARY VIEW distinct_table1 AS SELECT * FROM (SELECT *, ROW_NUMBER() OVER (PARTITION BY id ORDER BY change_date desc) AS rownum FROM table1) WHERE rownum = 1; CREATE TEMPORARY VIEW distinct_table2 AS SELECT * FROM (SELECT *, ROW_NUMBER() OVER (PARTITION BY id ORDER BY change_date desc) AS rownum FROM table2) WHERE rownum = 1; Insert into table3 Select t1.col1,t1.col2..,t1.coln,t2.col1,t2.col2,t2.coln From distinct_table1 t1 inner join distinct_table2 t2 on t1.id=t2.t1_id
- 当前算子名称生成示例:
- Source算子:
Source: table1[7] -> MiniBatchAssigner[8] -> Calc[9] - 去重SQL生成算子:
Rank[20] -> Calc[21] - Join语句生成算子:
Join[23] -> Calc[24]
- Source算子:
技术问题
- 为实现高效监控,是否支持将表名纳入算子名称中?
- 算子名称括号中的数字(如[7]、[20]等)代表什么?是否与TaskSlot分配有关联?
问题解答
1. 表名纳入算子名称的支持情况
Flink完全支持将表名纳入算子名称中,有两种实用方式:
- 全局配置开启:在Flink 1.13及以上版本中,设置配置项
table.exec.operator-name-include-table为true,执行计划中的算子会自动关联对应表名/视图名,比如去重的Rank算子会显示为Rank(distinct_table1)[20],Join算子会关联参与关联的表名。 - 查询Hint自定义命名:在SQL中使用
/*+ OPTIONS('operator.name'='xxx') */的Hint语法,手动指定包含表名的算子名称。例如修改去重逻辑的SQL:
CREATE TEMPORARY VIEW distinct_table1 AS SELECT * FROM (SELECT *, /*+ OPTIONS('operator.name'='table1_Rank') */ ROW_NUMBER() OVER (PARTITION BY id ORDER BY change_date desc) AS rownum FROM table1) WHERE rownum = 1;
生成的Rank算子名称会显示为table1_Rank[20],直接关联表名便于监控识别。
2. 算子名称中数字的含义及与TaskSlot的关系
算子名称括号里的数字是算子的唯一ID,是Flink生成执行计划时为每个算子节点分配的编号,用于内部区分不同算子、追踪执行计划链路、定位日志中的算子相关信息。
该数字和TaskSlot分配没有直接关联:TaskSlot是Flink的资源分配单元,每个Slot可运行多个算子的任务实例;而算子ID仅属于执行计划的节点标识,与资源分配逻辑完全独立。
内容的提问来源于stack exchange,提问作者Faisal Ahmed Siddiqui
相关产品推荐
相关产品推荐

