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

Apache Calcite 1.39.0流式聚合报错:GROUP BY需含单调表达式的解决

Apache Calcite流处理GROUP BY验证报错解决方法

问题描述

使用Apache Calcite 1.39.0进行有状态流处理开发,验证自定义Schema上的GROUP BY查询时执行报错,尝试标记字段单调属性但找不到对应重写方法。

代码示例

package calcite.streaming;

import org.apache.calcite.rel.type.RelDataType;
import org.apache.calcite.rel.type.RelDataTypeFactory;
import org.apache.calcite.schema.*;
import org.apache.calcite.schema.impl.AbstractSchema;
import org.apache.calcite.schema.impl.AbstractTable;
import org.apache.calcite.sql.SqlNode;
import org.apache.calcite.sql.parser.SqlParser;
import org.apache.calcite.tools.FrameworkConfig;
import org.apache.calcite.tools.Frameworks;
import org.apache.calcite.tools.Planner;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.sql.Timestamp;
import java.util.HashMap;
import java.util.Map;

public class App {
    private static final Logger logger = LoggerFactory.getLogger(App.class);
    
    public static void main(String[] args) throws Exception {
        // SQL query using SELECT STREAM
        
        String sql = "SELECT STREAM\n"+
                "CEIL(event_time TO HOUR) AS event_time,\n"+
                "productId,\n"+
                "COUNT(*) AS c,\n"+
                "SUM(units) AS units\n"+
                "FROM events\n"+
                "GROUP BY\n"+
                "CEIL(event_time TO HOUR),\n"+
                "productId\n";
                

        // Create root schema and add our custom schema
        SchemaPlus rootSchema = Frameworks.createRootSchema(true);
        rootSchema.add("streaming", new StreamingSchema());
        rootSchema.add("streaming2",new StreamingSchema());

        // Create framework config using that schema
        FrameworkConfig config = Frameworks.newConfigBuilder()
            .defaultSchema(rootSchema.getSubSchema("streaming"))
            .parserConfig(
                SqlParser.config().withCaseSensitive(false))
                .build();

        // Create planner
        Planner planner = Frameworks.getPlanner(config);

        // Step 1: Parse the SQL
        SqlNode parsed = planner.parse(sql);
        System.out.println("Parsed SQL:");
        System.out.println(parsed);
            SqlNode validated = planner.validate(parsed);
            System.out.println("Validated SQL:");
         System.out.println(validated);
    
        // Step 2: Validate the SQL
        try{
            //SqlNode validated = planner.validate(parsed);
            System.out.println("Validated SQL:");
         System.out.println(validated);
    
        }catch(Exception e){
            e.printStackTrace(System.out);
        }
    }
    // ---- StreamingTable ----
    public static class StreamingTable extends AbstractTable implements StreamableTable {
        @Override
        public RelDataType getRowType(RelDataTypeFactory typeFactory) {
            return typeFactory.builder()
                .add("amount", typeFactory.createJavaType(int.class))
                .add("name", typeFactory.createJavaType(String.class))
                .add("event_time", typeFactory.createJavaType(Timestamp.class))
                .add("productId", typeFactory.createJavaType(int.class))
                .add("units", typeFactory.createJavaType(int.class))
                .build();
        }

        @Override
        public Table stream() {
            return this;
        }
    }

    // ---- Schema that includes the streaming table ----
    public static class StreamingSchema extends AbstractSchema {
        @Override
        protected Map<String, Table> getTableMap() {
            Map<String, Table> tables = new HashMap<>();
            tables.put("events", new StreamingTable());
            return tables;
        }
    }
}

报错信息

Parsed SQL:
SELECT STREAM CEIL(`EVENT_TIME` TO HOUR) AS `EVENT_TIME`, `PRODUCTID`, COUNT(*) AS `C`, SUM(`UNITS`) AS `UNITS`
FROM `EVENTS`
GROUP BY CEIL(`EVENT_TIME` TO HOUR), `PRODUCTID`
Exception in thread "main" org.apache.calcite.tools.ValidationException: org.apache.calcite.runtime.CalciteContextException: From line 7, column 1 to line 9, column 9: Streaming aggregation requires at least one monotonic expression in GROUP BY clause
        at org.apache.calcite.prepare.PlannerImpl.validate(PlannerImpl.java:228)
        at calcite.streaming.App.main(App.java:57)
Caused by: org.apache.calcite.runtime.CalciteContextException: From line 7, column 1 to line 9, column 9: Streaming aggregation requires at least one monotonic expression in GROUP BY clause
        at java.base/jdk.internal.reflect.NativeConstructorAccessorImpl.newInstance0(Native Method)
        at java.base/jdk.internal.reflect.NativeConstructorAccessorImpl.newInstance(NativeConstructorAccessorImpl.java:62)
        at java.base/jdk.internal.reflect.DelegatingConstructorAccessorImpl.newInstance(DelegatingConstructorAccessorImpl.java:45)
        at java.base/java.lang.reflect.Constructor.newInstance(Constructor.java:490)
        at org.apache.calcite.runtime.Resources$ExInstWithCause.ex(Resources.java:511)
        at org.apache.calcite.sql.SqlUtil.newContextException(SqlUtil.java:954)
        at org.apache.calcite.sql.SqlUtil.newContextException(SqlUtil.java:939)
        at org.apache.calcite.sql.validate.SqlValidatorImpl.newValidationError(SqlValidatorImpl.java:6024)
        at org.apache.calcite.sql.validate.SqlValidatorImpl.validateModality(SqlValidatorImpl.java:4534)
        at org.apache.calcite.sql.validate.SqlValidatorImpl.validateModality(SqlValidatorImpl.java:4426)
        at org.apache.calcite.sql.validate.SqlValidatorImpl.validateQuery(SqlValidatorImpl.java:1189)
        at org.apache.calcite.sql.SqlSelect.validate(SqlSelect.java:282)
        at org.apache.calcite.sql.validate.SqlValidatorImpl.validateScopedExpression(SqlValidatorImpl.java:1145)
        at org.apache.calcite.sql.validate.SqlValidatorImpl.validate(SqlValidatorImpl.java:851)
        at org.apache.calcite.prepare.PlannerImpl.validate(PlannerImpl.java:226)
        ... 1 more
Caused by: org.apache.calcite.sql.validate.SqlValidatorException: Streaming aggregation requires at least one monotonic expression in GROUP BY clause
        at java.base/jdk.internal.reflect.NativeConstructorAccessorImpl.newInstance0(Native Method)
        at java.base/jdk.internal.reflect.NativeConstructorAccessorImpl.newInstance(NativeConstructorAccessorImpl.java:62)
        at java.base/jdk.internal.reflect.DelegatingConstructorAccessorImpl.newInstance(DelegatingConstructorAccessorImpl.java:45)
        at java.base/java.lang.reflect.Constructor.newInstance(Constructor.java:490)
        at org.apache.calcite.runtime.Resources$ExInstWithCause.ex(Resources.java:511)
        at org.apache.calcite.runtime.Resources$ExInst.ex(Resources.java:605)
        ... 11 more

解决方案

Calcite要求流聚合的GROUP BY子句必须包含至少一个单调递增/递减的表达式,确保流处理可按顺序计算无需回溯。解决步骤如下:

1. 让自定义表实现MonotonicityProvider接口

Calcite的MonotonicityProvider接口用于声明字段的单调性,修改StreamingTable类实现该接口并重写对应方法:

public static class StreamingTable extends AbstractTable implements StreamableTable, MonotonicityProvider {
    @Override
    public RelDataType getRowType(RelDataTypeFactory typeFactory) {
        return typeFactory.builder()
            .add("amount", typeFactory.createJavaType(int.class))
            .add("name", typeFactory.createJavaType(String.class))
            .add("event_time", typeFactory.createJavaType(Timestamp.class))
            .add("productId", typeFactory.createJavaType(int.class))
            .add("units", typeFactory.createJavaType(int.class))
            .build();
    }

    @Override
    public Table stream() {
        return this;
    }

    @Override
    public SqlMonotonicity getMonotonicity(RelDataTypeField field) {
        // 标记event_time字段为单调递增
        if ("event_time".equalsIgnoreCase(field.getName())) {
            return SqlMonotonicity.INCREASING;
        }
        // 其他字段返回非单调
        return SqlMonotonicity.NOT_MONOTONIC;
    }
}

注意:需重写的是getMonotonicity(RelDataTypeField field)方法,而非字符串参数版本,这是MonotonicityProvider接口的标准方法。

2. 确认时间表达式的单调性继承

GROUP BY中的CEIL(event_time TO HOUR)表达式,Calcite会自动基于event_time的单调性推断其也是单调递增的,无需额外配置,默认规则即可覆盖这种场景。

3. 验证修改

重新运行代码,Calcite会识别event_time的单调性,进而确认CEIL(event_time TO HOUR)符合流聚合的要求,验证阶段即可通过。

额外提示

  • 流处理中时间字段的单调性是窗口聚合正确性的核心,Calcite通过该检查避免无法高效处理的聚合逻辑。
  • 若数据流存在乱序情况,后续可考虑引入水印(Watermark)机制,但当前需先解决字段单调性标记的基础问题。

内容的提问来源于stack exchange,提问作者Mukund

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 02:28:11