【问题标题】:What is the relevant rules of Flink Window TVF and CEP SQL?Flink Window TVF和CEP SQL的相关规则是什么?
【发布时间】:2023-01-01 09:05:24
【问题描述】:

我正在尝试解析 Flink 窗口化 TVF sql 列级沿袭,我初始化了一个自定义 FlinkChainedProgram 并设置了一些 Opt 规则。

除了 Window TVF SQL 和 CEP SQL 外,大多数情况下工作正常。

例如,我得到一个合乎逻辑的计划

insert into sink_table(f1, f2, f3, f4) 
       SELECT cast(window_start as String),
              cast(window_start as String), 
              user_id, 
              cast(SUM(price) as Bigint) 
       FROM TABLE(TUMBLE(TABLE source_table, DESCRIPTOR(event_time), INTERVAL '10' MINUTES))
       GROUP BY window_start, window_end, GROUPING SETS ((user_id), ());

rel#1032:FlinkLogicalCalc.LOGICAL.any.None: 0.[NONE].[NONE](input=FlinkLogicalAggregate#1030,select=CAST(window_start) AS EXPR$0, CAST(window_start) AS EXPR$1, null:BIGINT AS EXPR$2, user_id, null:VARCHAR(2147483647) CHARACTER SET "UTF-16LE" AS EXPR$4, CAST($f4) AS EXPR$5)

如我们所见,优化的 RelNode 树包含空列,因此 MetadataQuery 无法获取原始列信息。

我应该在逻辑优化阶段设置什么规则来解析 Window TVF SQL 和 CEP SQL?谢谢

【问题讨论】:

    标签: flink-sql apache-calcite sql-parser


    【解决方案1】:

    我解决了Flink CEP SQL的字段血缘关系方法,在org.apache.calcite.rel.metadata.org.apache.calcite.rel.metadata. RelMdColumnOrigins中添加了getColumnOrigins(Match rel, RelMetadataQuery mq, int iOutputColumn)方法:

    /**
     * Support field blood relationship of CEP.
     * The first column is the field after PARTITION BY, and the other columns come from the measures in Match
     */
    public Set<RelColumnOrigin> getColumnOrigins(Match rel, RelMetadataQuery mq, int iOutputColumn) {
        if (iOutputColumn == 0) {
            return mq.getColumnOrigins(rel.getInput(), iOutputColumn);
        }
        final RelNode input = rel.getInput();
        RexNode rexNode = rel.getMeasures().values().asList().get(iOutputColumn - 1);
    
        RexPatternFieldRef rexPatternFieldRef = searchRexPatternFieldRef(rexNode);
        if (rexPatternFieldRef != null) {
            return mq.getColumnOrigins(input, rexPatternFieldRef.getIndex());
        }
        return null;
    }
    
    private RexPatternFieldRef searchRexPatternFieldRef(RexNode rexNode) {
        if (rexNode instanceof RexCall) {
            RexNode operand = ((RexCall) rexNode).getOperands().get(0);
            if (operand instanceof RexPatternFieldRef) {
                return (RexPatternFieldRef) operand;
            } else {
                // recursive search
                return searchRexPatternFieldRef(operand);
            }
        }
        return null;
    }
    

    来源地址:https://github.com/HamaWhiteGG/flink-sql-lineage/blob/main/src/main/java/org/apache/calcite/rel/metadata/RelMdColumnOrigins.java

    我已经给出了详细的测试用例,你可以参考:https://github.com/HamaWhiteGG/flink-sql-lineage/blob/main/src/test/java/com/dtwave/flink/lineage/cep/CepTest.java

    Flink CEP SQL 测试用例:

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2013-05-31
      • 1970-01-01
      • 1970-01-01
      • 2016-12-10
      • 2018-01-26
      • 1970-01-01
      相关资源
      最近更新 更多