From 1269db973099b3860978227d462707eabcfc64f4 Mon Sep 17 00:00:00 2001 From: Ferenc Csaky Date: Mon, 27 Jul 2026 14:14:43 +0200 Subject: [PATCH 1/2] refactor: Rationalize SQL dialect and statement code, to simplify adding new dialects --- .../datasqrl/calcite/SqrlConfigurations.java | 2 +- .../convert/AbstractSqlConverters.java | 75 ++++++++++++++++ .../calcite/convert/DuckDbRelToSqlNode.java | 41 --------- ...ToString.java => DuckDbSqlConverters.java} | 14 ++- .../convert/DuckdbSqlNodeToString.java | 43 ---------- .../convert/FlinkSqlConverters.java} | 18 ++-- .../calcite/convert/PostgresRelToSqlNode.java | 38 --------- ...qlNode.java => PostgresSqlConverters.java} | 17 ++-- .../convert/PostgresSqlNodeToString.java | 43 ---------- ...tring.java => SnowflakeSqlConverters.java} | 21 ++--- ...keRelToSqlNode.java => SqlConverters.java} | 37 ++++---- ...Factory.java => SqlConvertersFactory.java} | 9 +- .../dialect/BasePostgresSqlDialect.java | 10 +-- .../calcite/dialect/DuckDbSqlDialect.java | 31 ++----- .../dialect/ExtendedPostgresSqlDialect.java | 35 +++----- .../dialect/ExtendedSnowflakeSqlDialect.java | 28 ++---- .../dialect/SqlTranslationDispatcher.java | 62 ++++++++++++++ .../{ => postgres}/PostgresConformance.java | 2 +- .../config/JdbcEngineConfigDelegate.java | 85 ------------------- .../AbstractJdbcStatementFactory.java | 35 +++++--- .../relational/DuckDbStatementFactory.java | 7 +- .../relational/IcebergStatementFactory.java | 1 - .../relational/PostgresStatementFactory.java | 12 +-- .../relational/SnowflakeStatementFactory.java | 11 +-- .../relational/ddl/InsertStatement.java | 12 ++- .../stream/flink/sql/RelToFlinkSql.java | 6 -- .../stream/flink/plan/FlinkSqlNodesTest.java | 11 +-- 27 files changed, 263 insertions(+), 443 deletions(-) create mode 100644 sqrl-planner/src/main/java/com/datasqrl/calcite/convert/AbstractSqlConverters.java delete mode 100644 sqrl-planner/src/main/java/com/datasqrl/calcite/convert/DuckDbRelToSqlNode.java rename sqrl-planner/src/main/java/com/datasqrl/calcite/convert/{SqlNodeToString.java => DuckDbSqlConverters.java} (70%) delete mode 100644 sqrl-planner/src/main/java/com/datasqrl/calcite/convert/DuckdbSqlNodeToString.java rename sqrl-planner/src/main/java/com/datasqrl/{util/FlinkSqlNodeToString.java => calcite/convert/FlinkSqlConverters.java} (69%) delete mode 100644 sqrl-planner/src/main/java/com/datasqrl/calcite/convert/PostgresRelToSqlNode.java rename sqrl-planner/src/main/java/com/datasqrl/calcite/convert/{RelToSqlNode.java => PostgresSqlConverters.java} (67%) delete mode 100644 sqrl-planner/src/main/java/com/datasqrl/calcite/convert/PostgresSqlNodeToString.java rename sqrl-planner/src/main/java/com/datasqrl/calcite/convert/{SnowflakeSqlNodeToString.java => SnowflakeSqlConverters.java} (66%) rename sqrl-planner/src/main/java/com/datasqrl/calcite/convert/{SnowflakeRelToSqlNode.java => SqlConverters.java} (51%) rename sqrl-planner/src/main/java/com/datasqrl/calcite/convert/{SqlToStringFactory.java => SqlConvertersFactory.java} (72%) create mode 100644 sqrl-planner/src/main/java/com/datasqrl/calcite/dialect/SqlTranslationDispatcher.java rename sqrl-planner/src/main/java/com/datasqrl/calcite/dialect/{ => postgres}/PostgresConformance.java (99%) delete mode 100644 sqrl-planner/src/main/java/com/datasqrl/config/JdbcEngineConfigDelegate.java diff --git a/sqrl-planner/src/main/java/com/datasqrl/calcite/SqrlConfigurations.java b/sqrl-planner/src/main/java/com/datasqrl/calcite/SqrlConfigurations.java index 56e8d58791..77cb9b6b13 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/calcite/SqrlConfigurations.java +++ b/sqrl-planner/src/main/java/com/datasqrl/calcite/SqrlConfigurations.java @@ -20,7 +20,7 @@ public class SqrlConfigurations { - public static final UnaryOperator sqlToString = + public static final UnaryOperator SQL_TO_STRING = c -> c.withAlwaysUseParentheses(false) .withSelectListItemsOnSeparateLines(false) diff --git a/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/AbstractSqlConverters.java b/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/AbstractSqlConverters.java new file mode 100644 index 0000000000..dfe40322c9 --- /dev/null +++ b/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/AbstractSqlConverters.java @@ -0,0 +1,75 @@ +/* + * Copyright © 2021 DataSQRL (contact@datasqrl.com) + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package com.datasqrl.calcite.convert; + +import com.datasqrl.calcite.Dialect; +import com.datasqrl.calcite.DynamicParamSqlPrettyWriter; +import com.datasqrl.calcite.SqrlConfigurations; +import java.util.Map; +import lombok.AccessLevel; +import lombok.RequiredArgsConstructor; +import org.apache.calcite.rel.RelNode; +import org.apache.calcite.rel.rel2sql.RelToSqlConverterWithHints; +import org.apache.calcite.sql.CalciteFixes; +import org.apache.calcite.sql.SqlDialect; +import org.apache.calcite.sql.SqlNode; +import org.apache.calcite.sql.pretty.SqlPrettyWriter; + +@RequiredArgsConstructor(access = AccessLevel.PACKAGE) +abstract class AbstractSqlConverters implements SqlConverters { + + private final Dialect dialect; + private final SqlDialect calciteSqlDialect; + private final boolean appendSelectLists; + + @Override + public final SqlNode convert(RelNode relNode, Map tableNameMapping) { + var sqlNode = + new RelToSqlConverterWithHints(calciteSqlDialect, tableNameMapping) + .visitRoot(relNode) + .asStatement(); + + if (appendSelectLists) { + CalciteFixes.appendSelectLists(sqlNode); + } + + return sqlNode; + } + + @Override + public final String convert(SqlNode sqlNode) { + var writer = createWriter(); + sqlNode.unparse(writer, 0, 0); + + return writer.toSqlString().getSql(); + } + + @Override + public final Dialect getDialect() { + return dialect; + } + + protected SqlPrettyWriter createWriter() { + var baseConfig = SqlPrettyWriter.config().withDialect(calciteSqlDialect); + var config = SqrlConfigurations.SQL_TO_STRING.apply(baseConfig); + + return new DynamicParamSqlPrettyWriter(config); + } + + protected final SqlDialect getCalciteSqlDialect() { + return calciteSqlDialect; + } +} diff --git a/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/DuckDbRelToSqlNode.java b/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/DuckDbRelToSqlNode.java deleted file mode 100644 index d2d357cadf..0000000000 --- a/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/DuckDbRelToSqlNode.java +++ /dev/null @@ -1,41 +0,0 @@ -/* - * Copyright © 2021 DataSQRL (contact@datasqrl.com) - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package com.datasqrl.calcite.convert; - -import com.datasqrl.calcite.Dialect; -import com.datasqrl.calcite.dialect.DuckDbSqlDialect; -import java.util.Map; -import org.apache.calcite.rel.RelNode; -import org.apache.calcite.rel.rel2sql.RelToSqlConverterWithHints; -import org.apache.calcite.sql.CalciteFixes; - -public class DuckDbRelToSqlNode implements RelToSqlNode { - - @Override - public SqlNodes convert(RelNode relNode, Map tableNameMapping) { - var node = - new RelToSqlConverterWithHints(DuckDbSqlDialect.DEFAULT, tableNameMapping) - .visitRoot(relNode) - .asStatement(); - CalciteFixes.appendSelectLists(node); - return () -> node; - } - - @Override - public Dialect getDialect() { - return Dialect.DUCKDB; - } -} diff --git a/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/SqlNodeToString.java b/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/DuckDbSqlConverters.java similarity index 70% rename from sqrl-planner/src/main/java/com/datasqrl/calcite/convert/SqlNodeToString.java rename to sqrl-planner/src/main/java/com/datasqrl/calcite/convert/DuckDbSqlConverters.java index f04df5cb57..46c93c8951 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/SqlNodeToString.java +++ b/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/DuckDbSqlConverters.java @@ -16,15 +16,13 @@ package com.datasqrl.calcite.convert; import com.datasqrl.calcite.Dialect; -import com.datasqrl.calcite.convert.RelToSqlNode.SqlNodes; +import com.datasqrl.calcite.dialect.DuckDbSqlDialect; +import com.google.auto.service.AutoService; -/** Converts a SqlNode to a SQL string for a given dialect */ -public interface SqlNodeToString { - SqlStrings convert(SqlNodes sqlNode); +@AutoService(SqlConverters.class) +public class DuckDbSqlConverters extends AbstractSqlConverters { - Dialect getDialect(); - - interface SqlStrings { - String getSql(); + public DuckDbSqlConverters() { + super(Dialect.DUCKDB, DuckDbSqlDialect.DEFAULT, true); } } diff --git a/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/DuckdbSqlNodeToString.java b/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/DuckdbSqlNodeToString.java deleted file mode 100644 index 4f2fa4cd2e..0000000000 --- a/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/DuckdbSqlNodeToString.java +++ /dev/null @@ -1,43 +0,0 @@ -/* - * Copyright © 2021 DataSQRL (contact@datasqrl.com) - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package com.datasqrl.calcite.convert; - -import com.datasqrl.calcite.Dialect; -import com.datasqrl.calcite.DynamicParamSqlPrettyWriter; -import com.datasqrl.calcite.SqrlConfigurations; -import com.datasqrl.calcite.convert.RelToSqlNode.SqlNodes; -import com.datasqrl.calcite.dialect.DuckDbSqlDialect; -import com.google.auto.service.AutoService; -import org.apache.calcite.sql.pretty.SqlPrettyWriter; - -@AutoService(SqlNodeToString.class) -public class DuckdbSqlNodeToString implements SqlNodeToString { - - @Override - public SqlStrings convert(SqlNodes sqlNode) { - var config = - SqrlConfigurations.sqlToString.apply( - SqlPrettyWriter.config().withDialect(DuckDbSqlDialect.DEFAULT)); - var writer = new DynamicParamSqlPrettyWriter(config); - sqlNode.getSqlNode().unparse(writer, 0, 0); - return () -> writer.toSqlString().getSql(); - } - - @Override - public Dialect getDialect() { - return Dialect.DUCKDB; - } -} diff --git a/sqrl-planner/src/main/java/com/datasqrl/util/FlinkSqlNodeToString.java b/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/FlinkSqlConverters.java similarity index 69% rename from sqrl-planner/src/main/java/com/datasqrl/util/FlinkSqlNodeToString.java rename to sqrl-planner/src/main/java/com/datasqrl/calcite/convert/FlinkSqlConverters.java index 6741f68534..ffc750d193 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/util/FlinkSqlNodeToString.java +++ b/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/FlinkSqlConverters.java @@ -13,19 +13,25 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package com.datasqrl.util; +package com.datasqrl.calcite.convert; import com.datasqrl.calcite.Dialect; -import com.datasqrl.calcite.convert.RelToSqlNode.SqlNodes; -import com.datasqrl.calcite.convert.SqlNodeToString; import com.datasqrl.engine.stream.flink.sql.RelToFlinkSql; import com.google.auto.service.AutoService; +import java.util.Map; +import org.apache.calcite.rel.RelNode; +import org.apache.calcite.sql.SqlNode; -@AutoService(SqlNodeToString.class) -public class FlinkSqlNodeToString implements SqlNodeToString { +@AutoService(SqlConverters.class) +public class FlinkSqlConverters implements SqlConverters { @Override - public SqlStrings convert(SqlNodes sqlNode) { + public SqlNode convert(RelNode relNode, Map tableNameMapping) { + throw new UnsupportedOperationException(); + } + + @Override + public String convert(SqlNode sqlNode) { // TODO: Migrate remaining FlinkRelToSqlConverter to this paradigm return RelToFlinkSql.convertToString(sqlNode); } diff --git a/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/PostgresRelToSqlNode.java b/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/PostgresRelToSqlNode.java deleted file mode 100644 index 07b2747bb0..0000000000 --- a/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/PostgresRelToSqlNode.java +++ /dev/null @@ -1,38 +0,0 @@ -/* - * Copyright © 2021 DataSQRL (contact@datasqrl.com) - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package com.datasqrl.calcite.convert; - -import com.datasqrl.calcite.Dialect; -import com.datasqrl.calcite.dialect.ExtendedPostgresSqlDialect; -import java.util.Map; -import org.apache.calcite.rel.RelNode; -import org.apache.calcite.rel.rel2sql.RelToSqlConverter; -import org.apache.calcite.rel.rel2sql.RelToSqlConverterWithHints; - -public class PostgresRelToSqlNode implements RelToSqlNode { - - @Override - public SqlNodes convert(RelNode relNode, Map tableNameMapping) { - RelToSqlConverter converter = - new RelToSqlConverterWithHints(ExtendedPostgresSqlDialect.DEFAULT, tableNameMapping); - return () -> converter.visitRoot(relNode).asStatement(); - } - - @Override - public Dialect getDialect() { - return Dialect.POSTGRES; - } -} diff --git a/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/RelToSqlNode.java b/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/PostgresSqlConverters.java similarity index 67% rename from sqrl-planner/src/main/java/com/datasqrl/calcite/convert/RelToSqlNode.java rename to sqrl-planner/src/main/java/com/datasqrl/calcite/convert/PostgresSqlConverters.java index bf3f7b951d..f885d67eb6 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/RelToSqlNode.java +++ b/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/PostgresSqlConverters.java @@ -16,18 +16,13 @@ package com.datasqrl.calcite.convert; import com.datasqrl.calcite.Dialect; -import java.util.Map; -import org.apache.calcite.rel.RelNode; -import org.apache.calcite.sql.SqlNode; +import com.datasqrl.calcite.dialect.ExtendedPostgresSqlDialect; +import com.google.auto.service.AutoService; -/** Converts a RelNode to a SqlNode for a given dialect */ -public interface RelToSqlNode { +@AutoService(SqlConverters.class) +public class PostgresSqlConverters extends AbstractSqlConverters { - SqlNodes convert(RelNode relNode, Map tableNameMapping); - - Dialect getDialect(); - - interface SqlNodes { - SqlNode getSqlNode(); + public PostgresSqlConverters() { + super(Dialect.POSTGRES, ExtendedPostgresSqlDialect.DEFAULT, false); } } diff --git a/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/PostgresSqlNodeToString.java b/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/PostgresSqlNodeToString.java deleted file mode 100644 index 91a95c552f..0000000000 --- a/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/PostgresSqlNodeToString.java +++ /dev/null @@ -1,43 +0,0 @@ -/* - * Copyright © 2021 DataSQRL (contact@datasqrl.com) - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package com.datasqrl.calcite.convert; - -import com.datasqrl.calcite.Dialect; -import com.datasqrl.calcite.DynamicParamSqlPrettyWriter; -import com.datasqrl.calcite.SqrlConfigurations; -import com.datasqrl.calcite.convert.RelToSqlNode.SqlNodes; -import com.datasqrl.calcite.dialect.ExtendedPostgresSqlDialect; -import com.google.auto.service.AutoService; -import org.apache.calcite.sql.pretty.SqlPrettyWriter; - -@AutoService(SqlNodeToString.class) -public class PostgresSqlNodeToString implements SqlNodeToString { - - @Override - public SqlStrings convert(SqlNodes sqlNode) { - var config = - SqrlConfigurations.sqlToString.apply( - SqlPrettyWriter.config().withDialect(ExtendedPostgresSqlDialect.DEFAULT)); - var writer = new DynamicParamSqlPrettyWriter(config); - sqlNode.getSqlNode().unparse(writer, 0, 0); - return () -> writer.toSqlString().getSql(); - } - - @Override - public Dialect getDialect() { - return Dialect.POSTGRES; - } -} diff --git a/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/SnowflakeSqlNodeToString.java b/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/SnowflakeSqlConverters.java similarity index 66% rename from sqrl-planner/src/main/java/com/datasqrl/calcite/convert/SnowflakeSqlNodeToString.java rename to sqrl-planner/src/main/java/com/datasqrl/calcite/convert/SnowflakeSqlConverters.java index 4c22f349f5..4fef7542f1 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/SnowflakeSqlNodeToString.java +++ b/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/SnowflakeSqlConverters.java @@ -16,28 +16,25 @@ package com.datasqrl.calcite.convert; import com.datasqrl.calcite.Dialect; -import com.datasqrl.calcite.convert.RelToSqlNode.SqlNodes; import com.datasqrl.calcite.dialect.ExtendedSnowflakeSqlDialect; import com.google.auto.service.AutoService; import org.apache.calcite.sql.pretty.SqlPrettyWriter; -@AutoService(SqlNodeToString.class) -public class SnowflakeSqlNodeToString implements SqlNodeToString { +@AutoService(SqlConverters.class) +public class SnowflakeSqlConverters extends AbstractSqlConverters { + + public SnowflakeSqlConverters() { + super(Dialect.SNOWFLAKE, ExtendedSnowflakeSqlDialect.DEFAULT, true); + } @Override - public SqlStrings convert(SqlNodes sqlNode) { + protected SqlPrettyWriter createWriter() { var config = SqlPrettyWriter.config() - .withDialect(ExtendedSnowflakeSqlDialect.DEFAULT) + .withDialect(getCalciteSqlDialect()) .withQuoteAllIdentifiers(false) .withIndentation(0); - var prettyWriter = new SqlPrettyWriter(config); - sqlNode.getSqlNode().unparse(prettyWriter, 0, 0); - return () -> prettyWriter.toSqlString().getSql(); - } - @Override - public Dialect getDialect() { - return Dialect.SNOWFLAKE; + return new SqlPrettyWriter(config); } } diff --git a/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/SnowflakeRelToSqlNode.java b/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/SqlConverters.java similarity index 51% rename from sqrl-planner/src/main/java/com/datasqrl/calcite/convert/SnowflakeRelToSqlNode.java rename to sqrl-planner/src/main/java/com/datasqrl/calcite/convert/SqlConverters.java index e60f410b29..4ce3f60c5a 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/SnowflakeRelToSqlNode.java +++ b/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/SqlConverters.java @@ -16,26 +16,29 @@ package com.datasqrl.calcite.convert; import com.datasqrl.calcite.Dialect; -import com.datasqrl.calcite.dialect.ExtendedSnowflakeSqlDialect; import java.util.Map; import org.apache.calcite.rel.RelNode; -import org.apache.calcite.rel.rel2sql.RelToSqlConverterWithHints; -import org.apache.calcite.sql.CalciteFixes; +import org.apache.calcite.sql.SqlNode; -public class SnowflakeRelToSqlNode implements RelToSqlNode { +/** Provides conversions to SQL representations for a specific dialect. */ +public interface SqlConverters { - @Override - public SqlNodes convert(RelNode relNode, Map tableNameMapping) { - var node = - new RelToSqlConverterWithHints(ExtendedSnowflakeSqlDialect.DEFAULT, tableNameMapping) - .visitRoot(relNode) - .asStatement(); - CalciteFixes.appendSelectLists(node); - return () -> node; - } + /** + * Converts a relational plan to a SQL node for this converter's dialect. + * + * @param relNode relational plan to convert + * @param tableNameMapping mapping from planner table identifiers to physical table names + * @return the dialect-specific SQL node + */ + SqlNode convert(RelNode relNode, Map tableNameMapping); - @Override - public Dialect getDialect() { - return Dialect.SNOWFLAKE; - } + /** + * Serializes a SQL node as SQL for this converter's dialect. + * + * @param sqlNode SQL node to unparse + * @return the dialect-specific SQL string + */ + String convert(SqlNode sqlNode); + + Dialect getDialect(); } diff --git a/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/SqlToStringFactory.java b/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/SqlConvertersFactory.java similarity index 72% rename from sqrl-planner/src/main/java/com/datasqrl/calcite/convert/SqlToStringFactory.java rename to sqrl-planner/src/main/java/com/datasqrl/calcite/convert/SqlConvertersFactory.java index d8c2792dde..2b4d789294 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/SqlToStringFactory.java +++ b/sqrl-planner/src/main/java/com/datasqrl/calcite/convert/SqlConvertersFactory.java @@ -17,11 +17,14 @@ import com.datasqrl.calcite.Dialect; import com.datasqrl.util.ServiceLoaderDiscovery; +import lombok.AccessLevel; +import lombok.NoArgsConstructor; -public class SqlToStringFactory { +@NoArgsConstructor(access = AccessLevel.PRIVATE) +public final class SqlConvertersFactory { - public static SqlNodeToString get(Dialect dialect) { + public static SqlConverters get(Dialect dialect) { return ServiceLoaderDiscovery.get( - SqlNodeToString.class, e -> e.getDialect().name(), dialect.name()); + SqlConverters.class, converter -> converter.getDialect().name(), dialect.name()); } } diff --git a/sqrl-planner/src/main/java/com/datasqrl/calcite/dialect/BasePostgresSqlDialect.java b/sqrl-planner/src/main/java/com/datasqrl/calcite/dialect/BasePostgresSqlDialect.java index 3c10d8b9ee..aeab03f22b 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/calcite/dialect/BasePostgresSqlDialect.java +++ b/sqrl-planner/src/main/java/com/datasqrl/calcite/dialect/BasePostgresSqlDialect.java @@ -17,8 +17,7 @@ import static org.apache.calcite.sql.SqlKind.COLLECTION_TABLE; -import com.datasqrl.function.translation.SqlTranslation; -import java.util.Map; +import com.datasqrl.calcite.Dialect; import org.apache.calcite.sql.SqlCall; import org.apache.calcite.sql.SqlWriter; import org.apache.calcite.sql.dialect.PostgresqlSqlDialect; @@ -29,7 +28,7 @@ public BasePostgresSqlDialect(Context context) { super(context); } - protected abstract Map getTranslationMap(); + protected abstract Dialect getTranslationDialect(); @Override public void unparseCall(SqlWriter writer, SqlCall call, int leftPrec, int rightPrec) { @@ -38,9 +37,8 @@ public void unparseCall(SqlWriter writer, SqlCall call, int leftPrec, int rightP return; } - var operatorName = call.getOperator().getName().toLowerCase(); - if (getTranslationMap().containsKey(operatorName)) { - getTranslationMap().get(operatorName).unparse(call, writer, leftPrec, rightPrec); + if (SqlTranslationDispatcher.tryUnparseTranslatedCall( + getTranslationDialect(), call, writer, leftPrec, rightPrec)) { return; } try { diff --git a/sqrl-planner/src/main/java/com/datasqrl/calcite/dialect/DuckDbSqlDialect.java b/sqrl-planner/src/main/java/com/datasqrl/calcite/dialect/DuckDbSqlDialect.java index 045bfebad6..ab5a6097a1 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/calcite/dialect/DuckDbSqlDialect.java +++ b/sqrl-planner/src/main/java/com/datasqrl/calcite/dialect/DuckDbSqlDialect.java @@ -16,10 +16,6 @@ package com.datasqrl.calcite.dialect; import com.datasqrl.calcite.Dialect; -import com.datasqrl.function.translation.SqlTranslation; -import com.datasqrl.util.ServiceLoaderDiscovery; -import java.util.Map; -import java.util.stream.Collectors; import org.apache.calcite.avatica.util.Casing; import org.apache.calcite.rel.type.RelDataType; import org.apache.calcite.sql.SqlAlienSystemTypeNameSpec; @@ -29,24 +25,13 @@ public class DuckDbSqlDialect extends BasePostgresSqlDialect { - public static final SqlDialect.Context DEFAULT_CONTEXT; - public static final SqlDialect DEFAULT; + public static final SqlDialect.Context DEFAULT_CONTEXT = + SqlDialect.EMPTY_CONTEXT + .withDatabaseProduct(DatabaseProduct.POSTGRESQL) + .withIdentifierQuoteString("\"") + .withUnquotedCasing(Casing.TO_LOWER); - private static final Map TRANSLATION_MAP; - - static { - DEFAULT_CONTEXT = - SqlDialect.EMPTY_CONTEXT - .withDatabaseProduct(DatabaseProduct.POSTGRESQL) - .withIdentifierQuoteString("\"") - .withUnquotedCasing(Casing.TO_LOWER); - DEFAULT = new DuckDbSqlDialect(DEFAULT_CONTEXT); - - TRANSLATION_MAP = - ServiceLoaderDiscovery.getAll(SqlTranslation.class).stream() - .filter(f -> f.getDialect() == Dialect.DUCKDB) - .collect(Collectors.toMap(f -> f.getOperator().getName().toLowerCase(), f -> f)); - } + public static final SqlDialect DEFAULT = new DuckDbSqlDialect(DEFAULT_CONTEXT); public DuckDbSqlDialect(Context context) { super(context); @@ -69,7 +54,7 @@ public SqlDataTypeSpec getCastSpec(RelDataType type) { } @Override - protected Map getTranslationMap() { - return TRANSLATION_MAP; + protected Dialect getTranslationDialect() { + return Dialect.DUCKDB; } } diff --git a/sqrl-planner/src/main/java/com/datasqrl/calcite/dialect/ExtendedPostgresSqlDialect.java b/sqrl-planner/src/main/java/com/datasqrl/calcite/dialect/ExtendedPostgresSqlDialect.java index 141b39c9a2..750a74c073 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/calcite/dialect/ExtendedPostgresSqlDialect.java +++ b/sqrl-planner/src/main/java/com/datasqrl/calcite/dialect/ExtendedPostgresSqlDialect.java @@ -16,12 +16,9 @@ package com.datasqrl.calcite.dialect; import com.datasqrl.calcite.Dialect; +import com.datasqrl.calcite.dialect.postgres.PostgresConformance; import com.datasqrl.flinkrunner.stdlib.json.FlinkJsonType; import com.datasqrl.flinkrunner.stdlib.vector.FlinkVectorType; -import com.datasqrl.function.translation.SqlTranslation; -import com.datasqrl.util.ServiceLoaderDiscovery; -import java.util.Map; -import java.util.stream.Collectors; import org.apache.calcite.avatica.util.Casing; import org.apache.calcite.rel.type.RelDataType; import org.apache.calcite.rel.type.RelRecordType; @@ -36,24 +33,16 @@ public class ExtendedPostgresSqlDialect extends BasePostgresSqlDialect { - public static final Map translationMap = - ServiceLoaderDiscovery.getAll(SqlTranslation.class).stream() - .filter(f -> f.getDialect() == Dialect.POSTGRES) - .collect(Collectors.toMap(f -> f.getOperator().getName().toLowerCase(), f -> f)); + public static final ExtendedPostgresSqlDialect.Context DEFAULT_CONTEXT = + SqlDialect.EMPTY_CONTEXT + .withDatabaseProduct(DatabaseProduct.POSTGRESQL) + .withIdentifierQuoteString("\"") + .withUnquotedCasing(Casing.TO_LOWER) + .withDataTypeSystem(POSTGRESQL_TYPE_SYSTEM) + .withConformance(new PostgresConformance()); - public static final ExtendedPostgresSqlDialect.Context DEFAULT_CONTEXT; - public static final ExtendedPostgresSqlDialect DEFAULT; - - static { - DEFAULT_CONTEXT = - SqlDialect.EMPTY_CONTEXT - .withDatabaseProduct(DatabaseProduct.POSTGRESQL) - .withIdentifierQuoteString("\"") - .withUnquotedCasing(Casing.TO_LOWER) - .withDataTypeSystem(POSTGRESQL_TYPE_SYSTEM) - .withConformance(new PostgresConformance()); - DEFAULT = new ExtendedPostgresSqlDialect(DEFAULT_CONTEXT); - } + public static final ExtendedPostgresSqlDialect DEFAULT = + new ExtendedPostgresSqlDialect(DEFAULT_CONTEXT); public ExtendedPostgresSqlDialect(Context context) { super(context); @@ -146,7 +135,7 @@ public boolean supportsGroupByLiteral() { } @Override - protected Map getTranslationMap() { - return translationMap; + protected Dialect getTranslationDialect() { + return Dialect.POSTGRES; } } diff --git a/sqrl-planner/src/main/java/com/datasqrl/calcite/dialect/ExtendedSnowflakeSqlDialect.java b/sqrl-planner/src/main/java/com/datasqrl/calcite/dialect/ExtendedSnowflakeSqlDialect.java index e8b9af3433..136e77d488 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/calcite/dialect/ExtendedSnowflakeSqlDialect.java +++ b/sqrl-planner/src/main/java/com/datasqrl/calcite/dialect/ExtendedSnowflakeSqlDialect.java @@ -16,31 +16,19 @@ package com.datasqrl.calcite.dialect; import com.datasqrl.calcite.Dialect; -import com.datasqrl.function.translation.SqlTranslation; -import com.datasqrl.util.ServiceLoaderDiscovery; -import java.util.Map; -import java.util.stream.Collectors; import org.apache.calcite.sql.SqlCall; import org.apache.calcite.sql.SqlDialect; import org.apache.calcite.sql.SqlWriter; import org.apache.calcite.sql.dialect.SnowflakeSqlDialect; public class ExtendedSnowflakeSqlDialect extends SnowflakeSqlDialect { - public static final Map translationMap = - ServiceLoaderDiscovery.getAll(SqlTranslation.class).stream() - .filter(f -> f.getDialect() == Dialect.SNOWFLAKE) - .collect(Collectors.toMap(f -> f.getOperator().getName().toLowerCase(), f -> f)); - public static final SqlDialect.Context DEFAULT_CONTEXT; - public static final SqlDialect DEFAULT; + public static final SqlDialect.Context DEFAULT_CONTEXT = + SqlDialect.EMPTY_CONTEXT + .withDatabaseProduct(DatabaseProduct.SNOWFLAKE) + .withIdentifierQuoteString("\""); - static { - DEFAULT_CONTEXT = SqlDialect.EMPTY_CONTEXT.withDatabaseProduct(DatabaseProduct.SNOWFLAKE) - // .withIdentifierQuoteString("\"") - // .withUnquotedCasing(Casing.TO_UPPER) - ; - DEFAULT = new ExtendedSnowflakeSqlDialect(DEFAULT_CONTEXT); - } + public static final SqlDialect DEFAULT = new ExtendedSnowflakeSqlDialect(DEFAULT_CONTEXT); public ExtendedSnowflakeSqlDialect(Context context) { super(context); @@ -48,10 +36,8 @@ public ExtendedSnowflakeSqlDialect(Context context) { @Override public void unparseCall(SqlWriter writer, SqlCall call, int leftPrec, int rightPrec) { - if (translationMap.containsKey(call.getOperator().getName().toLowerCase())) { - translationMap - .get(call.getOperator().getName().toLowerCase()) - .unparse(call, writer, leftPrec, rightPrec); + if (SqlTranslationDispatcher.tryUnparseTranslatedCall( + Dialect.SNOWFLAKE, call, writer, leftPrec, rightPrec)) { return; } diff --git a/sqrl-planner/src/main/java/com/datasqrl/calcite/dialect/SqlTranslationDispatcher.java b/sqrl-planner/src/main/java/com/datasqrl/calcite/dialect/SqlTranslationDispatcher.java new file mode 100644 index 0000000000..e38f850560 --- /dev/null +++ b/sqrl-planner/src/main/java/com/datasqrl/calcite/dialect/SqlTranslationDispatcher.java @@ -0,0 +1,62 @@ +/* + * Copyright © 2021 DataSQRL (contact@datasqrl.com) + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package com.datasqrl.calcite.dialect; + +import com.datasqrl.calcite.Dialect; +import com.datasqrl.function.translation.SqlTranslation; +import com.datasqrl.util.ServiceLoaderDiscovery; +import java.util.EnumMap; +import java.util.Map; +import java.util.stream.Collectors; +import lombok.AccessLevel; +import lombok.NoArgsConstructor; +import org.apache.calcite.sql.SqlCall; +import org.apache.calcite.sql.SqlWriter; + +@NoArgsConstructor(access = AccessLevel.PRIVATE) +final class SqlTranslationDispatcher { + + private static final Map> DIALECT_TRANSLATIONS = + Map.copyOf( + ServiceLoaderDiscovery.getAll(SqlTranslation.class).stream() + .collect( + Collectors.groupingBy( + SqlTranslation::getDialect, + () -> new EnumMap<>(Dialect.class), + Collectors.collectingAndThen( + Collectors.toMap( + translation -> operatorKey(translation.getOperator().getName()), + translation -> translation), + Map::copyOf)))); + + static boolean tryUnparseTranslatedCall( + Dialect dialect, SqlCall call, SqlWriter writer, int leftPrec, int rightPrec) { + + var dialectTranslations = DIALECT_TRANSLATIONS.getOrDefault(dialect, Map.of()); + var translation = dialectTranslations.get(operatorKey(call.getOperator().getName())); + if (translation == null) { + return false; + } + + translation.unparse(call, writer, leftPrec, rightPrec); + + return true; + } + + private static String operatorKey(String operatorName) { + return operatorName.toLowerCase(); + } +} diff --git a/sqrl-planner/src/main/java/com/datasqrl/calcite/dialect/PostgresConformance.java b/sqrl-planner/src/main/java/com/datasqrl/calcite/dialect/postgres/PostgresConformance.java similarity index 99% rename from sqrl-planner/src/main/java/com/datasqrl/calcite/dialect/PostgresConformance.java rename to sqrl-planner/src/main/java/com/datasqrl/calcite/dialect/postgres/PostgresConformance.java index e86bc18e5b..d428fa3651 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/calcite/dialect/PostgresConformance.java +++ b/sqrl-planner/src/main/java/com/datasqrl/calcite/dialect/postgres/PostgresConformance.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package com.datasqrl.calcite.dialect; +package com.datasqrl.calcite.dialect.postgres; import org.apache.calcite.sql.fun.SqlLibrary; import org.apache.calcite.sql.validate.SqlConformance; diff --git a/sqrl-planner/src/main/java/com/datasqrl/config/JdbcEngineConfigDelegate.java b/sqrl-planner/src/main/java/com/datasqrl/config/JdbcEngineConfigDelegate.java deleted file mode 100644 index 29a9984293..0000000000 --- a/sqrl-planner/src/main/java/com/datasqrl/config/JdbcEngineConfigDelegate.java +++ /dev/null @@ -1,85 +0,0 @@ -/* - * Copyright © 2021 DataSQRL (contact@datasqrl.com) - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package com.datasqrl.config; - -import java.util.Map; -import java.util.regex.Pattern; - -/** Temporary wrapper around jdbc engine config until it can move to a template */ -public class JdbcEngineConfigDelegate { - public static final Pattern JDBC_URL_REGEX = - Pattern.compile("^jdbc:(.*?):\\/\\/([^/:]*)(?::(\\d+))?\\/([^/:?]*)(.*)$"); - public static final Pattern JDBC_DIALECT_REGEX = Pattern.compile("^jdbc:(.*?):(.*)$"); - - private final Map map; - private final String dialect; - private String host; - private int port; - private String database; - private final String url; - - public JdbcEngineConfigDelegate(ConnectorConf connectorConf) { - this.map = connectorConf.toMap(); - this.url = map.get("url"); - var matcher = JDBC_URL_REGEX.matcher(url); - if (matcher.find()) { - var dialect = matcher.group(1); - // connectorConfig.getErrorCollector().checkFatal(JdbcDialect.find(dialect).isPresent(), - // "Invalid database dialect: %s", dialect); - this.dialect = dialect; - this.host = matcher.group(2); - this.port = Integer.parseInt(matcher.group(3)); - this.database = matcher.group(4); - } else { - // Only extract the dialect - matcher = JDBC_DIALECT_REGEX.matcher(url); - if (matcher.find()) { - var dialect = matcher.group(1); - this.dialect = dialect; - } else { - throw new RuntimeException("Invalid database URL: %s".formatted(url)); - } - } - } - - public JdbcDialect getDialect() { - return JdbcDialect.find(dialect.toLowerCase()).get(); - } - - public String getHost() { - return this.host; - } - - public Integer getPort() { - return this.port; - } - - public String getUser() { - return (String) map.get("username"); - } - - public String getPassword() { - return (String) map.get("password"); - } - - public String getDatabase() { - return this.database; - } - - public String getUrl() { - return this.url; - } -} diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/AbstractJdbcStatementFactory.java b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/AbstractJdbcStatementFactory.java index cdcd374abd..7479956196 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/AbstractJdbcStatementFactory.java +++ b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/AbstractJdbcStatementFactory.java @@ -17,10 +17,10 @@ import static com.google.common.base.Preconditions.checkArgument; +import com.datasqrl.calcite.Dialect; import com.datasqrl.calcite.OperatorRuleTransformer; -import com.datasqrl.calcite.convert.RelToSqlNode; -import com.datasqrl.calcite.convert.RelToSqlNode.SqlNodes; -import com.datasqrl.calcite.convert.SqlNodeToString; +import com.datasqrl.calcite.convert.SqlConverters; +import com.datasqrl.calcite.convert.SqlConvertersFactory; import com.datasqrl.calcite.dialect.postgres.SqlCreatePostgresView; import com.datasqrl.canonicalizer.Name; import com.datasqrl.engine.database.relational.CreateTableJdbcStatement.CreateTableDdlFactory; @@ -49,7 +49,8 @@ import java.util.regex.Pattern; import java.util.stream.Collectors; import java.util.stream.Stream; -import lombok.AllArgsConstructor; +import lombok.AccessLevel; +import lombok.RequiredArgsConstructor; import org.apache.calcite.rel.RelNode; import org.apache.calcite.rel.type.RelDataType; import org.apache.calcite.rel.type.RelDataTypeField; @@ -64,17 +65,24 @@ import org.apache.calcite.sql.pretty.SqlPrettyWriter; import org.apache.flink.table.planner.plan.schema.RawRelDataType; -@AllArgsConstructor +@RequiredArgsConstructor(access = AccessLevel.PROTECTED) public abstract class AbstractJdbcStatementFactory implements JdbcStatementFactory { private static final Pattern POSITIONAL_ARG_PATTERN = Pattern.compile(Pattern.quote(SqrlStatementParser.POSITIONAL_ARGUMENT_PREFIX) + "(\\d+)"); protected final OperatorRuleTransformer dialectCallConverter; - protected final RelToSqlNode relToSqlConverter; - protected final SqlNodeToString sqlNodeToString; + protected final SqlConverters sqlConverters; protected final CreateTableDdlFactory createTableDdlFactory; + protected AbstractJdbcStatementFactory( + Dialect dialect, CreateTableDdlFactory createTableDdlFactory) { + this( + new OperatorRuleTransformer(dialect), + SqlConvertersFactory.get(dialect), + createTableDdlFactory); + } + @Override public QueryResult createQuery( Query query, boolean withView, Map tableIdMap) { @@ -134,14 +142,14 @@ protected QueryResult createQueryInternal( Map tableNameMapping, Documented.Documentation documentation) { var rewrittenRelNode = dialectCallConverter.convert(relNode); - var sqlNodes = relToSqlConverter.convert(rewrittenRelNode, tableNameMapping); - var sql = sqlNodeToString.convert(sqlNodes).getSql(); + var sqlNode = sqlConverters.convert(rewrittenRelNode, tableNameMapping); + var sql = sqlConverters.convert(sqlNode); var qBuilder = ExecutableJdbcReadQuery.builder(); qBuilder.sql(sql); JdbcStatement view = null; if (withView) { - view = getViewStatement(viewName, relNode.getRowType(), sqlNodes, documentation); + view = getViewStatement(viewName, relNode.getRowType(), sqlNode, documentation); } return new JdbcStatementFactory.QueryResult(qBuilder, view); } @@ -225,7 +233,8 @@ protected String createView( var createView = new SqlCreatePostgresView( SqlParserPos.ZERO, true, viewNameIdentifier, columnList, viewSqlNode); - return sqlNodeToString.convert(() -> createView).getSql(); + + return sqlConverters.convert(createView); } public static List quoteIdentifier(List columns) { @@ -299,7 +308,7 @@ private String replacePassthroughArg(MatchResult matchResult, String sql) { private JdbcStatement getViewStatement( String viewName, RelDataType rowType, - SqlNodes sqlNodes, + SqlNode sqlNode, Documented.Documentation documentation) { var viewNameIdentifier = new SqlIdentifier(viewName, SqlParserPos.ZERO); var columnList = @@ -308,7 +317,7 @@ private JdbcStatement getViewStatement( .map(f -> new SqlIdentifier(f.getName(), SqlParserPos.ZERO)) .collect(Collectors.toList()), SqlParserPos.ZERO); - var viewSql = createView(viewNameIdentifier, columnList, sqlNodes.getSqlNode()); + var viewSql = createView(viewNameIdentifier, columnList, sqlNode); return new GenericJdbcStatement( viewName, diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/DuckDbStatementFactory.java b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/DuckDbStatementFactory.java index cd1b6b1ce1..da581c76b2 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/DuckDbStatementFactory.java +++ b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/DuckDbStatementFactory.java @@ -23,9 +23,6 @@ import static com.datasqrl.function.CalciteFunctionUtil.lightweightOp; import com.datasqrl.calcite.Dialect; -import com.datasqrl.calcite.OperatorRuleTransformer; -import com.datasqrl.calcite.convert.DuckDbRelToSqlNode; -import com.datasqrl.calcite.convert.DuckdbSqlNodeToString; import com.datasqrl.calcite.dialect.DuckDbSqlDialect; import com.datasqrl.calcite.type.TypeFactory; import com.datasqrl.config.JdbcDialect; @@ -58,9 +55,7 @@ public class DuckDbStatementFactory extends AbstractJdbcStatementFactory { public DuckDbStatementFactory(EngineConfig engineConfig) { super( - new OperatorRuleTransformer(Dialect.DUCKDB), - new DuckDbRelToSqlNode(), - new DuckdbSqlNodeToString(), + Dialect.DUCKDB, new GenericCreateTableDdlFactory()); // Iceberg creates the tables, DuckDB only queries this.engineConfig = engineConfig; } diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/IcebergStatementFactory.java b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/IcebergStatementFactory.java index 3a484f3c81..955b02b94a 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/IcebergStatementFactory.java +++ b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/IcebergStatementFactory.java @@ -32,7 +32,6 @@ public IcebergStatementFactory() { super( new OperatorRuleTransformer(Dialect.POSTGRES), null, // Iceberg does not support queries - null, // Iceberg does not support queries new IcebergCreateTableDdlFactory()); } diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/PostgresStatementFactory.java b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/PostgresStatementFactory.java index dc464c4a85..bc365942c9 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/PostgresStatementFactory.java +++ b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/PostgresStatementFactory.java @@ -18,9 +18,6 @@ import static com.google.common.base.Preconditions.checkArgument; import com.datasqrl.calcite.Dialect; -import com.datasqrl.calcite.OperatorRuleTransformer; -import com.datasqrl.calcite.convert.PostgresRelToSqlNode; -import com.datasqrl.calcite.convert.PostgresSqlNodeToString; import com.datasqrl.calcite.dialect.ExtendedPostgresSqlDialect; import com.datasqrl.config.JdbcDialect; import com.datasqrl.engine.database.relational.CreateTableJdbcStatement.PartitionType; @@ -47,15 +44,10 @@ import org.apache.calcite.sql.SqlDataTypeSpec; import org.apache.calcite.sql.parser.SqlParserPos; -public class PostgresStatementFactory extends AbstractJdbcStatementFactory - implements JdbcStatementFactory { +public class PostgresStatementFactory extends AbstractJdbcStatementFactory { public PostgresStatementFactory() { - super( - new OperatorRuleTransformer(Dialect.POSTGRES), - new PostgresRelToSqlNode(), - new PostgresSqlNodeToString(), - new PostgresCreateTableDdlFactory(true)); + super(Dialect.POSTGRES, new PostgresCreateTableDdlFactory(true)); } @Override diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/SnowflakeStatementFactory.java b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/SnowflakeStatementFactory.java index 6a582391b8..a7f4be9707 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/SnowflakeStatementFactory.java +++ b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/SnowflakeStatementFactory.java @@ -16,9 +16,6 @@ package com.datasqrl.engine.database.relational; import com.datasqrl.calcite.Dialect; -import com.datasqrl.calcite.OperatorRuleTransformer; -import com.datasqrl.calcite.convert.SnowflakeRelToSqlNode; -import com.datasqrl.calcite.convert.SnowflakeSqlNodeToString; import com.datasqrl.calcite.dialect.ExtendedPostgresSqlDialect; import com.datasqrl.calcite.dialect.snowflake.SqlCreateIcebergTableFromObjectStorage; import com.datasqrl.config.JdbcDialect; @@ -39,11 +36,7 @@ public class SnowflakeStatementFactory extends AbstractJdbcStatementFactory { private final EngineConfig engineConfig; public SnowflakeStatementFactory(EngineConfig engineConfig) { - super( - new OperatorRuleTransformer(Dialect.SNOWFLAKE), - new SnowflakeRelToSqlNode(), - new SnowflakeSqlNodeToString(), - new GenericCreateTableDdlFactory()); // Iceberg does not support queries + super(Dialect.SNOWFLAKE, new GenericCreateTableDdlFactory()); this.engineConfig = engineConfig; } @@ -78,7 +71,7 @@ public String getSnowflakeCreateTable(String tableName) { null, null); - return sqlNodeToString.convert(() -> icebergTable).getSql(); + return sqlConverters.convert(icebergTable); } @Override diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/ddl/InsertStatement.java b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/ddl/InsertStatement.java index 5b4a8f0e8c..edc66658cb 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/ddl/InsertStatement.java +++ b/sqrl-planner/src/main/java/com/datasqrl/engine/database/relational/ddl/InsertStatement.java @@ -15,7 +15,8 @@ */ package com.datasqrl.engine.database.relational.ddl; -import com.datasqrl.calcite.convert.PostgresSqlNodeToString; +import com.datasqrl.calcite.Dialect; +import com.datasqrl.calcite.convert.SqlConvertersFactory; import com.datasqrl.sql.SqlDDLStatement; import com.fasterxml.jackson.annotation.JsonPropertyOrder; import java.util.ArrayList; @@ -23,6 +24,7 @@ import java.util.List; import java.util.stream.Collectors; import lombok.AllArgsConstructor; +import lombok.Getter; import org.apache.calcite.rel.type.RelDataType; import org.apache.calcite.rel.type.RelDataTypeField; import org.apache.calcite.sql.SqlDynamicParam; @@ -36,14 +38,10 @@ @AllArgsConstructor public class InsertStatement implements SqlDDLStatement { - String tableName; + @Getter String tableName; RelDataType tableSchema; - public String getTableName() { - return tableName; - } - public List getParams() { return tableSchema.getFieldList().stream() .map(RelDataTypeField::getName) @@ -77,7 +75,7 @@ public String getSql() { new SqlInsert(SqlParserPos.ZERO, SqlNodeList.EMPTY, targetTable, values, columns); // Convert the INSERT statement to a SQL string - var sql = addValuesKeyword(new PostgresSqlNodeToString().convert(() -> sqlInsert).getSql()); + var sql = addValuesKeyword(SqlConvertersFactory.get(Dialect.POSTGRES).convert(sqlInsert)); return sql; } diff --git a/sqrl-planner/src/main/java/com/datasqrl/engine/stream/flink/sql/RelToFlinkSql.java b/sqrl-planner/src/main/java/com/datasqrl/engine/stream/flink/sql/RelToFlinkSql.java index cfa6b42e39..857a7d3ca5 100644 --- a/sqrl-planner/src/main/java/com/datasqrl/engine/stream/flink/sql/RelToFlinkSql.java +++ b/sqrl-planner/src/main/java/com/datasqrl/engine/stream/flink/sql/RelToFlinkSql.java @@ -15,8 +15,6 @@ */ package com.datasqrl.engine.stream.flink.sql; -import com.datasqrl.calcite.convert.RelToSqlNode.SqlNodes; -import com.datasqrl.calcite.convert.SqlNodeToString.SqlStrings; import com.datasqrl.engine.stream.flink.sql.calcite.FlinkDialect; import java.util.List; import java.util.function.Function; @@ -43,10 +41,6 @@ public class RelToFlinkSql { .withDialect(FlinkDialect.DEFAULT) .withSelectFolding(null); - public static SqlStrings convertToString(SqlNodes sqlNode) { - return () -> convertToString(sqlNode.getSqlNode()); - } - public static List convertToSqlString(List sqlNode) { return sqlNode.stream().map(RelToFlinkSql::convertToString).toList(); } diff --git a/sqrl-planner/src/test/java/com/datasqrl/engine/stream/flink/plan/FlinkSqlNodesTest.java b/sqrl-planner/src/test/java/com/datasqrl/engine/stream/flink/plan/FlinkSqlNodesTest.java index 4071a13db4..3be2d0bb31 100644 --- a/sqrl-planner/src/test/java/com/datasqrl/engine/stream/flink/plan/FlinkSqlNodesTest.java +++ b/sqrl-planner/src/test/java/com/datasqrl/engine/stream/flink/plan/FlinkSqlNodesTest.java @@ -18,14 +18,12 @@ import static org.assertj.core.api.Assertions.assertThat; import com.datasqrl.calcite.Dialect; -import com.datasqrl.calcite.convert.SqlToStringFactory; +import com.datasqrl.calcite.convert.SqlConvertersFactory; import com.datasqrl.engine.stream.flink.FlinkSqlNodes; -import com.datasqrl.engine.stream.flink.FlinkSqlNodes.MetadataEntry; import java.util.Arrays; import java.util.HashMap; import java.util.List; import java.util.Map; -import java.util.Optional; import org.apache.calcite.sql.SqlDataTypeSpec; import org.apache.calcite.sql.SqlIdentifier; import org.apache.calcite.sql.SqlLiteral; @@ -43,13 +41,8 @@ class FlinkSqlNodesTest { - public record MockMetadataEntry( - Optional type, Optional attribute, Optional virtual) - implements MetadataEntry {} - private String unparse(SqlNode node) { - var sqlToString = SqlToStringFactory.get(Dialect.FLINK); - return sqlToString.convert(() -> node).getSql(); + return SqlConvertersFactory.get(Dialect.FLINK).convert(node); } @Test From e589ac243241ff1906c2468e350be1beedd09a9d Mon Sep 17 00:00:00 2001 From: Ferenc Csaky Date: Mon, 27 Jul 2026 15:28:33 +0200 Subject: [PATCH 2/2] update snapshots --- .../datasqrl/DAGPlannerTest/comprehensiveTest.txt | 4 ++-- .../datasqrl/DAGPlannerTest/databasejoin-w-hash.txt | 2 +- .../com/datasqrl/DAGPlannerTest/duplicatePKTest.txt | 2 +- .../DAGPlannerTest/joinTimestampPropagationTest.txt | 4 ++-- .../com/datasqrl/DAGPlannerTest/staticDataTest.txt | 8 ++++---- .../com/datasqrl/DAGPlannerTest/tableRowCounts.txt | 4 ++-- .../datasqrl/DAGPlannerTest/tableStreamJoinTest.txt | 12 ++++++------ .../comprehensiveTest-limit-offset-combinations.txt | 4 ++-- .../comprehensiveTest-parameters-order.txt | 4 ++-- .../GraphQLValidationTest/comprehensiveTest.txt | 4 ++-- .../analytics-only-package-snowflake.txt | 6 +++--- .../UseCaseCompileTest/pg-minimal-flink-package.txt | 2 +- 12 files changed, 28 insertions(+), 28 deletions(-) diff --git a/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/DAGPlannerTest/comprehensiveTest.txt b/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/DAGPlannerTest/comprehensiveTest.txt index 9aa62a3935..69de479963 100644 --- a/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/DAGPlannerTest/comprehensiveTest.txt +++ b/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/DAGPlannerTest/comprehensiveTest.txt @@ -1880,7 +1880,7 @@ END { "name" : "MissedTemporalJoin", "type" : "VIEW", - "sql" : "CREATE OR REPLACE VIEW \"MissedTemporalJoin\"(\"id\", \"customerid\", \"time\", \"entries\", \"customerid0\", \"timestamp\", \"name\") AS SELECT *\nFROM \"ExternalOrders_11\" AS \"ExternalOrders_112\"\n INNER JOIN \"ExplicitDistinct_10\" AS \"ExplicitDistinct_102\" ON \"ExternalOrders_112\".\"customerid\" = \"ExplicitDistinct_102\".\"customerid\"", + "sql" : "CREATE OR REPLACE VIEW \"MissedTemporalJoin\"(\"id\", \"customerid\", \"time\", \"entries\", \"customerid0\", \"timestamp\", \"name\") AS SELECT *\nFROM \"ExternalOrders_11\" AS \"ExternalOrders_110\"\n INNER JOIN \"ExplicitDistinct_10\" AS \"ExplicitDistinct_100\" ON \"ExternalOrders_110\".\"customerid\" = \"ExplicitDistinct_100\".\"customerid\"", "fields" : [ { "name" : "id", @@ -1942,7 +1942,7 @@ END { "name" : "SelectCustomers", "type" : "VIEW", - "sql" : "CREATE OR REPLACE VIEW \"SelectCustomers\"(\"customerid\", \"email\", \"name\", \"lastUpdated\", \"timestamp\") AS SELECT *\nFROM (SELECT \"customerid\", \"email\", \"name\", \"lastUpdated\", \"timestamp\"\n FROM \"SelectCustomers_14\"\n ORDER BY \"timestamp\" DESC NULLS LAST\n FETCH NEXT 10 ROWS ONLY) AS \"t1\"", + "sql" : "CREATE OR REPLACE VIEW \"SelectCustomers\"(\"customerid\", \"email\", \"name\", \"lastUpdated\", \"timestamp\") AS SELECT *\nFROM (SELECT \"customerid\", \"email\", \"name\", \"lastUpdated\", \"timestamp\"\n FROM \"SelectCustomers_14\"\n ORDER BY \"timestamp\" DESC NULLS LAST\n FETCH NEXT 10 ROWS ONLY) AS \"t\"", "description" : "This is for selected customers\n and their orders", "fields" : [ { diff --git a/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/DAGPlannerTest/databasejoin-w-hash.txt b/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/DAGPlannerTest/databasejoin-w-hash.txt index 57247a42f2..05fe94cb59 100644 --- a/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/DAGPlannerTest/databasejoin-w-hash.txt +++ b/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/DAGPlannerTest/databasejoin-w-hash.txt @@ -197,7 +197,7 @@ END { "name" : "OrderEntryJoin", "type" : "VIEW", - "sql" : "CREATE OR REPLACE VIEW \"OrderEntryJoin\"(\"id1\", \"id2\", \"multi\") AS SELECT \"t3\".\"id\" AS \"id1\", \"t4\".\"id\" AS \"id2\", \"t3\".\"quantity\" * \"t4\".\"quantity\" AS \"multi\"\nFROM (SELECT \"id\", \"customerid\", \"time\", \"productid\", \"quantity\", \"discount\"\n FROM \"_OrderEntries_1\") AS \"t3\"\n INNER JOIN (SELECT \"id\", \"customerid\", \"time\", \"productid\", \"quantity\", \"discount\"\n FROM \"_OrderEntries_1\") AS \"t4\" ON \"t3\".\"customerid\" = \"t4\".\"customerid\"", + "sql" : "CREATE OR REPLACE VIEW \"OrderEntryJoin\"(\"id1\", \"id2\", \"multi\") AS SELECT \"t\".\"id\" AS \"id1\", \"t0\".\"id\" AS \"id2\", \"t\".\"quantity\" * \"t0\".\"quantity\" AS \"multi\"\nFROM (SELECT \"id\", \"customerid\", \"time\", \"productid\", \"quantity\", \"discount\"\n FROM \"_OrderEntries_1\") AS \"t\"\n INNER JOIN (SELECT \"id\", \"customerid\", \"time\", \"productid\", \"quantity\", \"discount\"\n FROM \"_OrderEntries_1\") AS \"t0\" ON \"t\".\"customerid\" = \"t0\".\"customerid\"", "fields" : [ { "name" : "id1", diff --git a/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/DAGPlannerTest/duplicatePKTest.txt b/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/DAGPlannerTest/duplicatePKTest.txt index d2b1829b59..2d83109725 100644 --- a/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/DAGPlannerTest/duplicatePKTest.txt +++ b/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/DAGPlannerTest/duplicatePKTest.txt @@ -287,7 +287,7 @@ END { "name" : "JoinTable", "type" : "VIEW", - "sql" : "CREATE OR REPLACE VIEW \"JoinTable\"(\"_uuid\", \"userid\", \"event_time\", \"unique_id\", \"_uuid0\", \"userid0\", \"event_time0\", \"unique_id0\") AS SELECT *\nFROM (SELECT \"_uuid\", \"userid\", \"event_time\", \"unique_id\"\n FROM \"InputDistinct_1\") AS \"t2\"\n INNER JOIN \"_Input_2\" AS \"_Input_22\" ON \"t2\".\"userid\" = \"_Input_22\".\"userid\"", + "sql" : "CREATE OR REPLACE VIEW \"JoinTable\"(\"_uuid\", \"userid\", \"event_time\", \"unique_id\", \"_uuid0\", \"userid0\", \"event_time0\", \"unique_id0\") AS SELECT *\nFROM (SELECT \"_uuid\", \"userid\", \"event_time\", \"unique_id\"\n FROM \"InputDistinct_1\") AS \"t\"\n INNER JOIN \"_Input_2\" AS \"_Input_20\" ON \"t\".\"userid\" = \"_Input_20\".\"userid\"", "fields" : [ { "name" : "_uuid", diff --git a/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/DAGPlannerTest/joinTimestampPropagationTest.txt b/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/DAGPlannerTest/joinTimestampPropagationTest.txt index 50002c8f3c..9abe12ecf3 100644 --- a/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/DAGPlannerTest/joinTimestampPropagationTest.txt +++ b/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/DAGPlannerTest/joinTimestampPropagationTest.txt @@ -323,7 +323,7 @@ END { "name" : "OrderCustomer1", "type" : "VIEW", - "sql" : "CREATE OR REPLACE VIEW \"OrderCustomer1\"(\"id\", \"name\") AS SELECT \"Orders_22\".\"id\", \"Customer_12\".\"name\"\nFROM \"Orders_2\" AS \"Orders_22\"\n INNER JOIN \"Customer_1\" AS \"Customer_12\" ON \"Orders_22\".\"customerid\" = \"Customer_12\".\"customerid\"", + "sql" : "CREATE OR REPLACE VIEW \"OrderCustomer1\"(\"id\", \"name\") AS SELECT \"Orders_20\".\"id\", \"Customer_10\".\"name\"\nFROM \"Orders_2\" AS \"Orders_20\"\n INNER JOIN \"Customer_1\" AS \"Customer_10\" ON \"Orders_20\".\"customerid\" = \"Customer_10\".\"customerid\"", "fields" : [ { "name" : "id", @@ -340,7 +340,7 @@ END { "name" : "OrderCustomer2", "type" : "VIEW", - "sql" : "CREATE OR REPLACE VIEW \"OrderCustomer2\"(\"id\", \"name\", \"timestamp\") AS SELECT \"Orders_22\".\"id\", \"Customer_12\".\"name\", \"GREATEST\"(\"Orders_22\".\"_ingest_time\", \"Customer_12\".\"_ingest_time\") AS \"timestamp\"\nFROM \"Orders_2\" AS \"Orders_22\"\n INNER JOIN \"Customer_1\" AS \"Customer_12\" ON \"Orders_22\".\"customerid\" = \"Customer_12\".\"customerid\"", + "sql" : "CREATE OR REPLACE VIEW \"OrderCustomer2\"(\"id\", \"name\", \"timestamp\") AS SELECT \"Orders_20\".\"id\", \"Customer_10\".\"name\", \"GREATEST\"(\"Orders_20\".\"_ingest_time\", \"Customer_10\".\"_ingest_time\") AS \"timestamp\"\nFROM \"Orders_2\" AS \"Orders_20\"\n INNER JOIN \"Customer_1\" AS \"Customer_10\" ON \"Orders_20\".\"customerid\" = \"Customer_10\".\"customerid\"", "fields" : [ { "name" : "id", diff --git a/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/DAGPlannerTest/staticDataTest.txt b/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/DAGPlannerTest/staticDataTest.txt index 2a81180788..44ef5257bd 100644 --- a/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/DAGPlannerTest/staticDataTest.txt +++ b/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/DAGPlannerTest/staticDataTest.txt @@ -1458,7 +1458,7 @@ END { "name" : "OrderNumbers", "type" : "VIEW", - "sql" : "CREATE OR REPLACE VIEW \"OrderNumbers\"(\"id\", \"time\") AS SELECT \"Orders_42\".\"id\", \"Orders_42\".\"time\"\nFROM \"Orders_4\" AS \"Orders_42\"\n INNER JOIN (SELECT \"id\", CAST(\"id\" AS BIGINT) AS \"id0\"\n FROM \"Numbers_1\") AS \"t3\" ON \"Orders_42\".\"customerid\" = \"t3\".\"id0\"", + "sql" : "CREATE OR REPLACE VIEW \"OrderNumbers\"(\"id\", \"time\") AS SELECT \"Orders_40\".\"id\", \"Orders_40\".\"time\"\nFROM \"Orders_4\" AS \"Orders_40\"\n INNER JOIN (SELECT \"id\", CAST(\"id\" AS BIGINT) AS \"id0\"\n FROM \"Numbers_1\") AS \"t\" ON \"Orders_40\".\"customerid\" = \"t\".\"id0\"", "fields" : [ { "name" : "id", @@ -1475,7 +1475,7 @@ END { "name" : "OrderNumbers2", "type" : "VIEW", - "sql" : "CREATE OR REPLACE VIEW \"OrderNumbers2\"(\"id\", \"time\") AS SELECT \"Orders_42\".\"id\", \"Orders_42\".\"time\"\nFROM \"Orders_4\" AS \"Orders_42\"\n INNER JOIN (SELECT \"id\", CAST(\"id\" AS BIGINT) AS \"id0\"\n FROM (VALUES (1),\n (2),\n (3),\n (4),\n (4)) AS \"t\" (\"id\")) AS \"t5\" ON \"Orders_42\".\"customerid\" = \"t5\".\"id0\"", + "sql" : "CREATE OR REPLACE VIEW \"OrderNumbers2\"(\"id\", \"time\") AS SELECT \"Orders_40\".\"id\", \"Orders_40\".\"time\"\nFROM \"Orders_4\" AS \"Orders_40\"\n INNER JOIN (SELECT \"id\", CAST(\"id\" AS BIGINT) AS \"id0\"\n FROM (VALUES (1),\n (2),\n (3),\n (4),\n (4)) AS \"t\" (\"id\")) AS \"t0\" ON \"Orders_40\".\"customerid\" = \"t0\".\"id0\"", "fields" : [ { "name" : "id", @@ -1492,7 +1492,7 @@ END { "name" : "OrderPairs", "type" : "VIEW", - "sql" : "CREATE OR REPLACE VIEW \"OrderPairs\"(\"id\", \"time\", \"otherid\") AS SELECT \"Orders_42\".\"id\", \"Orders_42\".\"time\", \"t2\".\"id\" AS \"otherid\"\nFROM \"Orders_4\" AS \"Orders_42\",\n (VALUES (1, 1),\n (2, 2),\n (3, 3),\n (4, 4),\n (4, 4)) AS \"t2\" (\"id\", \"pk\")", + "sql" : "CREATE OR REPLACE VIEW \"OrderPairs\"(\"id\", \"time\", \"otherid\") AS SELECT \"Orders_40\".\"id\", \"Orders_40\".\"time\", \"t\".\"id\" AS \"otherid\"\nFROM \"Orders_4\" AS \"Orders_40\",\n (VALUES (1, 1),\n (2, 2),\n (3, 3),\n (4, 4),\n (4, 4)) AS \"t\" (\"id\", \"pk\")", "fields" : [ { "name" : "id", @@ -1514,7 +1514,7 @@ END { "name" : "OrderPairs2", "type" : "VIEW", - "sql" : "CREATE OR REPLACE VIEW \"OrderPairs2\"(\"id\", \"time\", \"pk\") AS SELECT \"Orders_42\".\"id\", \"Orders_42\".\"time\", \"t2\".\"pk\"\nFROM \"Orders_4\" AS \"Orders_42\",\n (VALUES (1, 1),\n (2, 2),\n (3, 3),\n (4, 4),\n (4, 4)) AS \"t2\" (\"id\", \"pk\")", + "sql" : "CREATE OR REPLACE VIEW \"OrderPairs2\"(\"id\", \"time\", \"pk\") AS SELECT \"Orders_40\".\"id\", \"Orders_40\".\"time\", \"t\".\"pk\"\nFROM \"Orders_4\" AS \"Orders_40\",\n (VALUES (1, 1),\n (2, 2),\n (3, 3),\n (4, 4),\n (4, 4)) AS \"t\" (\"id\", \"pk\")", "fields" : [ { "name" : "id", diff --git a/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/DAGPlannerTest/tableRowCounts.txt b/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/DAGPlannerTest/tableRowCounts.txt index b1f26cfc62..5e33b09daf 100644 --- a/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/DAGPlannerTest/tableRowCounts.txt +++ b/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/DAGPlannerTest/tableRowCounts.txt @@ -495,7 +495,7 @@ END { "name" : "MachineIntervalJoin", "type" : "VIEW", - "sql" : "CREATE OR REPLACE VIEW \"MachineIntervalJoin\"(\"sensorid\", \"temperature\", \"event_time\", \"machineid\") AS SELECT \"_SensorReading_32\".\"sensorid\", \"_SensorReading_32\".\"temperature\", \"_SensorReading_32\".\"event_time\", \"_Sensors_42\".\"machineid\"\nFROM \"_SensorReading_3\" AS \"_SensorReading_32\"\n INNER JOIN \"_Sensors_4\" AS \"_Sensors_42\" ON \"_SensorReading_32\".\"sensorid\" = \"_Sensors_42\".\"sensorid\"", + "sql" : "CREATE OR REPLACE VIEW \"MachineIntervalJoin\"(\"sensorid\", \"temperature\", \"event_time\", \"machineid\") AS SELECT \"_SensorReading_30\".\"sensorid\", \"_SensorReading_30\".\"temperature\", \"_SensorReading_30\".\"event_time\", \"_Sensors_40\".\"machineid\"\nFROM \"_SensorReading_3\" AS \"_SensorReading_30\"\n INNER JOIN \"_Sensors_4\" AS \"_Sensors_40\" ON \"_SensorReading_30\".\"sensorid\" = \"_Sensors_40\".\"sensorid\"", "fields" : [ { "name" : "sensorid", @@ -522,7 +522,7 @@ END { "name" : "MachineIntervalJoinLimit", "type" : "VIEW", - "sql" : "CREATE OR REPLACE VIEW \"MachineIntervalJoinLimit\"(\"sensorid\", \"temperature\", \"event_time\", \"machineid\") AS SELECT \"_SensorReading_32\".\"sensorid\", \"_SensorReading_32\".\"temperature\", \"_SensorReading_32\".\"event_time\", \"_Sensors_42\".\"machineid\"\nFROM \"_SensorReading_3\" AS \"_SensorReading_32\"\n INNER JOIN \"_Sensors_4\" AS \"_Sensors_42\" ON \"_SensorReading_32\".\"sensorid\" = \"_Sensors_42\".\"sensorid\" AND \"_SensorReading_32\".\"event_time\" >= \"_Sensors_42\".\"updatedTime\" AND \"_SensorReading_32\".\"event_time\" < \"_Sensors_42\".\"updatedTime\" + INTERVAL '10' MINUTE", + "sql" : "CREATE OR REPLACE VIEW \"MachineIntervalJoinLimit\"(\"sensorid\", \"temperature\", \"event_time\", \"machineid\") AS SELECT \"_SensorReading_30\".\"sensorid\", \"_SensorReading_30\".\"temperature\", \"_SensorReading_30\".\"event_time\", \"_Sensors_40\".\"machineid\"\nFROM \"_SensorReading_3\" AS \"_SensorReading_30\"\n INNER JOIN \"_Sensors_4\" AS \"_Sensors_40\" ON \"_SensorReading_30\".\"sensorid\" = \"_Sensors_40\".\"sensorid\" AND \"_SensorReading_30\".\"event_time\" >= \"_Sensors_40\".\"updatedTime\" AND \"_SensorReading_30\".\"event_time\" < \"_Sensors_40\".\"updatedTime\" + INTERVAL '10' MINUTE", "fields" : [ { "name" : "sensorid", diff --git a/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/DAGPlannerTest/tableStreamJoinTest.txt b/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/DAGPlannerTest/tableStreamJoinTest.txt index 0cdfb10ecf..8b1bed731a 100644 --- a/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/DAGPlannerTest/tableStreamJoinTest.txt +++ b/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/DAGPlannerTest/tableStreamJoinTest.txt @@ -620,7 +620,7 @@ END { "name" : "OrderCustomer", "type" : "VIEW", - "sql" : "CREATE OR REPLACE VIEW \"OrderCustomer\"(\"id\", \"name\", \"customerid\") AS SELECT \"Orders_32\".\"id\", \"Customer_12\".\"name\", \"Orders_32\".\"customerid\"\nFROM \"Orders_3\" AS \"Orders_32\"\n INNER JOIN \"Customer_1\" AS \"Customer_12\" ON \"Orders_32\".\"customerid\" = \"Customer_12\".\"customerid\"", + "sql" : "CREATE OR REPLACE VIEW \"OrderCustomer\"(\"id\", \"name\", \"customerid\") AS SELECT \"Orders_30\".\"id\", \"Customer_10\".\"name\", \"Orders_30\".\"customerid\"\nFROM \"Orders_3\" AS \"Orders_30\"\n INNER JOIN \"Customer_1\" AS \"Customer_10\" ON \"Orders_30\".\"customerid\" = \"Customer_10\".\"customerid\"", "fields" : [ { "name" : "id", @@ -642,7 +642,7 @@ END { "name" : "OrderCustomerConstant", "type" : "VIEW", - "sql" : "CREATE OR REPLACE VIEW \"OrderCustomerConstant\"(\"time\", \"id\", \"name\", \"customerid\") AS SELECT \"Orders_32\".\"time\", \"Orders_32\".\"id\", \"Customer_12\".\"name\", \"Orders_32\".\"customerid\"\nFROM \"Orders_3\" AS \"Orders_32\"\n INNER JOIN \"Customer_1\" AS \"Customer_12\" ON \"Orders_32\".\"customerid\" = \"Customer_12\".\"customerid\" AND \"Customer_12\".\"name\" = 'Robert' AND \"Orders_32\".\"id\" > 5", + "sql" : "CREATE OR REPLACE VIEW \"OrderCustomerConstant\"(\"time\", \"id\", \"name\", \"customerid\") AS SELECT \"Orders_30\".\"time\", \"Orders_30\".\"id\", \"Customer_10\".\"name\", \"Orders_30\".\"customerid\"\nFROM \"Orders_3\" AS \"Orders_30\"\n INNER JOIN \"Customer_1\" AS \"Customer_10\" ON \"Orders_30\".\"customerid\" = \"Customer_10\".\"customerid\" AND \"Customer_10\".\"name\" = 'Robert' AND \"Orders_30\".\"id\" > 5", "fields" : [ { "name" : "time", @@ -669,7 +669,7 @@ END { "name" : "OrderCustomerLeft", "type" : "VIEW", - "sql" : "CREATE OR REPLACE VIEW \"OrderCustomerLeft\"(\"cid\", \"ctime\", \"id\", \"name\", \"customerid\") AS SELECT COALESCE(\"Customer_12\".\"customerid\", 0) AS \"cid\", COALESCE(\"Customer_12\".\"lastUpdated\", 0) AS \"ctime\", \"Orders_32\".\"id\", \"Customer_12\".\"name\", \"Orders_32\".\"customerid\"\nFROM \"Orders_3\" AS \"Orders_32\"\n LEFT JOIN \"Customer_1\" AS \"Customer_12\" ON \"Orders_32\".\"customerid\" = \"Customer_12\".\"customerid\"", + "sql" : "CREATE OR REPLACE VIEW \"OrderCustomerLeft\"(\"cid\", \"ctime\", \"id\", \"name\", \"customerid\") AS SELECT COALESCE(\"Customer_10\".\"customerid\", 0) AS \"cid\", COALESCE(\"Customer_10\".\"lastUpdated\", 0) AS \"ctime\", \"Orders_30\".\"id\", \"Customer_10\".\"name\", \"Orders_30\".\"customerid\"\nFROM \"Orders_3\" AS \"Orders_30\"\n LEFT JOIN \"Customer_1\" AS \"Customer_10\" ON \"Orders_30\".\"customerid\" = \"Customer_10\".\"customerid\"", "fields" : [ { "name" : "cid", @@ -701,7 +701,7 @@ END { "name" : "OrderCustomerLeftExcluded", "type" : "VIEW", - "sql" : "CREATE OR REPLACE VIEW \"OrderCustomerLeftExcluded\"(\"time\", \"id\", \"customerid\") AS SELECT \"Orders_32\".\"time\", \"Orders_32\".\"id\", \"Orders_32\".\"customerid\"\nFROM \"Orders_3\" AS \"Orders_32\"\n LEFT JOIN \"Customer_1\" AS \"Customer_12\" ON \"Orders_32\".\"customerid\" = \"Customer_12\".\"customerid\"\nWHERE \"Customer_12\".\"customerid\" IS NULL AND \"Customer_12\".\"lastUpdated\" IS NULL", + "sql" : "CREATE OR REPLACE VIEW \"OrderCustomerLeftExcluded\"(\"time\", \"id\", \"customerid\") AS SELECT \"Orders_30\".\"time\", \"Orders_30\".\"id\", \"Orders_30\".\"customerid\"\nFROM \"Orders_3\" AS \"Orders_30\"\n LEFT JOIN \"Customer_1\" AS \"Customer_10\" ON \"Orders_30\".\"customerid\" = \"Customer_10\".\"customerid\"\nWHERE \"Customer_10\".\"customerid\" IS NULL AND \"Customer_10\".\"lastUpdated\" IS NULL", "fields" : [ { "name" : "time", @@ -745,7 +745,7 @@ END { "name" : "OrderCustomerRight", "type" : "VIEW", - "sql" : "CREATE OR REPLACE VIEW \"OrderCustomerRight\"(\"ouuid\", \"otime\", \"id\", \"name\", \"customerid\") AS SELECT COALESCE(\"Orders_32\".\"id\", 0) AS \"ouuid\", COALESCE(\"Orders_32\".\"time\", PROCTIME()) AS \"otime\", \"Orders_32\".\"id\", \"Customer_12\".\"name\", \"Orders_32\".\"customerid\"\nFROM \"Orders_3\" AS \"Orders_32\"\n RIGHT JOIN \"Customer_1\" AS \"Customer_12\" ON \"Orders_32\".\"customerid\" = \"Customer_12\".\"customerid\"", + "sql" : "CREATE OR REPLACE VIEW \"OrderCustomerRight\"(\"ouuid\", \"otime\", \"id\", \"name\", \"customerid\") AS SELECT COALESCE(\"Orders_30\".\"id\", 0) AS \"ouuid\", COALESCE(\"Orders_30\".\"time\", PROCTIME()) AS \"otime\", \"Orders_30\".\"id\", \"Customer_10\".\"name\", \"Orders_30\".\"customerid\"\nFROM \"Orders_3\" AS \"Orders_30\"\n RIGHT JOIN \"Customer_1\" AS \"Customer_10\" ON \"Orders_30\".\"customerid\" = \"Customer_10\".\"customerid\"", "fields" : [ { "name" : "ouuid", @@ -777,7 +777,7 @@ END { "name" : "OrderCustomerRightExcluded", "type" : "VIEW", - "sql" : "CREATE OR REPLACE VIEW \"OrderCustomerRightExcluded\"(\"lastUpdated\", \"customerid\", \"name\") AS SELECT \"Customer_12\".\"lastUpdated\", \"Customer_12\".\"customerid\", \"Customer_12\".\"name\"\nFROM \"Orders_3\" AS \"Orders_32\"\n RIGHT JOIN \"Customer_1\" AS \"Customer_12\" ON \"Orders_32\".\"customerid\" = \"Customer_12\".\"customerid\"\nWHERE \"Orders_32\".\"id\" IS NULL AND \"Orders_32\".\"time\" IS NULL", + "sql" : "CREATE OR REPLACE VIEW \"OrderCustomerRightExcluded\"(\"lastUpdated\", \"customerid\", \"name\") AS SELECT \"Customer_10\".\"lastUpdated\", \"Customer_10\".\"customerid\", \"Customer_10\".\"name\"\nFROM \"Orders_3\" AS \"Orders_30\"\n RIGHT JOIN \"Customer_1\" AS \"Customer_10\" ON \"Orders_30\".\"customerid\" = \"Customer_10\".\"customerid\"\nWHERE \"Orders_30\".\"id\" IS NULL AND \"Orders_30\".\"time\" IS NULL", "fields" : [ { "name" : "lastUpdated", diff --git a/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/GraphQLValidationTest/comprehensiveTest-limit-offset-combinations.txt b/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/GraphQLValidationTest/comprehensiveTest-limit-offset-combinations.txt index 90d28f02f0..9e9a86e006 100644 --- a/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/GraphQLValidationTest/comprehensiveTest-limit-offset-combinations.txt +++ b/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/GraphQLValidationTest/comprehensiveTest-limit-offset-combinations.txt @@ -1549,8 +1549,8 @@ CREATE TABLE IF NOT EXISTS "UnnestOrders" ("id" BIGINT NOT NULL, "customerid" BI CREATE INDEX IF NOT EXISTS "SelectCustomers_hash_c2" ON "SelectCustomers" USING hash ("name") >>>postgres-views.sql CREATE OR REPLACE VIEW "MissedTemporalJoin"("id", "customerid", "time", "entries", "customerid0", "timestamp", "name") AS SELECT * -FROM "ExternalOrders" AS "ExternalOrders0" - INNER JOIN "ExplicitDistinct" AS "ExplicitDistinct0" ON "ExternalOrders0"."customerid" = "ExplicitDistinct0"."customerid" +FROM "ExternalOrders" + INNER JOIN "ExplicitDistinct" ON "ExternalOrders"."customerid" = "ExplicitDistinct"."customerid" >>>vertx.json { "models" : { diff --git a/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/GraphQLValidationTest/comprehensiveTest-parameters-order.txt b/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/GraphQLValidationTest/comprehensiveTest-parameters-order.txt index 504f2e1225..ca94187fb9 100644 --- a/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/GraphQLValidationTest/comprehensiveTest-parameters-order.txt +++ b/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/GraphQLValidationTest/comprehensiveTest-parameters-order.txt @@ -1549,8 +1549,8 @@ CREATE TABLE IF NOT EXISTS "UnnestOrders" ("id" BIGINT NOT NULL, "customerid" BI CREATE INDEX IF NOT EXISTS "SelectCustomers_hash_c2" ON "SelectCustomers" USING hash ("name") >>>postgres-views.sql CREATE OR REPLACE VIEW "MissedTemporalJoin"("id", "customerid", "time", "entries", "customerid0", "timestamp", "name") AS SELECT * -FROM "ExternalOrders" AS "ExternalOrders0" - INNER JOIN "ExplicitDistinct" AS "ExplicitDistinct0" ON "ExternalOrders0"."customerid" = "ExplicitDistinct0"."customerid" +FROM "ExternalOrders" + INNER JOIN "ExplicitDistinct" ON "ExternalOrders"."customerid" = "ExplicitDistinct"."customerid" >>>vertx.json { "models" : { diff --git a/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/GraphQLValidationTest/comprehensiveTest.txt b/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/GraphQLValidationTest/comprehensiveTest.txt index 137140b91d..5c36bc6b0b 100644 --- a/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/GraphQLValidationTest/comprehensiveTest.txt +++ b/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/GraphQLValidationTest/comprehensiveTest.txt @@ -1549,8 +1549,8 @@ CREATE TABLE IF NOT EXISTS "UnnestOrders" ("id" BIGINT NOT NULL, "customerid" BI CREATE INDEX IF NOT EXISTS "SelectCustomers_hash_c2" ON "SelectCustomers" USING hash ("name") >>>postgres-views.sql CREATE OR REPLACE VIEW "MissedTemporalJoin"("id", "customerid", "time", "entries", "customerid0", "timestamp", "name") AS SELECT * -FROM "ExternalOrders" AS "ExternalOrders0" - INNER JOIN "ExplicitDistinct" AS "ExplicitDistinct0" ON "ExternalOrders0"."customerid" = "ExplicitDistinct0"."customerid" +FROM "ExternalOrders" + INNER JOIN "ExplicitDistinct" ON "ExternalOrders"."customerid" = "ExplicitDistinct"."customerid" >>>vertx.json { "models" : { diff --git a/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/UseCaseCompileTest/analytics-only-package-snowflake.txt b/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/UseCaseCompileTest/analytics-only-package-snowflake.txt index f98e0d1e27..6ae044514f 100644 --- a/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/UseCaseCompileTest/analytics-only-package-snowflake.txt +++ b/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/UseCaseCompileTest/analytics-only-package-snowflake.txt @@ -366,7 +366,7 @@ CREATE OR REPLACE ICEBERG TABLE ApplicationStatus EXTERNAL_VOLUME = 'MyNewVolume CREATE OR REPLACE ICEBERG TABLE _Applications EXTERNAL_VOLUME = 'MyNewVolume' CATALOG = 'MyCatalog' CATALOG_TABLE_NAME = '_Applications'; CREATE OR REPLACE ICEBERG TABLE _LoanTypes EXTERNAL_VOLUME = 'MyNewVolume' CATALOG = 'MyCatalog' CATALOG_TABLE_NAME = '_LoanTypes' >>>iceberg-snowflake-views.sql -CREATE OR REPLACE VIEW ApplicationInfo(id, customer_id, loan_type_id, amount, duration, application_date, updated_at, id1, id2) AS SELECT _Applications.id, _Applications.customer_id, _Applications.loan_type_id, _Applications.amount, _Applications.duration, _Applications.application_date, _Applications.updated_at, _LoanTypes.id AS id1, _LoanTypes0.id AS id2 +CREATE OR REPLACE VIEW ApplicationInfo(id, customer_id, loan_type_id, amount, duration, application_date, updated_at, id1, id2) AS SELECT "_Applications"."id", "_Applications"."customer_id", "_Applications"."loan_type_id", "_Applications"."amount", "_Applications"."duration", "_Applications"."application_date", "_Applications"."updated_at", "_LoanTypes"."id" AS "id1", "_LoanTypes0"."id" AS "id2" FROM _Applications -INNER JOIN _LoanTypes ON _Applications.loan_type_id = _LoanTypes.id -LEFT JOIN _LoanTypes AS _LoanTypes0 ON _Applications.customer_id = _LoanTypes0.id +INNER JOIN _LoanTypes ON "_Applications"."loan_type_id" = "_LoanTypes"."id" +LEFT JOIN _LoanTypes AS "_LoanTypes0" ON "_Applications"."customer_id" = "_LoanTypes0"."id" diff --git a/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/UseCaseCompileTest/pg-minimal-flink-package.txt b/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/UseCaseCompileTest/pg-minimal-flink-package.txt index 11b87ec6b8..fc2fc2ba4f 100644 --- a/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/UseCaseCompileTest/pg-minimal-flink-package.txt +++ b/sqrl-testing/sqrl-testing-integration/src/test/resources/snapshots/com/datasqrl/UseCaseCompileTest/pg-minimal-flink-package.txt @@ -152,7 +152,7 @@ CREATE OR REPLACE VIEW "SensorAggregate"("sensorid", "maxTemp", "avgTemp") AS SE FROM (SELECT "sensorid", MAX("temperature") AS "maxTemp", AVG("temperature") AS "avgTemp" FROM "SensorReading" GROUP BY "sensorid" - ORDER BY "sensorid" NULLS FIRST) AS "t5" + ORDER BY "sensorid" NULLS FIRST) AS "t1" >>>vertx.json { "models" : {