gianm commented on code in PR #19830:
URL: https://github.com/apache/druid/pull/19830#discussion_r3994879403
##########
docs/development/extensions-core/catalog.md:
##########
@@ -43,6 +43,208 @@ allowing queries to be more concise, and simpler to write.
This also allows the
written into a defined column of the table is consistent with that columns
definition, minimizing errors where unexpected
data is written into a particular column of the table.
+### Effects of a table definition
+
+A table definition takes effect when it is read, not when it is written:
changing one, whether through SQL DDL or the
+REST API, rewrites nothing. What it affects:
+
+- **SELECT queries** report each declared column with its declared type, in
declared order. Segments that store a
+ column as a different physical type are converted to the declared type at
query time. Columns present in segments
+ but not declared in the catalog remain queryable, and follow the declared
columns with their physical types.
+- **INSERT and REPLACE** (the catalog applies only to SQL-based ingestion)
validate against the definition rather
+ than merely defaulting from it. A query column targeting a declared column
must be convertible to the declared type
+ without changing it: inserting a `BIGINT` value into a column declared
`VARCHAR` stores the value as a string, while
+ inserting a `VARCHAR` value into a column declared `BIGINT` is an error.
Columns the query produces that the table
+ does not declare are rejected when the table is
[`sealed`](#table-properties) and ingested normally otherwise.
+ Table properties such as `segmentGranularity` and `clusterKeys` act as
defaults that an individual statement may
+ override, such as with its own `PARTITIONED BY`.
+- **Streaming and native batch ingestion** do not consult the catalog; their
specs are unaffected by any table
+ definition.
+
+### SQL DDL
+
+Tables can be defined with SQL instead of by posting a table specification.
`CREATE TABLE` and `ALTER TABLE` are
+submitted to the Broker like any other SQL statement, and write the same
catalog metadata the REST API does. They
+return no rows.
+
+These statements change catalog metadata only. They never create, modify, or
delete segments: defining a table does
+not ingest anything, and altering a column does not rewrite existing data.
Changes take effect for subsequent queries
+and ingestion, as described in [Effects of a table
definition](#effects-of-a-table-definition).
+
+These statements are disabled by default. Set
`druid.sql.planner.enableCatalogDdl` to `true` on the Broker to enable
+them. They require both `READ` and `WRITE` permission on the datasource, the
same permissions the catalog API
+requires for its write operations, so enabling them lets anyone who can ingest
into a datasource also change its
+catalog definition; leave them disabled if you manage catalog entries with
your own tooling. The setting cannot be
+overridden per query.
+
+The `druid-catalog` extension must be loaded on both the Broker and the
Coordinator; without it, these statements
+report that the extension is not available.
+
+```sql
+CREATE [OR REPLACE] TABLE [IF NOT EXISTS] <table>
+ [ [ SEALED ] ( <table element> [, ...] ) ]
+ [ PARTITIONED BY <granularity> ]
+ [ CLUSTERED BY <column> [, ...] ]
+
+<table element> ::=
+ <column> <type>
+ | PROJECTION <name> AS ( <select> )
+```
+
+`OR REPLACE` replaces the specification of an existing table; `IF NOT EXISTS`
leaves an existing table unchanged.
+The two cannot be combined. `PARTITIONED BY` sets
[`segmentGranularity`](#table-properties) and `CLUSTERED BY` sets
+`clusterKeys`, both of which a later `INSERT` or `REPLACE` inherits unless it
states its own. `SEALED` sets
+[`sealed`](#table-properties), which requires every ingested column to be
declared. It is written just before the
+column list because that is what it describes: the list is the table's whole
schema, so the statement does not accept
+`SEALED` without one. (The `sealed` property itself can still be set on any
table through `SET PROPERTIES` or the
+REST API.)
+
+Note that the table-level `CLUSTERED BY` is a sort order applied to each
ingestion, which is a different thing from
+the `CLUSTERED BY` inside a [`__base` projection](#the-base-table), which
defines how segments physically group rows.
+
+Column types are written as SQL types, such as `VARCHAR`, `BIGINT`, `DOUBLE`,
or `VARCHAR ARRAY`. The `__time` column
+is written as `TIMESTAMP`. Types that have no SQL spelling, such as complex
types, use `TYPE('...')` with the Druid
+native type string:
+
+```sql
+CREATE TABLE "druid"."visits" (
+ __time TIMESTAMP,
+ user_id VARCHAR,
+ pages_visited BIGINT,
+ sketch TYPE('COMPLEX<thetaSketch>')
+)
+PARTITIONED BY DAY
+CLUSTERED BY user_id
+```
+
+`ALTER TABLE` supports one change per statement, so that each statement is a
single atomic catalog operation:
+
+```sql
+ALTER TABLE <table> ADD COLUMN <column> <type>
+ALTER TABLE <table> DROP COLUMN <column>
+ALTER TABLE <table> ALTER COLUMN <column> SET DATA TYPE <type>
+ALTER TABLE <table> ADD [IF NOT EXISTS] PROJECTION <name> AS ( <select> )
Review Comment:
Please change it to `ADD PROJECTION [IF NOT EXISTS] <name>`, as this is more
commonly the way such syntaxes are done in SQL systems. It also better matches
`DROP PROJECTION [IF EXISTS]`.
##########
sql/src/main/java/org/apache/druid/sql/calcite/planner/CatalogColumnTypes.java:
##########
@@ -0,0 +1,145 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you 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 org.apache.druid.sql.calcite.planner;
+
+import org.apache.calcite.avatica.SqlType;
+import org.apache.calcite.sql.SqlDataTypeSpec;
+import org.apache.calcite.sql.SqlIdentifier;
+import org.apache.calcite.sql.SqlTypeNameSpec;
+import org.apache.calcite.sql.SqlUserDefinedTypeNameSpec;
+import org.apache.calcite.sql.type.SqlTypeName;
+import org.apache.druid.catalog.model.Columns;
+import org.apache.druid.java.util.common.IAE;
+import org.apache.druid.java.util.common.StringUtils;
+import org.apache.druid.segment.column.ColumnType;
+
+/**
+ * Converts a SQL type as written in a statement into the type string stored
in the catalog. Calcite accepts any
+ * type spelling; Druid accepts only its own supported types and their
aliases, plus the {@code TYPE('...')} escape
+ * hatch for native type strings that have no SQL spelling, such as {@code
COMPLEX<json>}.
+ * <p>
+ * Druid has its own rules for nullability, so any nullability clause is
ignored.
+ */
+public class CatalogColumnTypes
+{
+ /**
+ * Which statement the type was written in. The two differ in what they
accept, so they are named rather than
+ * flagged.
+ */
+ private enum Target
+ {
+ /**
+ * The {@code EXTEND} clause, describing columns read from an external
input source. Those are read as their
+ * underlying storage type, so {@code TIMESTAMP} is not among them, and
the escape hatch is limited to complex
+ * types.
+ */
+ EXTERNAL,
+
+ /**
+ * A catalog DDL statement, which additionally accepts {@code TIMESTAMP}
(how {@code __time} is spelled in SQL)
+ * and any native type string through the escape hatch.
+ */
+ CATALOG
+ }
+
+ private CatalogColumnTypes()
+ {
+ // No instantiation.
+ }
+
+ public static String forExternalColumn(String name, SqlDataTypeSpec dataType)
+ {
+ return convert(name, dataType, Target.EXTERNAL);
+ }
+
+ public static String forCatalogColumn(String name, SqlDataTypeSpec dataType)
+ {
+ final String typeString = convert(name, dataType, Target.CATALOG);
+ // the catalog rejects unparseable types at write time, but catching it
here attributes the error to the statement
+ // rather than to a Coordinator round trip.
+ if (Columns.druidTypeFromString(typeString) == null) {
+ throw unsupportedType(name, dataType);
+ }
+ return typeString;
+ }
+
+ private static String convert(String name, SqlDataTypeSpec dataType, Target
target)
+ {
+ final SqlTypeNameSpec spec = dataType.getTypeNameSpec();
+ if (spec == null) {
+ throw unsupportedType(name, dataType);
+ }
+ final SqlIdentifier typeNameIdentifier = spec.getTypeName();
+ if (typeNameIdentifier == null || !typeNameIdentifier.isSimple()) {
+ throw unsupportedType(name, dataType);
+ }
+ final String simpleName = typeNameIdentifier.getSimple();
+
+ if (spec instanceof SqlUserDefinedTypeNameSpec) {
+ // The TYPE('...') escape hatch names a Druid native type. Parse and
validate rather than passing the raw
+ // string downstream, where a malformed type string would silently
resolve to a different type, and return the
+ // canonical form.
+ if (target == Target.EXTERNAL &&
!StringUtils.toLowerCase(simpleName).startsWith("complex<")) {
Review Comment:
Any reason to not allow non-complex types here? Seems like it would be fine
to allow arrays and primitives.
##########
sql/src/main/java/org/apache/druid/sql/calcite/planner/CatalogColumnTypes.java:
##########
@@ -0,0 +1,145 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you 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 org.apache.druid.sql.calcite.planner;
+
+import org.apache.calcite.avatica.SqlType;
+import org.apache.calcite.sql.SqlDataTypeSpec;
+import org.apache.calcite.sql.SqlIdentifier;
+import org.apache.calcite.sql.SqlTypeNameSpec;
+import org.apache.calcite.sql.SqlUserDefinedTypeNameSpec;
+import org.apache.calcite.sql.type.SqlTypeName;
+import org.apache.druid.catalog.model.Columns;
+import org.apache.druid.java.util.common.IAE;
+import org.apache.druid.java.util.common.StringUtils;
+import org.apache.druid.segment.column.ColumnType;
+
+/**
+ * Converts a SQL type as written in a statement into the type string stored
in the catalog. Calcite accepts any
+ * type spelling; Druid accepts only its own supported types and their
aliases, plus the {@code TYPE('...')} escape
+ * hatch for native type strings that have no SQL spelling, such as {@code
COMPLEX<json>}.
+ * <p>
+ * Druid has its own rules for nullability, so any nullability clause is
ignored.
+ */
+public class CatalogColumnTypes
+{
+ /**
+ * Which statement the type was written in. The two differ in what they
accept, so they are named rather than
+ * flagged.
+ */
+ private enum Target
+ {
+ /**
+ * The {@code EXTEND} clause, describing columns read from an external
input source. Those are read as their
+ * underlying storage type, so {@code TIMESTAMP} is not among them, and
the escape hatch is limited to complex
+ * types.
+ */
+ EXTERNAL,
+
+ /**
+ * A catalog DDL statement, which additionally accepts {@code TIMESTAMP}
(how {@code __time} is spelled in SQL)
+ * and any native type string through the escape hatch.
+ */
+ CATALOG
+ }
+
+ private CatalogColumnTypes()
+ {
+ // No instantiation.
+ }
+
+ public static String forExternalColumn(String name, SqlDataTypeSpec dataType)
+ {
+ return convert(name, dataType, Target.EXTERNAL);
+ }
+
+ public static String forCatalogColumn(String name, SqlDataTypeSpec dataType)
+ {
+ final String typeString = convert(name, dataType, Target.CATALOG);
+ // the catalog rejects unparseable types at write time, but catching it
here attributes the error to the statement
+ // rather than to a Coordinator round trip.
+ if (Columns.druidTypeFromString(typeString) == null) {
+ throw unsupportedType(name, dataType);
+ }
+ return typeString;
+ }
+
+ private static String convert(String name, SqlDataTypeSpec dataType, Target
target)
+ {
+ final SqlTypeNameSpec spec = dataType.getTypeNameSpec();
+ if (spec == null) {
+ throw unsupportedType(name, dataType);
+ }
+ final SqlIdentifier typeNameIdentifier = spec.getTypeName();
+ if (typeNameIdentifier == null || !typeNameIdentifier.isSimple()) {
+ throw unsupportedType(name, dataType);
+ }
+ final String simpleName = typeNameIdentifier.getSimple();
+
+ if (spec instanceof SqlUserDefinedTypeNameSpec) {
+ // The TYPE('...') escape hatch names a Druid native type. Parse and
validate rather than passing the raw
+ // string downstream, where a malformed type string would silently
resolve to a different type, and return the
+ // canonical form.
+ if (target == Target.EXTERNAL &&
!StringUtils.toLowerCase(simpleName).startsWith("complex<")) {
+ throw unsupportedType(name, dataType);
+ }
+ final ColumnType nativeType = ColumnType.fromString(simpleName);
+ if (nativeType == null) {
Review Comment:
Should check that the complex type is valid, or else this will cause
failures later when we attempt to apply the table spec, leading to a broken
table.
##########
sql/src/test/java/org/apache/druid/sql/calcite/CalciteCatalogDdlTest.java:
##########
@@ -0,0 +1,1089 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you 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 org.apache.druid.sql.calcite;
+
+import com.google.common.collect.ImmutableList;
+import com.google.common.collect.ImmutableMap;
+import com.google.common.collect.ImmutableSet;
+import org.apache.calcite.sql.SqlExplain;
+import org.apache.calcite.sql.SqlExplainFormat;
+import org.apache.calcite.sql.SqlExplainLevel;
+import org.apache.calcite.sql.SqlNode;
+import org.apache.calcite.sql.parser.SqlParserPos;
+import org.apache.druid.catalog.model.ClusteredValueGroupsBaseTableMetadata;
+import org.apache.druid.catalog.model.ColumnSpec;
+import org.apache.druid.catalog.model.DatasourceBaseTableMetadata;
+import org.apache.druid.catalog.model.DatasourceProjectionMetadata;
+import org.apache.druid.catalog.model.TableId;
+import org.apache.druid.catalog.model.TableMetadata;
+import org.apache.druid.catalog.model.TableSpec;
+import org.apache.druid.catalog.model.table.ClusterKeySpec;
+import org.apache.druid.catalog.model.table.DatasourceDefn;
+import org.apache.druid.data.input.impl.AggregateProjectionSpec;
+import org.apache.druid.error.DruidException;
+import org.apache.druid.java.util.common.granularity.Granularities;
+import org.apache.druid.server.security.Action;
+import org.apache.druid.server.security.AuthConfig;
+import org.apache.druid.server.security.AuthenticationResult;
+import org.apache.druid.server.security.ForbiddenException;
+import org.apache.druid.server.security.Resource;
+import org.apache.druid.server.security.ResourceAction;
+import org.apache.druid.server.security.ResourceType;
+import org.apache.druid.sql.DirectStatement;
+import org.apache.druid.sql.SqlQueryPlus;
+import
org.apache.druid.sql.calcite.CalciteCatalogDdlTest.CatalogDdlComponentSupplier;
+import org.apache.druid.sql.calcite.parser.DruidSqlParser;
+import org.apache.druid.sql.calcite.planner.CatalogTableWriter;
+import org.apache.druid.sql.calcite.planner.DruidPlanner;
+import org.apache.druid.sql.calcite.planner.PlannerConfig;
+import org.apache.druid.sql.calcite.planner.PlannerFactory;
+import org.apache.druid.sql.calcite.util.CalciteTests;
+import
org.apache.druid.sql.calcite.util.SqlTestFramework.StandardComponentSupplier;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import javax.annotation.Nullable;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * Tests that catalog DDL statements plan into the catalog operations they
claim to, using a writer that records
+ * calls instead of contacting a Coordinator.
+ */
[email protected](CatalogDdlComponentSupplier.class)
+public class CalciteCatalogDdlTest extends BaseCalciteQueryTest
+{
+ private static final RecordingCatalogTableWriter WRITER = new
RecordingCatalogTableWriter();
+
+ public static class CatalogDdlComponentSupplier extends
StandardComponentSupplier
+ {
+ public CatalogDdlComponentSupplier(TempDirProducer tempFolderProducer)
+ {
+ super(tempFolderProducer);
+ }
+
+ @Override
+ public CatalogTableWriter createCatalogTableWriter()
+ {
+ return WRITER;
+ }
+ }
+
+ @BeforeEach
+ public void resetWriter()
+ {
+ WRITER.reset();
+ }
+
+ @Test
+ public void testCreateTable()
+ {
+ execute("CREATE TABLE tbl (__time TIMESTAMP, page VARCHAR, cnt BIGINT)");
+
+ assertEquals(1, WRITER.calls.size());
+ final RecordingCatalogTableWriter.Call call = WRITER.calls.get(0);
+ assertEquals("createTable", call.operation);
+ assertEquals(TableId.datasource("tbl"), call.tableId);
+ assertEquals(DatasourceDefn.TABLE_TYPE, call.spec.type());
+ assertEquals(ImmutableMap.of(), call.spec.properties());
+ assertEquals(
+ ImmutableList.of(
+ new ColumnSpec("__time", "TIMESTAMP", null),
+ new ColumnSpec("page", "VARCHAR", null),
+ new ColumnSpec("cnt", "BIGINT", null)
+ ),
+ call.spec.columns()
+ );
+ assertFalse(call.ifNotExists);
+ assertFalse(call.replace);
+ }
+
+ @Test
+ public void testCreateTableWithPartitioningAndClustering()
+ {
+ execute("CREATE TABLE tbl (page VARCHAR, cnt BIGINT) PARTITIONED BY DAY
CLUSTERED BY page, cnt");
+
+ final TableSpec spec = WRITER.calls.get(0).spec;
+ assertEquals("P1D",
spec.properties().get(DatasourceDefn.SEGMENT_GRANULARITY_PROPERTY));
+ assertEquals(
+ ImmutableList.of(new ClusterKeySpec("page", false), new
ClusterKeySpec("cnt", false)),
+ spec.properties().get(DatasourceDefn.CLUSTER_KEYS_PROPERTY)
+ );
+ }
+
+ @Test
+ public void testCreateTablePartitionedByAll()
+ {
+ execute("CREATE TABLE tbl (page VARCHAR) PARTITIONED BY ALL TIME");
+ assertEquals("ALL",
WRITER.calls.get(0).spec.properties().get(DatasourceDefn.SEGMENT_GRANULARITY_PROPERTY));
+ }
+
+ @Test
+ public void testCreateTableTypeCanonicalization()
+ {
+ execute(
+ "CREATE TABLE tbl (a CHAR, b INTEGER, c REAL, d DOUBLE, e VARCHAR
ARRAY, f TYPE('complex<json>'))"
+ );
+ assertEquals(
+ ImmutableList.of(
+ new ColumnSpec("a", "VARCHAR", null),
+ new ColumnSpec("b", "BIGINT", null),
+ new ColumnSpec("c", "FLOAT", null),
+ new ColumnSpec("d", "DOUBLE", null),
+ new ColumnSpec("e", "VARCHAR ARRAY", null),
+ new ColumnSpec("f", "COMPLEX<json>", null)
+ ),
+ WRITER.calls.get(0).spec.columns()
+ );
+ }
+
+ @Test
+ public void testCreateTableFlags()
+ {
+ execute("CREATE OR REPLACE TABLE tbl (a VARCHAR)");
+ assertTrue(WRITER.calls.get(0).replace);
+
+ WRITER.reset();
+ execute("CREATE TABLE IF NOT EXISTS tbl (a VARCHAR)");
+ assertTrue(WRITER.calls.get(0).ifNotExists);
+ }
+
+ @Test
+ public void testCreateTableInDruidSchema()
+ {
+ execute("CREATE TABLE druid.tbl (a VARCHAR)");
+ assertEquals(TableId.datasource("tbl"), WRITER.calls.get(0).tableId);
+ }
+
+ @Test
+ public void testResourceActionsAreDatasourceReadAndWrite()
+ {
+ // READ and WRITE, matching the Coordinator catalog API's requirement for
write operations: the write is forwarded
+ // with escalated credentials, so this statement-level check is the only
one the issuing user faces.
+ final DirectStatement stmt = statement("CREATE TABLE tbl (a VARCHAR)");
+ stmt.execute();
+ assertEquals(
+ ImmutableSet.of(
+ new ResourceAction(new Resource("tbl", ResourceType.DATASOURCE),
Action.READ),
+ new ResourceAction(new Resource("tbl", ResourceType.DATASOURCE),
Action.WRITE)
+ ),
+ ImmutableSet.copyOf(stmt.resources())
+ );
+ }
+
+ /**
+ * The grammar cannot express EXPLAIN of a DDL statement, so the planner's
defensive guard is exercised directly:
+ * parse the CREATE, wrap it in a hand-built SqlExplain, and hand that to
the planner as if the parser had produced
+ * it.
+ */
+ @Test
+ public void testExplainOfDdlIsRejectedByThePlanner()
+ {
+ final SqlNode create = DruidSqlParser.parse("CREATE TABLE tbl (a
VARCHAR)", true).getMainStatement();
+ final SqlExplain explain = new SqlExplain(
+ SqlParserPos.ZERO,
+ create,
+ SqlExplainLevel.ALL_ATTRIBUTES.symbol(SqlParserPos.ZERO),
+ SqlExplain.Depth.PHYSICAL.symbol(SqlParserPos.ZERO),
+ SqlExplainFormat.TEXT.symbol(SqlParserPos.ZERO),
+ 0
+ );
+ final PlannerFactory plannerFactory = queryFramework()
+
.plannerFixture(PlannerConfig.builder().enableCatalogDdl(true).build(), new
AuthConfig())
+ .plannerFactory();
+ try (DruidPlanner planner = plannerFactory.createPlanner(
+ queryFramework().engine(),
+ "EXPLAIN PLAN FOR CREATE TABLE tbl (a VARCHAR)",
+ explain,
+ CalciteTests.SUPER_USER_AUTH_RESULT,
+ Collections.emptySet(),
+ Collections.emptyMap(),
+ null
+ )) {
+ final DruidException e = assertThrows(DruidException.class,
planner::validate);
+ assertTrue(e.getMessage().contains("EXPLAIN is not supported for
[CREATE_TABLE]"), e.getMessage());
+ }
+ assertTrue(WRITER.calls.isEmpty());
+ }
+
+ /**
+ * An unauthorized statement must not touch the catalog: authorization runs
before planning, and the write itself is
+ * deferred to the result's run, so a rejected statement leaves no trace.
Covers both denial modes: no permissions at
+ * all, and READ without WRITE (DDL requires both).
+ */
+ @Test
+ public void testUnauthorizedDdlNeverReachesTheWriter()
+ {
+ assertThrows(
+ ForbiddenException.class,
+ () -> statement("CREATE TABLE forbiddenDatasource (a VARCHAR)",
CalciteTests.REGULAR_USER_AUTH_RESULT)
+ .execute()
+ );
+ assertThrows(
+ ForbiddenException.class,
+ () -> statement("CREATE TABLE readOnlyTbl (a VARCHAR)",
CalciteTests.REGULAR_USER_AUTH_RESULT).execute()
+ );
+ assertTrue(WRITER.calls.isEmpty());
+ }
+
+ /**
+ * Planning a DDL statement is free of side effects; the catalog write
happens when the result runs.
+ */
+ @Test
+ public void testWriteHappensOnRunNotOnPlan()
+ {
+ final DirectStatement.ResultSet resultSet = statement("CREATE TABLE tbl (a
VARCHAR)").plan();
+ assertTrue(WRITER.calls.isEmpty());
+
+ resultSet.run();
+ assertEquals(1, WRITER.calls.size());
+ assertEquals("createTable", WRITER.calls.get(0).operation);
+ }
+
+ @Test
+ public void testDdlReturnsNoRows()
+ {
+ final DirectStatement stmt = statement("CREATE TABLE tbl (a VARCHAR)");
+ final List<Object[]> results = stmt.execute().getResults().toList();
+ assertEquals(ImmutableList.of(), results);
+ }
+
+ @Test
+ public void testAlterTableAddColumn()
+ {
+ // ADD and ALTER differ only in which existence outcome is an error, and
that rule is enforced inside the
+ // Coordinator's update transaction, so all the statement does is pick the
verb. EditorTest covers the rules.
+ execute("ALTER TABLE tbl ADD COLUMN b BIGINT");
+
+ final RecordingCatalogTableWriter.Call call =
WRITER.lastCall("addColumns");
+ assertEquals(ImmutableList.of(new ColumnSpec("b", "BIGINT", null)),
call.columns);
+ }
+
+ @Test
+ public void testAlterTableDropColumn()
+ {
+ execute("ALTER TABLE tbl DROP COLUMN gone");
+ assertEquals(ImmutableList.of("gone"),
WRITER.lastCall("dropColumns").droppedColumns);
+ }
+
+ @Test
+ public void testAlterTableAlterColumn()
+ {
+ execute("ALTER TABLE tbl ALTER COLUMN cnt SET DATA TYPE DOUBLE");
+ assertEquals(
+ ImmutableList.of(new ColumnSpec("cnt", "DOUBLE", null)),
+ WRITER.lastCall("alterColumns").columns
+ );
+ }
+
+ @Test
+ public void testAlterTableSetProperties()
+ {
+ execute("ALTER TABLE tbl SET PROPERTIES (targetSegmentRows = 3000000,
sealed = TRUE, description = 'hi')");
+
+ final Map<String, Object> properties =
WRITER.lastCall("updateProperties").properties;
+ assertEquals(3000000L, properties.get("targetSegmentRows"));
+ assertEquals(true, properties.get("sealed"));
+ assertEquals("hi", properties.get("description"));
+ }
+
+ @Test
+ public void testAlterTableSetPropertyToNullRemovesIt()
+ {
+ execute("ALTER TABLE tbl SET PROPERTIES (description = NULL)");
+ final Map<String, Object> properties =
WRITER.lastCall("updateProperties").properties;
+ assertTrue(properties.containsKey("description"));
+ assertNull(properties.get("description"));
+ }
+
+ @Test
+ public void testCreateTableRejectsDuplicateColumn()
+ {
+ final DruidException e = assertThrows(
+ DruidException.class,
+ () -> execute("CREATE TABLE tbl (a VARCHAR, a BIGINT)")
+ );
+ assertTrue(e.getMessage().contains("Column [a] is declared more than
once"));
+ }
+
+ @Test
+ public void testCreateTableRejectsUnsupportedType()
+ {
+ final DruidException e = assertThrows(
+ DruidException.class,
+ () -> execute("CREATE TABLE tbl (a TYPE('NOT_A_TYPE'))")
+ );
+ assertTrue(e.getMessage().contains("unsupported type"));
+ }
+
+ /**
+ * Any spelling that resolves to a LONG is accepted for the time column, and
is stored as written.
+ */
+ @Test
+ public void testCreateTableTimeColumnSpellings()
+ {
+ execute("CREATE TABLE tbl (__time BIGINT)");
+ assertEquals(ImmutableList.of(new ColumnSpec("__time", "BIGINT", null)),
WRITER.calls.get(0).spec.columns());
+
+ WRITER.reset();
+ execute("CREATE TABLE tbl (__time TYPE('LONG'))");
+ assertEquals(ImmutableList.of(new ColumnSpec("__time", "LONG", null)),
WRITER.calls.get(0).spec.columns());
+ }
+
+ @Test
+ public void testCreateTableRejectsNonLongTimeColumn()
+ {
+ final DruidException e = assertThrows(
+ DruidException.class,
+ () -> execute("CREATE TABLE tbl (__time VARCHAR)")
+ );
+ assertTrue(e.getMessage().contains("Column [__time] must have type"));
+ }
+
+ @Test
+ public void testCreateTableRejectsNonDruidSchema()
+ {
+ final DruidException e = assertThrows(
+ DruidException.class,
+ () -> execute("CREATE TABLE lookup.tbl (a VARCHAR)")
+ );
+ assertTrue(e.getMessage().contains("is not a Druid datasource"));
+ }
+
+ @Test
+ public void testCreateTableRejectsBothReplaceAndIfNotExists()
+ {
+ final DruidException e = assertThrows(
+ DruidException.class,
+ () -> execute("CREATE OR REPLACE TABLE IF NOT EXISTS tbl (a VARCHAR)")
+ );
+ assertTrue(e.getMessage().contains("Cannot specify both OR REPLACE and IF
NOT EXISTS"));
+ }
+
+ @Test
+ public void testCreateTableRejectsClusteringExpression()
+ {
+ final DruidException e = assertThrows(
+ DruidException.class,
+ () -> execute("CREATE TABLE tbl (a VARCHAR) CLUSTERED BY a DESC")
+ );
+ assertTrue(e.getMessage().contains("must be a column name"));
+ }
+
+ /**
+ * The feature is off unless an operator turns it on, so that upgrading a
cluster does not silently widen what a
+ * datasource WRITE permission allows.
+ */
+ @Test
+ public void testDdlIsDisabledByDefault()
+ {
+ final DirectStatement stmt =
getSqlStatementFactory(PlannerConfig.builder().build(), new AuthConfig())
+ .directStatement(
+ SqlQueryPlus.builder("CREATE TABLE tbl (a VARCHAR)")
+ .auth(CalciteTests.SUPER_USER_AUTH_RESULT)
+ .build()
+ );
+ final DruidException e = assertThrows(DruidException.class, stmt::execute);
+ assertTrue(e.getMessage().contains("druid.sql.planner.enableCatalogDdl"),
e.getMessage());
+ assertEquals(ImmutableList.of(), WRITER.calls);
+ }
+
+ /**
+ * The stored specification must be the one the planner would produce for
the equivalent query, since that is what
+ * makes a projection match at query time. Pinned as JSON so a change in
planner output is visible here.
+ */
+ @Test
+ public void testCreateTableWithProjection() throws Exception
+ {
+ execute(
+ "CREATE TABLE tbl (__time TIMESTAMP, page VARCHAR, cnt BIGINT,"
+ + " PROJECTION daily AS (SELECT TIME_FLOOR(__time, 'P1D'), page,
SUM(cnt) AS total GROUP BY 1, 2))"
+ );
+
+ assertEquals(
+ "[{\"spec\":{\"type\":\"aggregate\",\"name\":\"daily\","
+ + "\"virtualColumns\":[{\"type\":\"expression\",\"name\":\"v0\","
+ +
"\"expression\":\"timestamp_floor(\\\"__time\\\",'P1D',null,'UTC')\",\"outputType\":\"LONG\"}],"
+ +
"\"groupingColumns\":[{\"type\":\"long\",\"name\":\"v0\",\"multiValueHandling\":\"SORTED_ARRAY\","
+ +
"\"createBitmapIndex\":false},{\"type\":\"string\",\"name\":\"page\","
+ +
"\"multiValueHandling\":\"SORTED_ARRAY\",\"createBitmapIndex\":true}],"
+ +
"\"aggregators\":[{\"type\":\"longSum\",\"name\":\"total\",\"fieldName\":\"cnt\"}],"
+ + "\"ordering\":[{\"columnName\":\"v0\",\"order\":\"ascending\"},"
+ + "{\"columnName\":\"page\",\"order\":\"ascending\"}]}}]",
+ projectionsJson()
+ );
+ }
+
+ /**
+ * A projection body is planned under the statement's own context, so a SET
clause that changes how the equivalent
+ * query would plan changes the stored definition the same way. Here the
session time zone reaches the TIME_FLOOR.
+ */
+ @Test
+ public void testProjectionBodyHonorsStatementContext() throws Exception
+ {
+ execute(
+ "SET sqlTimeZone = 'America/Los_Angeles';\n"
+ + "CREATE TABLE tbl (__time TIMESTAMP, page VARCHAR, cnt BIGINT,"
+ + " PROJECTION daily AS (SELECT TIME_FLOOR(__time, 'P1D'), page,
SUM(cnt) AS total GROUP BY 1, 2))"
+ );
+
+ assertTrue(
+
projectionsJson().contains("timestamp_floor(\\\"__time\\\",'P1D',null,'America/Los_Angeles')"),
+ projectionsJson()
+ );
+ }
+
+ /**
+ * The overrides the lift depends on are applied on top of the statement's
context, so a SET clause cannot put the
+ * planner into a shape the lift does not understand.
+ */
+ @Test
+ public void testProjectionBodyContextCannotOverrideDeterministicOverrides()
throws Exception
+ {
+ execute(
+ "SET sqlUseGranularity = TRUE;\n"
+ + "CREATE TABLE tbl (__time TIMESTAMP, page VARCHAR, cnt BIGINT,"
+ + " PROJECTION daily AS (SELECT TIME_FLOOR(__time, 'P1D'), page,
SUM(cnt) AS total GROUP BY 1, 2))"
+ );
+
+ // Still lifted as an ordinary grouping column rather than a query
granularity, exactly as without the SET.
+ assertTrue(
+
projectionsJson().contains("timestamp_floor(\\\"__time\\\",'P1D',null,'UTC')"),
+ projectionsJson()
+ );
+ }
+
+ /**
+ * A projection defined with TIME_FLOOR must carry a granularity the segment
layer can recover, which is how the
+ * projection gets matched to time-grouped queries.
+ */
+ @Test
+ public void testProjectionGranularityIsRecoverable()
+ {
+ execute(
+ "CREATE TABLE tbl (__time TIMESTAMP, page VARCHAR, cnt BIGINT,"
+ + " PROJECTION hourly AS (SELECT TIME_FLOOR(__time, 'PT1H'), page,
SUM(cnt) AS total GROUP BY 1, 2))"
+ );
+
+ final AggregateProjectionSpec spec = projection(0).getSpec();
+ final String timeColumn = spec.toMetadataSchema().getTimeColumnName();
+ assertEquals("v0", timeColumn);
+ assertEquals(
+ Granularities.HOUR,
+
Granularities.fromVirtualColumn(spec.getVirtualColumns().getVirtualColumn(timeColumn))
+ );
+ }
+
+ @Test
+ public void testProjectionWithFilter()
+ {
+ execute(
+ "CREATE TABLE tbl (__time TIMESTAMP, page VARCHAR, cnt BIGINT,"
+ + " PROJECTION filtered AS (SELECT page, SUM(cnt) AS total WHERE page
<> 'skip' GROUP BY page))"
+ );
+ assertEquals("!page = skip",
projection(0).getSpec().getFilter().toString());
+ }
+
+ /**
+ * A time bound written in the body is moved into the query's intervals
during planning, and has to be put back:
+ * a projection stores a filter, not an interval.
+ */
+ @Test
+ public void testProjectionWithTimeFilter()
+ {
+ execute(
+ "CREATE TABLE tbl (__time TIMESTAMP, page VARCHAR, cnt BIGINT,"
+ + " PROJECTION recent AS (SELECT page, SUM(cnt) AS total"
+ + " WHERE __time >= TIMESTAMP '2020-01-01 00:00:00' GROUP BY page))"
+ );
+ assertNotNull(projection(0).getSpec().getFilter(), "time filter must
survive as a filter");
+
assertTrue(projection(0).getSpec().getFilter().getRequiredColumns().contains("__time"));
+ }
+
+ @Test
+ public void testProjectionSelectDistinct()
+ {
+ execute("CREATE TABLE tbl (a VARCHAR, PROJECTION d AS (SELECT DISTINCT
a))");
+ final AggregateProjectionSpec spec = projection(0).getSpec();
+ assertEquals(1, spec.getGroupingColumns().size());
+ assertEquals("a", spec.getGroupingColumns().get(0).getName());
+ assertEquals(0, spec.getAggregators().length);
+ }
+
+ @Test
+ public void testMultipleProjections()
+ {
+ execute(
+ "CREATE TABLE tbl (a VARCHAR, b BIGINT,"
+ + " PROJECTION p1 AS (SELECT a, SUM(b) AS s GROUP BY a),"
+ + " PROJECTION p2 AS (SELECT b, COUNT(*) AS c GROUP BY b))"
+ );
+ assertEquals(List.of("p1", "p2"),
List.of(projection(0).getSpec().getName(), projection(1).getSpec().getName()));
+ }
+
+ @Test
+ public void testAlterTableAddProjection()
+ {
+ WRITER.existing.put(TableId.datasource("tbl"), tableWithColumns("a"));
+ execute("ALTER TABLE tbl ADD PROJECTION p AS (SELECT a, COUNT(*) AS c
GROUP BY a)");
+
+ final RecordingCatalogTableWriter.Call call =
WRITER.lastCall("addProjection");
+ assertEquals("p", call.projection.getSpec().getName());
+ assertFalse(call.ifNotExists);
+ }
+
+ @Test
+ public void testAlterTableAddProjectionIfNotExists()
+ {
+ WRITER.existing.put(TableId.datasource("tbl"), tableWithColumns("a"));
+ execute("ALTER TABLE tbl ADD IF NOT EXISTS PROJECTION p AS (SELECT a GROUP
BY a)");
+ assertTrue(WRITER.lastCall("addProjection").ifNotExists);
+ }
+
+ @Test
+ public void testAlterTableDropProjection()
+ {
+ execute("ALTER TABLE tbl DROP PROJECTION p");
+ final RecordingCatalogTableWriter.Call call =
WRITER.lastCall("dropProjection");
+ assertEquals("p", call.projectionName);
+ assertFalse(call.ifExists);
+
+ WRITER.reset();
+ execute("ALTER TABLE tbl DROP PROJECTION IF EXISTS p");
+ assertTrue(WRITER.lastCall("dropProjection").ifExists);
+ }
+
+ @Test
+ public void testProjectionRejectsUnaliasedAggregate()
+ {
+ final DruidException e = assertThrows(
+ DruidException.class,
+ () -> execute("CREATE TABLE tbl (a VARCHAR, b BIGINT, PROJECTION p AS
(SELECT a, SUM(b) GROUP BY a))")
+ );
+ assertTrue(e.getMessage().contains("no name"), e.getMessage());
+ }
+
+ @Test
+ public void testProjectionRejectsPostAggregation()
+ {
+ final DruidException e = assertThrows(
+ DruidException.class,
+ () -> execute("CREATE TABLE tbl (a VARCHAR, b BIGINT, PROJECTION p AS
(SELECT a, AVG(b) AS m GROUP BY a))")
+ );
+ assertTrue(e.getMessage().contains("expression over aggregates"),
e.getMessage());
+ }
+
+ /**
+ * Subqueries survive the grammar (they are just expressions in the select
list or WHERE), so the translator walks
+ * the body and rejects them with a pointed message rather than letting the
planner produce something the lift
+ * cannot express.
+ */
+ @Test
+ public void testProjectionRejectsSubquery()
+ {
+ for (String body : new String[]{
+ "SELECT a, (SELECT 1) AS s GROUP BY 1, 2",
+ "SELECT a WHERE a IN (SELECT a) GROUP BY a",
+ }) {
+ final DruidException e = assertThrows(
+ DruidException.class,
+ () -> execute("CREATE TABLE tbl (a VARCHAR, PROJECTION p AS (" +
body + "))"),
+ body
+ );
+ assertTrue(e.getMessage().contains("contains a subquery"), body + " -> "
+ e.getMessage());
+ }
+ }
+
+ /**
+ * Window functions survive the grammar as ordinary select items, and are
rejected by the nested planner because the
+ * projection engine does not declare {@code
EngineFeature.WINDOW_FUNCTIONS}; the error is attributed to the
+ * projection by name.
+ */
+ @Test
+ public void testProjectionRejectsWindowFunctions()
+ {
+ final DruidException e = assertThrows(
+ DruidException.class,
+ () -> execute(
+ "CREATE TABLE tbl (a VARCHAR, b BIGINT,"
+ + " PROJECTION p AS (SELECT a, SUM(b) OVER (PARTITION BY a) AS w
GROUP BY a, b))"
+ )
+ );
+ assertTrue(e.getMessage().contains("Cannot define projection [p]"),
e.getMessage());
+ assertTrue(e.getMessage().contains("window functions"), e.getMessage());
+ }
+
+ @Test
+ public void testProjectionRejectsUnknownColumn()
+ {
+ final DruidException e = assertThrows(
+ DruidException.class,
+ () -> execute("CREATE TABLE tbl (a VARCHAR, PROJECTION p AS (SELECT
nope GROUP BY nope))")
+ );
+ assertTrue(e.getMessage().contains("nope"), e.getMessage());
+ }
+
+ @Test
+ public void testProjectionRejectsNonAggregatingBody()
+ {
+ final DruidException e = assertThrows(
+ DruidException.class,
+ () -> execute("CREATE TABLE tbl (a VARCHAR, PROJECTION p AS (SELECT
a))")
+ );
+ assertTrue(e.getMessage().contains("does not aggregate"), e.getMessage());
+ }
+
+ /**
+ * {@code __base} names the table's own layout and is handled separately;
every other name beginning with the
+ * reserved prefix stays unavailable.
+ */
+ @Test
+ public void testProjectionRejectsOtherReservedNames()
+ {
+ final DruidException e = assertThrows(
+ DruidException.class,
+ () -> execute("CREATE TABLE tbl (a VARCHAR, PROJECTION __other AS
(SELECT a GROUP BY a))")
+ );
+ assertTrue(e.getMessage().contains("reserved name"), e.getMessage());
+ }
+
+ @Test
+ public void testProjectionRejectsDuplicateName()
+ {
+ final DruidException e = assertThrows(
+ DruidException.class,
+ () -> execute(
+ "CREATE TABLE tbl (a VARCHAR, PROJECTION p AS (SELECT a GROUP BY
a),"
+ + " PROJECTION p AS (SELECT a GROUP BY a))"
+ )
+ );
+ assertTrue(e.getMessage().contains("declared more than once"),
e.getMessage());
+ }
+
+ /**
+ * The reserved {@code __base} projection describes the table's own layout,
so it becomes the baseTable property
+ * rather than one of the projections. A computed column becomes a virtual
column materializing the declared column
+ * it fills.
+ */
+ @Test
+ public void testCreateTableWithBaseProjection() throws Exception
+ {
+ execute(
+ "CREATE TABLE tbl SEALED ("
+ + " tenant VARCHAR,"
+ + " bucket BIGINT,"
+ + " __time TIMESTAMP,"
+ + " user_id BIGINT,"
+ + " PROJECTION __base AS ("
+ + " SELECT tenant, ABS(user_id) AS bucket, __time, user_id"
+ + " CLUSTERED BY tenant, bucket"
+ + " )"
+ + ") PARTITIONED BY DAY"
+ );
+
+ final TableSpec spec = WRITER.calls.get(0).spec;
+ assertEquals(true, spec.properties().get(DatasourceDefn.SEALED_PROPERTY));
+
assertNull(spec.properties().get(DatasourceDefn.PROJECTIONS_KEYS_PROPERTY));
+ assertEquals(
+ "{\"clusteringColumns\":[\"tenant\",\"bucket\"],"
+ + "\"virtualColumns\":[{\"type\":\"expression\",\"name\":\"bucket\","
+ + "\"expression\":\"abs(\\\"user_id\\\")\",\"outputType\":\"LONG\"}],"
+ + "\"type\":\"clusteredValueGroups\"}",
+ queryFramework().queryJsonMapper()
+
.writeValueAsString(spec.properties().get(DatasourceDefn.BASE_TABLE_PROPERTY))
+ );
+ }
+
+ /**
+ * A base table column is stored under the name it declares, so a body item
that renames another column is rejected:
+ * only a virtual column materializes a name the body did not select, and a
bare reference produces none.
+ * <p>
+ * Selecting the same column twice also makes the scan's deduplicated column
list shorter than the select list, so
+ * the source of each item is taken from the query's output signature.
Reading the scan's columns positionally used
+ * to run off the end of that list and fail as an internal error.
+ */
+ @Test
+ public void testBaseProjectionRenamedColumnFails()
+ {
+ final DruidException e = assertThrows(
+ DruidException.class,
+ () -> execute(
+ "CREATE TABLE tbl SEALED (id BIGINT, copy BIGINT, __time
TIMESTAMP,"
+ + " PROJECTION __base AS (SELECT id, id AS copy, __time CLUSTERED
BY id))"
+ )
+ );
+ assertTrue(
+ e.getMessage().contains("its column 2 selects [id] but declares it as
[copy]"),
+ e.getMessage()
+ );
+ }
+
+ @Test
+ public void testBaseProjectionWithoutComputedColumns()
+ {
+ execute(
+ "CREATE TABLE tbl SEALED (tenant VARCHAR, __time TIMESTAMP, v BIGINT,"
+ + " PROJECTION __base AS (SELECT tenant, __time, v CLUSTERED BY
tenant))"
+ );
+ final ClusteredValueGroupsBaseTableMetadata baseTable = baseTable();
+ assertEquals(List.of("tenant"), baseTable.getClusteringColumns());
+ assertEquals(0, baseTable.getVirtualColumns().getVirtualColumns().length);
+ }
+
+ /**
+ * A base table and aggregate projections are different catalog entities and
may coexist.
+ */
+ @Test
+ public void testBaseProjectionAlongsideAggregateProjection()
+ {
+ execute(
+ "CREATE TABLE tbl SEALED (tenant VARCHAR, __time TIMESTAMP, v BIGINT,"
+ + " PROJECTION __base AS (SELECT tenant, __time, v CLUSTERED BY
tenant),"
+ + " PROJECTION by_tenant AS (SELECT tenant, SUM(v) AS sum_v GROUP BY
tenant))"
+ );
+ final TableSpec spec = WRITER.calls.get(0).spec;
+ assertNotNull(spec.properties().get(DatasourceDefn.BASE_TABLE_PROPERTY));
Review Comment:
This test class has a lot of `assertNotNull` that really should be asserting
equality to the expected thing. Please tighten them up.
##########
sql/src/main/java/org/apache/druid/sql/calcite/planner/CatalogDdlHandler.java:
##########
@@ -0,0 +1,702 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you 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 org.apache.druid.sql.calcite.planner;
+
+import com.google.common.collect.ImmutableList;
+import com.google.common.collect.ImmutableSet;
+import com.google.common.collect.Iterables;
+import org.apache.calcite.jdbc.CalciteSchema;
+import org.apache.calcite.rel.type.RelDataType;
+import org.apache.calcite.rel.type.RelDataTypeFactory;
+import org.apache.calcite.sql.SqlIdentifier;
+import org.apache.calcite.sql.SqlLiteral;
+import org.apache.calcite.sql.SqlNode;
+import org.apache.calcite.sql.SqlNodeList;
+import org.apache.calcite.sql.type.SqlTypeName;
+import org.apache.druid.catalog.model.ClusteredValueGroupsBaseTableMetadata;
+import org.apache.druid.catalog.model.ColumnSpec;
+import org.apache.druid.catalog.model.Columns;
+import org.apache.druid.catalog.model.DatasourceProjectionMetadata;
+import org.apache.druid.catalog.model.TableId;
+import org.apache.druid.catalog.model.TableMetadata;
+import org.apache.druid.catalog.model.TableSpec;
+import org.apache.druid.catalog.model.table.ClusterKeySpec;
+import org.apache.druid.catalog.model.table.DatasourceDefn;
+import org.apache.druid.common.utils.IdUtils;
+import org.apache.druid.error.DruidException;
+import org.apache.druid.error.InvalidSqlInput;
+import org.apache.druid.java.util.common.IAE;
+import org.apache.druid.java.util.common.granularity.Granularities;
+import org.apache.druid.java.util.common.granularity.Granularity;
+import org.apache.druid.java.util.common.granularity.PeriodGranularity;
+import org.apache.druid.java.util.common.guava.Sequences;
+import org.apache.druid.query.explain.ExplainAttributes;
+import org.apache.druid.segment.column.ColumnType;
+import org.apache.druid.segment.projections.Projections;
+import org.apache.druid.server.QueryResponse;
+import org.apache.druid.server.security.Action;
+import org.apache.druid.server.security.Resource;
+import org.apache.druid.server.security.ResourceAction;
+import org.apache.druid.server.security.ResourceType;
+import org.apache.druid.sql.calcite.parser.DruidSqlAlterTable;
+import org.apache.druid.sql.calcite.parser.DruidSqlColumnDeclaration;
+import org.apache.druid.sql.calcite.parser.DruidSqlCreateTable;
+import org.apache.druid.sql.calcite.parser.DruidSqlParser;
+import org.apache.druid.sql.calcite.parser.DruidSqlPropertyAssignment;
+import org.apache.druid.sql.calcite.parser.SqlGranularityLiteral;
+import org.apache.druid.sql.calcite.parser.SqlProjectionSpec;
+import org.apache.druid.sql.calcite.run.EngineFeature;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.HashSet;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+
+/**
+ * Handles the catalog DDL statements: {@code CREATE TABLE} and {@code ALTER
TABLE}.
+ * <p>
+ * These statements are metadata operations, not queries. They are validated
here, converted to a catalog
+ * {@link TableSpec} or column/property edit, and applied through {@link
CatalogTableWriter}, which forwards them to
+ * the Coordinator. No Calcite validation or query planning takes place, and
no rows are returned.
+ * <p>
+ * Validation is deliberately split. This class checks what it can attribute
to a position in the statement (type
+ * spellings, duplicate columns, granularity, clustering) so that the error
names the offending SQL. The Coordinator
+ * remains authoritative: {@code DatasourceDefn.validate} runs on write, and
its message is surfaced verbatim.
+ */
+public abstract class CatalogDdlHandler extends
SqlStatementHandler.BaseStatementHandler
+{
+ /**
+ * DDL produces no rows. A single-column type still has to be declared,
because JDBC clients ask for a result set
+ * signature when preparing the statement.
+ */
+ private static final RelDataType RESULT_TYPE = resultType();
+
+ protected final SqlIdentifier tableIdentifier;
+ protected TableId tableId;
+
+ protected CatalogDdlHandler(SqlStatementHandler.HandlerContext
handlerContext, SqlIdentifier tableIdentifier)
+ {
+ super(handlerContext);
+ this.tableIdentifier = tableIdentifier;
+ }
+
+ /**
+ * The runtime property that gates these statements. Read from {@link
PlannerConfig} rather than the query context
+ * so that a user cannot turn the feature on for their own statement.
+ */
+ public static final String ENABLE_CATALOG_DDL_PROPERTY =
"druid.sql.planner.enableCatalogDdl";
+
+ /**
+ * The reserved name of the base-table projection, which describes the
physical layout of the table itself. Handled
+ * as a separate catalog property, not as one of the aggregate projections.
+ */
+ public static final String BASE_PROJECTION_NAME = "__base";
Review Comment:
Use `Projections.BASE_TABLE_PROJECTION_NAME`?
##########
server/src/main/java/org/apache/druid/catalog/model/table/DatasourceDefn.java:
##########
@@ -110,6 +124,115 @@ public void validate(ResolvedTable table)
// fail fast instead of surfacing layout problems at ingest time.
baseTable.createSpec(table.spec().columns());
}
+ validateProjections(table);
+ }
+
+ /**
+ * Cross-validate the declared projections. Names must be unique, a
projection must not be coarser than the segments
+ * it lives in, and the types it groups by must agree with the types the
table declares. For a sealed table the
+ * declared columns are the whole schema, so a projection that reads a
column the table does not declare can never be
+ * built and is rejected; for a non-sealed table ingestion may add columns
the catalog has not seen, so only the
+ * projections' internal consistency is checked.
+ */
+ private void validateProjections(ResolvedTable table)
+ {
+ final List<DatasourceProjectionMetadata> projections =
table.decodeProperty(PROJECTIONS_KEYS_PROPERTY);
+ if (projections == null || projections.isEmpty()) {
+ return;
+ }
+
+ final List<AggregateProjectionSpec> specs = new
ArrayList<>(projections.size());
+ for (DatasourceProjectionMetadata projection : projections) {
+ if (projection == null || projection.getSpec() == null) {
+ throw InvalidInput.exception("Projections must each have a [spec]");
+ }
+ specs.add(projection.getSpec());
+ }
+
+ final String granularity =
table.stringProperty(SEGMENT_GRANULARITY_PROPERTY);
+ DataSchema.validateProjections(
+ specs,
+ granularity == null ? null :
CatalogUtils.asDruidGranularity(granularity)
+ );
+
+ validateProjectionGroupingTypes(table, specs);
+
+ if (!table.booleanProperty(SEALED_PROPERTY) || table.spec().columns() ==
null) {
+ return;
+ }
+ final Set<String> declared = new
HashSet<>(CatalogUtils.columnNames(table.spec().columns()));
+ declared.add(Columns.TIME_COLUMN);
+ for (AggregateProjectionSpec spec : specs) {
+ final Set<String> available = new HashSet<>(declared);
+ for (VirtualColumn virtualColumn :
spec.getVirtualColumns().getVirtualColumns()) {
+ available.add(virtualColumn.getOutputName());
+ }
+ for (String required : requiredColumns(spec)) {
+ if (!available.contains(required)) {
+ throw InvalidInput.exception(
+ "Projection [%s] references column [%s], which table [%s] does
not declare",
+ spec.getName(),
+ required,
+ table.spec().type()
Review Comment:
Should be table name, not type, in this position.
##########
docs/development/extensions-core/catalog.md:
##########
@@ -43,6 +43,208 @@ allowing queries to be more concise, and simpler to write.
This also allows the
written into a defined column of the table is consistent with that columns
definition, minimizing errors where unexpected
data is written into a particular column of the table.
+### Effects of a table definition
+
+A table definition takes effect when it is read, not when it is written:
changing one, whether through SQL DDL or the
+REST API, rewrites nothing. What it affects:
+
+- **SELECT queries** report each declared column with its declared type, in
declared order. Segments that store a
+ column as a different physical type are converted to the declared type at
query time. Columns present in segments
+ but not declared in the catalog remain queryable, and follow the declared
columns with their physical types.
+- **INSERT and REPLACE** (the catalog applies only to SQL-based ingestion)
validate against the definition rather
+ than merely defaulting from it. A query column targeting a declared column
must be convertible to the declared type
+ without changing it: inserting a `BIGINT` value into a column declared
`VARCHAR` stores the value as a string, while
+ inserting a `VARCHAR` value into a column declared `BIGINT` is an error.
Columns the query produces that the table
+ does not declare are rejected when the table is
[`sealed`](#table-properties) and ingested normally otherwise.
+ Table properties such as `segmentGranularity` and `clusterKeys` act as
defaults that an individual statement may
+ override, such as with its own `PARTITIONED BY`.
+- **Streaming and native batch ingestion** do not consult the catalog; their
specs are unaffected by any table
+ definition.
+
+### SQL DDL
+
+Tables can be defined with SQL instead of by posting a table specification.
`CREATE TABLE` and `ALTER TABLE` are
+submitted to the Broker like any other SQL statement, and write the same
catalog metadata the REST API does. They
+return no rows.
+
+These statements change catalog metadata only. They never create, modify, or
delete segments: defining a table does
+not ingest anything, and altering a column does not rewrite existing data.
Changes take effect for subsequent queries
+and ingestion, as described in [Effects of a table
definition](#effects-of-a-table-definition).
+
+These statements are disabled by default. Set
`druid.sql.planner.enableCatalogDdl` to `true` on the Broker to enable
+them. They require both `READ` and `WRITE` permission on the datasource, the
same permissions the catalog API
+requires for its write operations, so enabling them lets anyone who can ingest
into a datasource also change its
+catalog definition; leave them disabled if you manage catalog entries with
your own tooling. The setting cannot be
+overridden per query.
+
+The `druid-catalog` extension must be loaded on both the Broker and the
Coordinator; without it, these statements
+report that the extension is not available.
+
+```sql
+CREATE [OR REPLACE] TABLE [IF NOT EXISTS] <table>
+ [ [ SEALED ] ( <table element> [, ...] ) ]
+ [ PARTITIONED BY <granularity> ]
+ [ CLUSTERED BY <column> [, ...] ]
+
+<table element> ::=
+ <column> <type>
+ | PROJECTION <name> AS ( <select> )
+```
+
+`OR REPLACE` replaces the specification of an existing table; `IF NOT EXISTS`
leaves an existing table unchanged.
+The two cannot be combined. `PARTITIONED BY` sets
[`segmentGranularity`](#table-properties) and `CLUSTERED BY` sets
+`clusterKeys`, both of which a later `INSERT` or `REPLACE` inherits unless it
states its own. `SEALED` sets
+[`sealed`](#table-properties), which requires every ingested column to be
declared. It is written just before the
+column list because that is what it describes: the list is the table's whole
schema, so the statement does not accept
+`SEALED` without one. (The `sealed` property itself can still be set on any
table through `SET PROPERTIES` or the
+REST API.)
+
+Note that the table-level `CLUSTERED BY` is a sort order applied to each
ingestion, which is a different thing from
+the `CLUSTERED BY` inside a [`__base` projection](#the-base-table), which
defines how segments physically group rows.
+
+Column types are written as SQL types, such as `VARCHAR`, `BIGINT`, `DOUBLE`,
or `VARCHAR ARRAY`. The `__time` column
+is written as `TIMESTAMP`. Types that have no SQL spelling, such as complex
types, use `TYPE('...')` with the Druid
+native type string:
+
+```sql
+CREATE TABLE "druid"."visits" (
+ __time TIMESTAMP,
+ user_id VARCHAR,
+ pages_visited BIGINT,
+ sketch TYPE('COMPLEX<thetaSketch>')
+)
+PARTITIONED BY DAY
+CLUSTERED BY user_id
+```
+
+`ALTER TABLE` supports one change per statement, so that each statement is a
single atomic catalog operation:
+
+```sql
+ALTER TABLE <table> ADD COLUMN <column> <type>
+ALTER TABLE <table> DROP COLUMN <column>
+ALTER TABLE <table> ALTER COLUMN <column> SET DATA TYPE <type>
+ALTER TABLE <table> ADD [IF NOT EXISTS] PROJECTION <name> AS ( <select> )
+ALTER TABLE <table> DROP PROJECTION [IF EXISTS] <name>
+ALTER TABLE <table> SET PROPERTIES ( <property> = <value> [, ...] )
+```
+
+`ADD COLUMN` fails if the column already exists, and `ALTER COLUMN` fails if
it does not, so a misspelled column name
+is reported rather than quietly creating or replacing a column. Each statement
is also checked against the rest of the
+table definition, not only the part it changes: adding a column, changing a
type, or setting a property is rejected if
+the resulting table would be invalid, such as a segment granularity coarser
than a projection the table declares.
+
+`DROP COLUMN` removes the column's declaration, not its data. A table's SQL
schema is the declared columns followed by
+any other columns present in its segments, so a dropped column that still has
data remains queryable: it loses its
+declared type and position and appears after the declared columns, like any
column the catalog does not know about.
+Whether it can still be ingested then follows the usual rule: rejected if the
table is sealed, accepted as an
+undeclared column otherwise. To remove a column from query results without
rewriting data, use the catalog's
+`hiddenColumns` property instead, which currently has no SQL spelling and is
set through the REST API.
+
+#### Projections
+
+A table may declare [projections](../../querying/projections.md), which are
pre-aggregated views stored inside each
+segment. A projection is written as a `SELECT` over the table's own columns,
with no `FROM` clause:
+
+```sql
+CREATE TABLE "druid"."visits" (
+ __time TIMESTAMP,
+ user_id VARCHAR,
+ user_agent VARCHAR,
+ pages_visited BIGINT,
+ PROJECTION daily_by_agent AS (
+ SELECT TIME_FLOOR(__time, 'P1D'), user_agent, SUM(pages_visited) AS
total_pages
+ WHERE user_agent IS NOT NULL
+ GROUP BY 1, 2
+ )
+)
+PARTITIONED BY DAY
+```
+
+The body is planned exactly as the equivalent query would be, so a projection
matches the queries it was written to
+serve. Every aggregate needs an alias, which becomes the name of the stored
column. Time granularity is expressed
+with `TIME_FLOOR`, as it would be in a query.
+
+Because the body is planned like a query, it is planned under the statement's
own query context, including any `SET`
+clauses. Only context parameters that affect planning can change the stored
definition; parameters that only affect
+query execution have no effect, since the body is planned and stored rather
than run. The context itself is not part
+of the definition: what the catalog stores is the projection the body planned
to, so nothing from the statement's
+context is carried over to queries that later use it. Note also that a
projection is matched to a query by its shape,
+so a definition planned under a context that changes that shape only matches
queries run under the same context.
+
+A projection body accepts a select list, an optional `WHERE` and an optional
`GROUP BY`. It cannot use `ORDER BY`,
+`LIMIT` or `HAVING`: a projection's ordering follows its grouping columns and
is not something you choose. It also
+cannot use joins, subqueries, or expressions computed over aggregates, such as
`SUM(x) / COUNT(x)`: store the two
+aggregates as separate columns and divide at query time. A projection can
store the same
+[aggregation functions that are supported for rollup at ingestion
time](../../multi-stage-query/concepts.md#rollup).
+
+Projections may also be added to and removed from an existing table:
+
+```sql
+ALTER TABLE "druid"."visits" ADD [IF NOT EXISTS] PROJECTION by_agent AS (
+ SELECT user_agent, SUM(pages_visited) AS total_pages GROUP BY user_agent
+)
+ALTER TABLE "druid"."visits" DROP PROJECTION [IF EXISTS] by_agent
+```
+
+Both take effect for subsequent ingestion. Segments already built keep
whatever projections they were built with, so
+dropping a projection does not rewrite data.
+
+#### The base table
+
+The reserved projection name `__base` describes the table's own physical
layout rather than an additional
+pre-aggregation. Defining it makes the table a 'clustered' table: rows of
segments are stored grouped by the clustering
+columns.
+
+Its body lists the columns in the order segments store them, so it must name
every declared column, in declared
+order. An item written as `<expr> AS <name>` makes that column computed at
ingest time, from the columns it reads:
+
+```sql
+CREATE TABLE "druid"."events" SEALED (
+ tenant VARCHAR,
+ bucket BIGINT,
+ __time TIMESTAMP,
+ user_id BIGINT,
+ payload TYPE('COMPLEX<json>'),
+ PROJECTION __base AS (
+ SELECT tenant, ABS(user_id) % 128 AS bucket, __time, user_id, payload
Review Comment:
I don't think `ABS(user_id) % 128` is valid Druid SQL. (Should use the `MOD`
function.)
##########
sql/src/test/java/org/apache/druid/sql/calcite/CalciteCatalogDdlTest.java:
##########
@@ -0,0 +1,1089 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you 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 org.apache.druid.sql.calcite;
+
+import com.google.common.collect.ImmutableList;
+import com.google.common.collect.ImmutableMap;
+import com.google.common.collect.ImmutableSet;
+import org.apache.calcite.sql.SqlExplain;
+import org.apache.calcite.sql.SqlExplainFormat;
+import org.apache.calcite.sql.SqlExplainLevel;
+import org.apache.calcite.sql.SqlNode;
+import org.apache.calcite.sql.parser.SqlParserPos;
+import org.apache.druid.catalog.model.ClusteredValueGroupsBaseTableMetadata;
+import org.apache.druid.catalog.model.ColumnSpec;
+import org.apache.druid.catalog.model.DatasourceBaseTableMetadata;
+import org.apache.druid.catalog.model.DatasourceProjectionMetadata;
+import org.apache.druid.catalog.model.TableId;
+import org.apache.druid.catalog.model.TableMetadata;
+import org.apache.druid.catalog.model.TableSpec;
+import org.apache.druid.catalog.model.table.ClusterKeySpec;
+import org.apache.druid.catalog.model.table.DatasourceDefn;
+import org.apache.druid.data.input.impl.AggregateProjectionSpec;
+import org.apache.druid.error.DruidException;
+import org.apache.druid.java.util.common.granularity.Granularities;
+import org.apache.druid.server.security.Action;
+import org.apache.druid.server.security.AuthConfig;
+import org.apache.druid.server.security.AuthenticationResult;
+import org.apache.druid.server.security.ForbiddenException;
+import org.apache.druid.server.security.Resource;
+import org.apache.druid.server.security.ResourceAction;
+import org.apache.druid.server.security.ResourceType;
+import org.apache.druid.sql.DirectStatement;
+import org.apache.druid.sql.SqlQueryPlus;
+import
org.apache.druid.sql.calcite.CalciteCatalogDdlTest.CatalogDdlComponentSupplier;
+import org.apache.druid.sql.calcite.parser.DruidSqlParser;
+import org.apache.druid.sql.calcite.planner.CatalogTableWriter;
+import org.apache.druid.sql.calcite.planner.DruidPlanner;
+import org.apache.druid.sql.calcite.planner.PlannerConfig;
+import org.apache.druid.sql.calcite.planner.PlannerFactory;
+import org.apache.druid.sql.calcite.util.CalciteTests;
+import
org.apache.druid.sql.calcite.util.SqlTestFramework.StandardComponentSupplier;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import javax.annotation.Nullable;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * Tests that catalog DDL statements plan into the catalog operations they
claim to, using a writer that records
+ * calls instead of contacting a Coordinator.
+ */
[email protected](CatalogDdlComponentSupplier.class)
+public class CalciteCatalogDdlTest extends BaseCalciteQueryTest
+{
+ private static final RecordingCatalogTableWriter WRITER = new
RecordingCatalogTableWriter();
+
+ public static class CatalogDdlComponentSupplier extends
StandardComponentSupplier
+ {
+ public CatalogDdlComponentSupplier(TempDirProducer tempFolderProducer)
+ {
+ super(tempFolderProducer);
+ }
+
+ @Override
+ public CatalogTableWriter createCatalogTableWriter()
+ {
+ return WRITER;
+ }
+ }
+
+ @BeforeEach
+ public void resetWriter()
+ {
+ WRITER.reset();
+ }
+
+ @Test
+ public void testCreateTable()
+ {
+ execute("CREATE TABLE tbl (__time TIMESTAMP, page VARCHAR, cnt BIGINT)");
+
+ assertEquals(1, WRITER.calls.size());
+ final RecordingCatalogTableWriter.Call call = WRITER.calls.get(0);
+ assertEquals("createTable", call.operation);
+ assertEquals(TableId.datasource("tbl"), call.tableId);
+ assertEquals(DatasourceDefn.TABLE_TYPE, call.spec.type());
+ assertEquals(ImmutableMap.of(), call.spec.properties());
+ assertEquals(
+ ImmutableList.of(
+ new ColumnSpec("__time", "TIMESTAMP", null),
+ new ColumnSpec("page", "VARCHAR", null),
+ new ColumnSpec("cnt", "BIGINT", null)
+ ),
+ call.spec.columns()
+ );
+ assertFalse(call.ifNotExists);
+ assertFalse(call.replace);
+ }
+
+ @Test
+ public void testCreateTableWithPartitioningAndClustering()
+ {
+ execute("CREATE TABLE tbl (page VARCHAR, cnt BIGINT) PARTITIONED BY DAY
CLUSTERED BY page, cnt");
+
+ final TableSpec spec = WRITER.calls.get(0).spec;
+ assertEquals("P1D",
spec.properties().get(DatasourceDefn.SEGMENT_GRANULARITY_PROPERTY));
+ assertEquals(
+ ImmutableList.of(new ClusterKeySpec("page", false), new
ClusterKeySpec("cnt", false)),
+ spec.properties().get(DatasourceDefn.CLUSTER_KEYS_PROPERTY)
+ );
+ }
+
+ @Test
+ public void testCreateTablePartitionedByAll()
+ {
+ execute("CREATE TABLE tbl (page VARCHAR) PARTITIONED BY ALL TIME");
+ assertEquals("ALL",
WRITER.calls.get(0).spec.properties().get(DatasourceDefn.SEGMENT_GRANULARITY_PROPERTY));
+ }
+
+ @Test
+ public void testCreateTableTypeCanonicalization()
+ {
+ execute(
+ "CREATE TABLE tbl (a CHAR, b INTEGER, c REAL, d DOUBLE, e VARCHAR
ARRAY, f TYPE('complex<json>'))"
+ );
+ assertEquals(
+ ImmutableList.of(
+ new ColumnSpec("a", "VARCHAR", null),
+ new ColumnSpec("b", "BIGINT", null),
+ new ColumnSpec("c", "FLOAT", null),
+ new ColumnSpec("d", "DOUBLE", null),
+ new ColumnSpec("e", "VARCHAR ARRAY", null),
+ new ColumnSpec("f", "COMPLEX<json>", null)
+ ),
+ WRITER.calls.get(0).spec.columns()
+ );
+ }
+
+ @Test
+ public void testCreateTableFlags()
+ {
+ execute("CREATE OR REPLACE TABLE tbl (a VARCHAR)");
+ assertTrue(WRITER.calls.get(0).replace);
+
+ WRITER.reset();
+ execute("CREATE TABLE IF NOT EXISTS tbl (a VARCHAR)");
+ assertTrue(WRITER.calls.get(0).ifNotExists);
+ }
+
+ @Test
+ public void testCreateTableInDruidSchema()
+ {
+ execute("CREATE TABLE druid.tbl (a VARCHAR)");
+ assertEquals(TableId.datasource("tbl"), WRITER.calls.get(0).tableId);
+ }
+
+ @Test
+ public void testResourceActionsAreDatasourceReadAndWrite()
+ {
+ // READ and WRITE, matching the Coordinator catalog API's requirement for
write operations: the write is forwarded
+ // with escalated credentials, so this statement-level check is the only
one the issuing user faces.
+ final DirectStatement stmt = statement("CREATE TABLE tbl (a VARCHAR)");
+ stmt.execute();
+ assertEquals(
+ ImmutableSet.of(
+ new ResourceAction(new Resource("tbl", ResourceType.DATASOURCE),
Action.READ),
+ new ResourceAction(new Resource("tbl", ResourceType.DATASOURCE),
Action.WRITE)
+ ),
+ ImmutableSet.copyOf(stmt.resources())
+ );
+ }
+
+ /**
+ * The grammar cannot express EXPLAIN of a DDL statement, so the planner's
defensive guard is exercised directly:
+ * parse the CREATE, wrap it in a hand-built SqlExplain, and hand that to
the planner as if the parser had produced
+ * it.
+ */
+ @Test
+ public void testExplainOfDdlIsRejectedByThePlanner()
+ {
+ final SqlNode create = DruidSqlParser.parse("CREATE TABLE tbl (a
VARCHAR)", true).getMainStatement();
+ final SqlExplain explain = new SqlExplain(
+ SqlParserPos.ZERO,
+ create,
+ SqlExplainLevel.ALL_ATTRIBUTES.symbol(SqlParserPos.ZERO),
+ SqlExplain.Depth.PHYSICAL.symbol(SqlParserPos.ZERO),
+ SqlExplainFormat.TEXT.symbol(SqlParserPos.ZERO),
+ 0
+ );
+ final PlannerFactory plannerFactory = queryFramework()
+
.plannerFixture(PlannerConfig.builder().enableCatalogDdl(true).build(), new
AuthConfig())
+ .plannerFactory();
+ try (DruidPlanner planner = plannerFactory.createPlanner(
+ queryFramework().engine(),
+ "EXPLAIN PLAN FOR CREATE TABLE tbl (a VARCHAR)",
+ explain,
+ CalciteTests.SUPER_USER_AUTH_RESULT,
+ Collections.emptySet(),
+ Collections.emptyMap(),
+ null
+ )) {
+ final DruidException e = assertThrows(DruidException.class,
planner::validate);
+ assertTrue(e.getMessage().contains("EXPLAIN is not supported for
[CREATE_TABLE]"), e.getMessage());
+ }
+ assertTrue(WRITER.calls.isEmpty());
+ }
+
+ /**
+ * An unauthorized statement must not touch the catalog: authorization runs
before planning, and the write itself is
+ * deferred to the result's run, so a rejected statement leaves no trace.
Covers both denial modes: no permissions at
+ * all, and READ without WRITE (DDL requires both).
+ */
+ @Test
+ public void testUnauthorizedDdlNeverReachesTheWriter()
+ {
+ assertThrows(
+ ForbiddenException.class,
+ () -> statement("CREATE TABLE forbiddenDatasource (a VARCHAR)",
CalciteTests.REGULAR_USER_AUTH_RESULT)
+ .execute()
+ );
+ assertThrows(
+ ForbiddenException.class,
+ () -> statement("CREATE TABLE readOnlyTbl (a VARCHAR)",
CalciteTests.REGULAR_USER_AUTH_RESULT).execute()
+ );
+ assertTrue(WRITER.calls.isEmpty());
+ }
+
+ /**
+ * Planning a DDL statement is free of side effects; the catalog write
happens when the result runs.
+ */
+ @Test
+ public void testWriteHappensOnRunNotOnPlan()
+ {
+ final DirectStatement.ResultSet resultSet = statement("CREATE TABLE tbl (a
VARCHAR)").plan();
+ assertTrue(WRITER.calls.isEmpty());
+
+ resultSet.run();
+ assertEquals(1, WRITER.calls.size());
+ assertEquals("createTable", WRITER.calls.get(0).operation);
+ }
+
+ @Test
+ public void testDdlReturnsNoRows()
+ {
+ final DirectStatement stmt = statement("CREATE TABLE tbl (a VARCHAR)");
+ final List<Object[]> results = stmt.execute().getResults().toList();
+ assertEquals(ImmutableList.of(), results);
+ }
+
+ @Test
+ public void testAlterTableAddColumn()
+ {
+ // ADD and ALTER differ only in which existence outcome is an error, and
that rule is enforced inside the
+ // Coordinator's update transaction, so all the statement does is pick the
verb. EditorTest covers the rules.
+ execute("ALTER TABLE tbl ADD COLUMN b BIGINT");
+
+ final RecordingCatalogTableWriter.Call call =
WRITER.lastCall("addColumns");
+ assertEquals(ImmutableList.of(new ColumnSpec("b", "BIGINT", null)),
call.columns);
+ }
+
+ @Test
+ public void testAlterTableDropColumn()
+ {
+ execute("ALTER TABLE tbl DROP COLUMN gone");
+ assertEquals(ImmutableList.of("gone"),
WRITER.lastCall("dropColumns").droppedColumns);
+ }
+
+ @Test
+ public void testAlterTableAlterColumn()
+ {
+ execute("ALTER TABLE tbl ALTER COLUMN cnt SET DATA TYPE DOUBLE");
+ assertEquals(
+ ImmutableList.of(new ColumnSpec("cnt", "DOUBLE", null)),
+ WRITER.lastCall("alterColumns").columns
+ );
+ }
+
+ @Test
+ public void testAlterTableSetProperties()
+ {
+ execute("ALTER TABLE tbl SET PROPERTIES (targetSegmentRows = 3000000,
sealed = TRUE, description = 'hi')");
+
+ final Map<String, Object> properties =
WRITER.lastCall("updateProperties").properties;
+ assertEquals(3000000L, properties.get("targetSegmentRows"));
+ assertEquals(true, properties.get("sealed"));
+ assertEquals("hi", properties.get("description"));
+ }
+
+ @Test
+ public void testAlterTableSetPropertyToNullRemovesIt()
+ {
+ execute("ALTER TABLE tbl SET PROPERTIES (description = NULL)");
+ final Map<String, Object> properties =
WRITER.lastCall("updateProperties").properties;
+ assertTrue(properties.containsKey("description"));
+ assertNull(properties.get("description"));
+ }
+
+ @Test
+ public void testCreateTableRejectsDuplicateColumn()
+ {
+ final DruidException e = assertThrows(
+ DruidException.class,
+ () -> execute("CREATE TABLE tbl (a VARCHAR, a BIGINT)")
+ );
+ assertTrue(e.getMessage().contains("Column [a] is declared more than
once"));
+ }
+
+ @Test
+ public void testCreateTableRejectsUnsupportedType()
+ {
+ final DruidException e = assertThrows(
+ DruidException.class,
+ () -> execute("CREATE TABLE tbl (a TYPE('NOT_A_TYPE'))")
+ );
+ assertTrue(e.getMessage().contains("unsupported type"));
+ }
+
+ /**
+ * Any spelling that resolves to a LONG is accepted for the time column, and
is stored as written.
+ */
+ @Test
+ public void testCreateTableTimeColumnSpellings()
+ {
+ execute("CREATE TABLE tbl (__time BIGINT)");
+ assertEquals(ImmutableList.of(new ColumnSpec("__time", "BIGINT", null)),
WRITER.calls.get(0).spec.columns());
+
+ WRITER.reset();
+ execute("CREATE TABLE tbl (__time TYPE('LONG'))");
+ assertEquals(ImmutableList.of(new ColumnSpec("__time", "LONG", null)),
WRITER.calls.get(0).spec.columns());
+ }
+
+ @Test
+ public void testCreateTableRejectsNonLongTimeColumn()
+ {
+ final DruidException e = assertThrows(
+ DruidException.class,
+ () -> execute("CREATE TABLE tbl (__time VARCHAR)")
+ );
+ assertTrue(e.getMessage().contains("Column [__time] must have type"));
+ }
+
+ @Test
+ public void testCreateTableRejectsNonDruidSchema()
+ {
+ final DruidException e = assertThrows(
+ DruidException.class,
+ () -> execute("CREATE TABLE lookup.tbl (a VARCHAR)")
+ );
+ assertTrue(e.getMessage().contains("is not a Druid datasource"));
+ }
+
+ @Test
+ public void testCreateTableRejectsBothReplaceAndIfNotExists()
+ {
+ final DruidException e = assertThrows(
+ DruidException.class,
+ () -> execute("CREATE OR REPLACE TABLE IF NOT EXISTS tbl (a VARCHAR)")
+ );
+ assertTrue(e.getMessage().contains("Cannot specify both OR REPLACE and IF
NOT EXISTS"));
+ }
+
+ @Test
+ public void testCreateTableRejectsClusteringExpression()
+ {
+ final DruidException e = assertThrows(
+ DruidException.class,
+ () -> execute("CREATE TABLE tbl (a VARCHAR) CLUSTERED BY a DESC")
+ );
+ assertTrue(e.getMessage().contains("must be a column name"));
+ }
+
+ /**
+ * The feature is off unless an operator turns it on, so that upgrading a
cluster does not silently widen what a
+ * datasource WRITE permission allows.
+ */
+ @Test
+ public void testDdlIsDisabledByDefault()
+ {
+ final DirectStatement stmt =
getSqlStatementFactory(PlannerConfig.builder().build(), new AuthConfig())
+ .directStatement(
+ SqlQueryPlus.builder("CREATE TABLE tbl (a VARCHAR)")
+ .auth(CalciteTests.SUPER_USER_AUTH_RESULT)
+ .build()
+ );
+ final DruidException e = assertThrows(DruidException.class, stmt::execute);
+ assertTrue(e.getMessage().contains("druid.sql.planner.enableCatalogDdl"),
e.getMessage());
+ assertEquals(ImmutableList.of(), WRITER.calls);
+ }
+
+ /**
+ * The stored specification must be the one the planner would produce for
the equivalent query, since that is what
+ * makes a projection match at query time. Pinned as JSON so a change in
planner output is visible here.
+ */
+ @Test
+ public void testCreateTableWithProjection() throws Exception
+ {
+ execute(
+ "CREATE TABLE tbl (__time TIMESTAMP, page VARCHAR, cnt BIGINT,"
+ + " PROJECTION daily AS (SELECT TIME_FLOOR(__time, 'P1D'), page,
SUM(cnt) AS total GROUP BY 1, 2))"
+ );
+
+ assertEquals(
+ "[{\"spec\":{\"type\":\"aggregate\",\"name\":\"daily\","
+ + "\"virtualColumns\":[{\"type\":\"expression\",\"name\":\"v0\","
+ +
"\"expression\":\"timestamp_floor(\\\"__time\\\",'P1D',null,'UTC')\",\"outputType\":\"LONG\"}],"
+ +
"\"groupingColumns\":[{\"type\":\"long\",\"name\":\"v0\",\"multiValueHandling\":\"SORTED_ARRAY\","
+ +
"\"createBitmapIndex\":false},{\"type\":\"string\",\"name\":\"page\","
+ +
"\"multiValueHandling\":\"SORTED_ARRAY\",\"createBitmapIndex\":true}],"
+ +
"\"aggregators\":[{\"type\":\"longSum\",\"name\":\"total\",\"fieldName\":\"cnt\"}],"
+ + "\"ordering\":[{\"columnName\":\"v0\",\"order\":\"ascending\"},"
+ + "{\"columnName\":\"page\",\"order\":\"ascending\"}]}}]",
+ projectionsJson()
+ );
+ }
+
+ /**
+ * A projection body is planned under the statement's own context, so a SET
clause that changes how the equivalent
+ * query would plan changes the stored definition the same way. Here the
session time zone reaches the TIME_FLOOR.
+ */
+ @Test
+ public void testProjectionBodyHonorsStatementContext() throws Exception
+ {
+ execute(
+ "SET sqlTimeZone = 'America/Los_Angeles';\n"
+ + "CREATE TABLE tbl (__time TIMESTAMP, page VARCHAR, cnt BIGINT,"
+ + " PROJECTION daily AS (SELECT TIME_FLOOR(__time, 'P1D'), page,
SUM(cnt) AS total GROUP BY 1, 2))"
+ );
+
+ assertTrue(
+
projectionsJson().contains("timestamp_floor(\\\"__time\\\",'P1D',null,'America/Los_Angeles')"),
+ projectionsJson()
+ );
+ }
+
+ /**
+ * The overrides the lift depends on are applied on top of the statement's
context, so a SET clause cannot put the
+ * planner into a shape the lift does not understand.
+ */
+ @Test
+ public void testProjectionBodyContextCannotOverrideDeterministicOverrides()
throws Exception
+ {
+ execute(
+ "SET sqlUseGranularity = TRUE;\n"
+ + "CREATE TABLE tbl (__time TIMESTAMP, page VARCHAR, cnt BIGINT,"
+ + " PROJECTION daily AS (SELECT TIME_FLOOR(__time, 'P1D'), page,
SUM(cnt) AS total GROUP BY 1, 2))"
+ );
+
+ // Still lifted as an ordinary grouping column rather than a query
granularity, exactly as without the SET.
+ assertTrue(
+
projectionsJson().contains("timestamp_floor(\\\"__time\\\",'P1D',null,'UTC')"),
+ projectionsJson()
+ );
+ }
+
+ /**
+ * A projection defined with TIME_FLOOR must carry a granularity the segment
layer can recover, which is how the
+ * projection gets matched to time-grouped queries.
+ */
+ @Test
+ public void testProjectionGranularityIsRecoverable()
+ {
+ execute(
+ "CREATE TABLE tbl (__time TIMESTAMP, page VARCHAR, cnt BIGINT,"
+ + " PROJECTION hourly AS (SELECT TIME_FLOOR(__time, 'PT1H'), page,
SUM(cnt) AS total GROUP BY 1, 2))"
+ );
+
+ final AggregateProjectionSpec spec = projection(0).getSpec();
+ final String timeColumn = spec.toMetadataSchema().getTimeColumnName();
+ assertEquals("v0", timeColumn);
+ assertEquals(
+ Granularities.HOUR,
+
Granularities.fromVirtualColumn(spec.getVirtualColumns().getVirtualColumn(timeColumn))
+ );
+ }
+
+ @Test
+ public void testProjectionWithFilter()
+ {
+ execute(
+ "CREATE TABLE tbl (__time TIMESTAMP, page VARCHAR, cnt BIGINT,"
+ + " PROJECTION filtered AS (SELECT page, SUM(cnt) AS total WHERE page
<> 'skip' GROUP BY page))"
+ );
+ assertEquals("!page = skip",
projection(0).getSpec().getFilter().toString());
+ }
+
+ /**
+ * A time bound written in the body is moved into the query's intervals
during planning, and has to be put back:
+ * a projection stores a filter, not an interval.
+ */
+ @Test
+ public void testProjectionWithTimeFilter()
+ {
+ execute(
+ "CREATE TABLE tbl (__time TIMESTAMP, page VARCHAR, cnt BIGINT,"
+ + " PROJECTION recent AS (SELECT page, SUM(cnt) AS total"
+ + " WHERE __time >= TIMESTAMP '2020-01-01 00:00:00' GROUP BY page))"
+ );
+ assertNotNull(projection(0).getSpec().getFilter(), "time filter must
survive as a filter");
Review Comment:
Isn't it bad that time filters are retained as filters? Typically when a
real SQL query is actually issued, `__time` filters will be moved to
`intervals`, so they will not actually show up in the `filter`. So, does that
mean they won't be able to match the stored projection? Or, is there something
that causes this to be handled well?
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]