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
相关产品推荐
相关产品推荐

