diff --git a/docs/development/extensions-core/catalog.md b/docs/development/extensions-core/catalog.md index 20842c927ae4..01e6202d59c7 100644 --- a/docs/development/extensions-core/catalog.md +++ b/docs/development/extensions-core/catalog.md @@ -196,11 +196,68 @@ 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. +pre-aggregation. Its body lists the columns in the order segments store them, so it must name every declared column. +The shape of the body chooses the layout: a bare select declares a plain table, `GROUP BY` declares a +[rollup](../../ingestion/rollup.md) table, and `CLUSTERED BY` declares a clustered table. Whatever the layout, a +`__base` body cannot filter as the base table represents every ingested row. -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 ` AS ` makes that column computed at ingest time, from the columns it reads: +##### Plain tables + +By default, a table stores rows as they arrive, sorted by `__time`. Declaring a plain `__base` projection makes the +declared column order the segment sort order, and allows declaring a query granularity. The only expression it +accepts is `TIME_FLOOR(__time, )` selected as `__time`, which declares the table's query granularity (the +granularity ingested timestamps are floored to). A query granularity must be a period granularity in the UTC time zone, +matching the UTC-stored `__time` column: + +```sql +CREATE TABLE "druid"."visits" ( + user_id VARCHAR, + __time TIMESTAMP, + pages_visited BIGINT, + PROJECTION __base AS ( + SELECT user_id, TIME_FLOOR(__time, 'PT1H') AS __time, pages_visited + ) +) +PARTITIONED BY DAY +``` + +##### Rollup tables + +A rollup table aggregates rows with identical grouping values into a single row at ingest time, expressed as a +`GROUP BY` query. The query divides the table's declared columns into grouping columns and metric columns. A column +written as an aggregate is a metric column, and every other column is a grouping column, selected as-is. The time +column may instead be grouped as `TIME_FLOOR(__time, )`, which floors rows to that granularity before grouping +and declares the table's query granularity. Each aggregate is named for the metric column it fills, and reads that +same column. In the example below, ingestion supplies `total` values and the table sums them: + +```sql +CREATE TABLE "druid"."visits" ( + user_id VARCHAR, + __time TIMESTAMP, + total BIGINT, + PROJECTION __base AS ( + SELECT user_id, TIME_FLOOR(__time, 'PT1H') AS __time, SUM(total) AS total + GROUP BY 1, 2 + ) +) +PARTITIONED BY DAY +``` + +The metric columns must be declared after the grouping columns, because the declared order is the physical order. A +rollup body accepts the same +[aggregation functions that are supported for rollup at ingestion time](../../multi-stage-query/concepts.md#rollup) +and, like an aggregate projection, cannot compute expressions over aggregates. The body describes a transformation of +the table's own columns (the query that rolls the table up onto itself) not of whatever input first produced the +rows. This is what lets the table aggregate its rows again later, for example when compaction merges segments. An +aggregate whose output it cannot re-aggregate, such as `COUNT(*)`, is rejected: two counts combine by adding, so store +a count by summing a count column, as in `SUM(cnt) AS cnt`. A body that groups without any +aggregates collapses duplicate rows. A rollup table cannot also declare `CLUSTERED BY`. + +##### Clustered tables + +A clustered table stores the rows of each segment internally partitioned by the clustering columns, which must be the +leading columns of the table. An item written as ` AS ` makes that column computed at ingest time, from the +columns it reads: ```sql CREATE TABLE "druid"."events" SEALED ( @@ -217,15 +274,13 @@ CREATE TABLE "druid"."events" SEALED ( PARTITIONED BY DAY ``` -The clustering columns must be the leading columns of the table, because the declared order is the physical order. -`SEALED` is optional: a column the query produces but the table does not declare is stored after the declared -layout, in the order it arrives. Declare `SEALED` to reject such columns instead. - Computed columns, like `bucket` in the example above, are computed based on inputs provided by the `INSERT` or -`REPLACE`. In terms of the example, the `INSERT` or `REPLACE` command should provide `user_id`, not `bucket`. +`REPLACE`. In terms of the example, the `INSERT` or `REPLACE` command should provide `user_id`, not `bucket`. Only a +clustered table materializes computed columns. -Unlike an aggregate projection, a `__base` body cannot filter or group: the base table stores every ingested row. -It is the only projection that chooses a clustering. +For every layout, `SEALED` is optional: a column the ingestion query produces but the table does not declare is +stored after the declared layout, in the order it arrives (for a rollup table, as an additional grouping column). +Declare `SEALED` to reject such columns instead. `ALTER TABLE ... ADD PROJECTION __base AS ( ... )` gives an existing table a layout, and `ALTER TABLE ... DROP PROJECTION __base` removes it. Both affect future segments only. diff --git a/embedded-tests/src/test/java/org/apache/druid/testing/embedded/catalog/CatalogDdlAndIngestTest.java b/embedded-tests/src/test/java/org/apache/druid/testing/embedded/catalog/CatalogDdlAndIngestTest.java index b280a7a15a02..d20893891e60 100644 --- a/embedded-tests/src/test/java/org/apache/druid/testing/embedded/catalog/CatalogDdlAndIngestTest.java +++ b/embedded-tests/src/test/java/org/apache/druid/testing/embedded/catalog/CatalogDdlAndIngestTest.java @@ -456,6 +456,101 @@ public void testCreateClusteredBaseTableWithComposedComputedColumn() ); } + /** + * A {@code __base} projection with {@code GROUP BY} declares the rollup layout: the declared columns split into + * grouping columns and (by aggregator name) metric columns, {@code TIME_FLOOR(__time, )} declares the query + * granularity, and rows with identical grouping values after flooring are aggregated into one at ingest time. + */ + @Test + public void testCreateRollupBaseTableThenIngestAndQuery() + { + final String tableName = dataSource; + + cluster.callApi().runSql( + "CREATE TABLE \"%s\" (\n" + + " tenant VARCHAR,\n" + + " __time TIMESTAMP,\n" + + " total BIGINT,\n" + + " PROJECTION __base AS (\n" + + " SELECT tenant, TIME_FLOOR(__time, 'PT1H') AS __time, SUM(total) AS total\n" + + " GROUP BY 1, 2\n" + + " )\n" + + ")\n" + + "PARTITIONED BY DAY", + tableName + ); + + ingest( + "INSERT INTO \"%s\"\n" + + "SELECT b AS tenant, TIME_FLOOR(TIME_PARSE(a), 'PT1H') AS __time, SUM(c) AS total\n" + + "FROM TABLE(\n" + + " EXTERN(\n" + + " '{\"type\":\"inline\",\"data\":\"2022-12-26T05:10:00,bbb,4" + + "\\n2022-12-26T05:20:00,aaa,1\\n2022-12-26T05:40:00,aaa,2\\n2022-12-26T09:10:00,aaa,8\"}',\n" + + " '{\"type\":\"csv\",\"findColumnsFromHeader\":false,\"columns\":[\"a\",\"b\",\"c\"]}'\n" + + " )\n" + + ") EXTEND (a VARCHAR, b VARCHAR, c BIGINT)\n" + + "GROUP BY 1, 2", + tableName + ); + + // The two same-hour 'aaa' rows collapsed into one with total summed; rows come back in grouping order. + cluster.callApi().verifySqlQuery( + "SELECT * FROM %s", + tableName, + "aaa,2022-12-26T05:00:00.000Z,3\n" + + "aaa,2022-12-26T09:00:00.000Z,8\n" + + "bbb,2022-12-26T05:00:00.000Z,4" + ); + } + + /** + * The canonical count workflow for a rollup table: the table cannot declare {@code COUNT(*)} (an aggregator must + * combine its own output), so it declares {@code SUM(cnt) AS cnt} and the ingestion query computes the raw count — + * {@code COUNT(*) AS cnt} — which the table then sums on any re-aggregation. + */ + @Test + public void testCreateRollupBaseTableThenIngestCounts() + { + final String tableName = dataSource; + + cluster.callApi().runSql( + "CREATE TABLE \"%s\" (\n" + + " tenant VARCHAR,\n" + + " __time TIMESTAMP,\n" + + " cnt BIGINT,\n" + + " PROJECTION __base AS (\n" + + " SELECT tenant, TIME_FLOOR(__time, 'PT1H') AS __time, SUM(cnt) AS cnt\n" + + " GROUP BY 1, 2\n" + + " )\n" + + ")\n" + + "PARTITIONED BY DAY", + tableName + ); + + ingest( + "INSERT INTO \"%s\"\n" + + "SELECT b AS tenant, TIME_FLOOR(TIME_PARSE(a), 'PT1H') AS __time, COUNT(*) AS cnt\n" + + "FROM TABLE(\n" + + " EXTERN(\n" + + " '{\"type\":\"inline\",\"data\":\"2022-12-26T05:10:00,bbb" + + "\\n2022-12-26T05:20:00,aaa\\n2022-12-26T05:40:00,aaa\\n2022-12-26T09:10:00,aaa\"}',\n" + + " '{\"type\":\"csv\",\"findColumnsFromHeader\":false,\"columns\":[\"a\",\"b\"]}'\n" + + " )\n" + + ") EXTEND (a VARCHAR, b VARCHAR)\n" + + "GROUP BY 1, 2", + tableName + ); + + cluster.callApi().verifySqlQuery( + "SELECT * FROM %s", + tableName, + "aaa,2022-12-26T05:00:00.000Z,2\n" + + "aaa,2022-12-26T09:00:00.000Z,1\n" + + "bbb,2022-12-26T05:00:00.000Z,1" + ); + } + /** * A {@code __base} projection without {@code CLUSTERED BY} declares the plain-table layout: declared column order is * the segment storage and sort order, and {@code TIME_FLOOR(__time, )} declares the table's query diff --git a/embedded-tests/src/test/java/org/apache/druid/testing/embedded/indexing/IndexTaskTest.java b/embedded-tests/src/test/java/org/apache/druid/testing/embedded/indexing/IndexTaskTest.java index 2b63fde274a7..617dc30dac6b 100644 --- a/embedded-tests/src/test/java/org/apache/druid/testing/embedded/indexing/IndexTaskTest.java +++ b/embedded-tests/src/test/java/org/apache/druid/testing/embedded/indexing/IndexTaskTest.java @@ -20,6 +20,7 @@ package org.apache.druid.testing.embedded.indexing; import org.apache.druid.data.input.impl.LongDimensionSchema; +import org.apache.druid.data.input.impl.RollupTableProjectionSpec; import org.apache.druid.data.input.impl.StringDimensionSchema; import org.apache.druid.data.input.impl.TableProjectionSpec; import org.apache.druid.indexer.granularity.SegmentGranularitySpec; @@ -29,6 +30,7 @@ import org.apache.druid.java.util.common.StringUtils; import org.apache.druid.java.util.common.granularity.Granularities; import org.apache.druid.java.util.common.guava.Comparators; +import org.apache.druid.query.aggregation.LongSumAggregatorFactory; import org.apache.druid.testing.embedded.EmbeddedBroker; import org.apache.druid.testing.embedded.EmbeddedClusterApis; import org.apache.druid.testing.embedded.EmbeddedCoordinator; @@ -162,6 +164,56 @@ public void test_runIndexTask_withBaseTableSpec() ); } + /** + * Batch ingestion driven by a {@link RollupTableProjectionSpec} base table: the declared columns are the grouping + * columns (item before __time — declared order is the sort order), the query-granularity carrier floors + * {@code __time} to the hour, and rows with identical grouping values after flooring are aggregated into one by the + * spec's aggregators. + */ + @Test + @Timeout(60) + public void test_runIndexTask_withRollupBaseTableSpec() + { + final String taskId = EmbeddedClusterApis.newTaskId(dataSource); + final RollupTableProjectionSpec baseTable = RollupTableProjectionSpec + .builder() + .groupingColumns( + new StringDimensionSchema("item"), + new LongDimensionSchema("__time") + ) + .aggregators(new LongSumAggregatorFactory("total", "total")) + .build() + .withQueryGranularity(Granularities.HOUR); + final IndexTask task = TaskBuilder + .ofTypeIndex() + .isoTimestampColumn("time") + .csvInputFormatWithColumns("time", "item", "total") + .inlineInputSourceWithData( + "2025-06-01T05:10:00.000Z,apple,1" + + "\n2025-06-01T05:20:00.000Z,apple,2" + + "\n2025-06-01T05:30:00.000Z,banana,4" + + "\n2025-06-01T09:10:00.000Z,apple,8" + ) + .dataSchema( + schema -> schema + .withBaseTable(baseTable) + .withSegmentGranularity(new SegmentGranularitySpec(Granularities.DAY, null)) + ) + .dataSource(dataSource) + .withId(taskId); + + cluster.callApi().runTask(task, overlord); + cluster.callApi().waitForAllSegmentsToBeAvailable(dataSource, coordinator, broker); + + // The two 5-o'clock apple rows collapsed into one with total summed; scan order is the declared grouping order. + Assertions.assertEquals( + "apple,2025-06-01T05:00:00.000Z,3" + + "\napple,2025-06-01T09:00:00.000Z,8" + + "\nbanana,2025-06-01T05:00:00.000Z,4", + cluster.runSql("SELECT item, __time, total FROM %s", dataSource) + ); + } + @Test @Timeout(20) public void test_run10Tasks_concurrently() diff --git a/processing/src/main/java/org/apache/druid/data/input/impl/AdaptedBaseTableProjectionSpec.java b/processing/src/main/java/org/apache/druid/data/input/impl/AdaptedBaseTableProjectionSpec.java index 2828d2cda533..f2f47492a124 100644 --- a/processing/src/main/java/org/apache/druid/data/input/impl/AdaptedBaseTableProjectionSpec.java +++ b/processing/src/main/java/org/apache/druid/data/input/impl/AdaptedBaseTableProjectionSpec.java @@ -112,6 +112,12 @@ public Granularity getQueryGranularity() return granularitySpec.getQueryGranularity(); } + @Override + public boolean isRollup() + { + return granularitySpec.isRollup(); + } + @Override public BaseTableProjectionSpec withQueryGranularity(@Nullable Granularity queryGranularity) { diff --git a/processing/src/main/java/org/apache/druid/data/input/impl/BaseTableProjectionSpec.java b/processing/src/main/java/org/apache/druid/data/input/impl/BaseTableProjectionSpec.java index e9d42175e4c3..6f488bd69eda 100644 --- a/processing/src/main/java/org/apache/druid/data/input/impl/BaseTableProjectionSpec.java +++ b/processing/src/main/java/org/apache/druid/data/input/impl/BaseTableProjectionSpec.java @@ -21,6 +21,7 @@ import com.fasterxml.jackson.annotation.JsonSubTypes; import com.fasterxml.jackson.annotation.JsonTypeInfo; +import com.google.common.collect.Sets; import org.apache.druid.error.InvalidInput; import org.apache.druid.java.util.common.granularity.Granularities; import org.apache.druid.java.util.common.granularity.Granularity; @@ -32,8 +33,10 @@ import org.apache.druid.segment.projections.BaseTableProjectionSchema; import javax.annotation.Nullable; +import java.util.ArrayList; import java.util.Collections; import java.util.List; +import java.util.Set; /** * Spec describing the shape of the 'base' table schema for a {@link org.apache.druid.segment.indexing.DataSchema}. This @@ -51,6 +54,7 @@ @JsonTypeInfo(use = JsonTypeInfo.Id.NAME, property = "type") @JsonSubTypes({ @JsonSubTypes.Type(name = TableProjectionSpec.TYPE_NAME, value = TableProjectionSpec.class), + @JsonSubTypes.Type(name = RollupTableProjectionSpec.TYPE_NAME, value = RollupTableProjectionSpec.class), @JsonSubTypes.Type( name = ClusteredValueGroupsBaseTableProjectionSpec.TYPE_NAME, value = ClusteredValueGroupsBaseTableProjectionSpec.class @@ -80,6 +84,14 @@ default Granularity getQueryGranularity() return granularity == null ? Granularities.NONE : granularity; } + /** + * Returns true if rows with identical grouping values are aggregated into one at ingest time. + */ + default boolean isRollup() + { + return false; + } + /** * Returns a copy of this spec with the given query granularity applied (implementation-defined representation). Used * by the ingestion and compaction config paths to attach the query-derived granularity to the operator-supplied spec. @@ -146,4 +158,87 @@ static void validateGranularity(VirtualColumn granularityColumn, String typeName } validateQueryGranularity(granularity, typeName); } + + /** + * Validates a declared column list: non-empty, unique names, an explicit {@link ColumnHolder#TIME_COLUMN_NAME} entry + * to define the time position, and no entry named for the query-granularity carrier (which is a virtual column, not + * a stored column). + */ + static void validateDeclaredColumns( + @Nullable List columns, + String columnsProperty, + String typeName + ) + { + if (columns == null || columns.isEmpty()) { + throw InvalidInput.exception("'%s' must be non-empty for [%s] base table", columnsProperty, typeName); + } + final Set seen = Sets.newHashSetWithExpectedSize(columns.size()); + for (DimensionSchema d : columns) { + if (!seen.add(d.getName())) { + throw InvalidInput.exception("'%s' contains duplicate name [%s]", columnsProperty, d.getName()); + } + } + boolean foundTime = false; + for (DimensionSchema column : columns) { + final String name = column.getName(); + // The query-granularity virtual column is a granularity carrier in virtualColumns (it floors the stored __time + // column); it is not itself a stored column, so it must not be declared as one. + if (Granularities.GRANULARITY_VIRTUAL_COLUMN_NAME.equals(name)) { + throw InvalidInput.exception( + "[%s] is the query-granularity virtual column, not a stored column; declare it in 'virtualColumns' and use" + + " [%s] as the time column in '%s'", + Granularities.GRANULARITY_VIRTUAL_COLUMN_NAME, + ColumnHolder.TIME_COLUMN_NAME, + columnsProperty + ); + } + if (ColumnHolder.TIME_COLUMN_NAME.equals(name)) { + foundTime = true; + } + } + if (!foundTime) { + throw InvalidInput.exception( + "[%s] base table must include [%s] in '%s' to define the time position", + typeName, + ColumnHolder.TIME_COLUMN_NAME, + columnsProperty + ); + } + } + + /** + * Validates that {@code virtualColumns} holds nothing but the {@link Granularities#GRANULARITY_VIRTUAL_COLUMN_NAME} + * carrier, for the layouts whose write path does not evaluate spec virtual columns (any other virtual column would + * be dead metadata whose column never gets materialized; computed columns belong in a {@code transformSpec}). The + * carrier itself is validated by {@link #validateGranularity}. + */ + static void validateGranularityOnlyVirtualColumns(VirtualColumns virtualColumns, String typeName) + { + for (VirtualColumn virtualColumn : virtualColumns.getVirtualColumns()) { + if (!Granularities.GRANULARITY_VIRTUAL_COLUMN_NAME.equals(virtualColumn.getOutputName())) { + throw InvalidInput.exception( + "virtual column [%s] is not supported: a [%s] base table stores only declared columns and does not" + + " materialize virtual columns; use a transformSpec to compute columns at ingest, or the [%s] virtual" + + " column to carry query granularity", + virtualColumn.getOutputName(), + typeName, + Granularities.GRANULARITY_VIRTUAL_COLUMN_NAME + ); + } + validateGranularity(virtualColumn, typeName); + } + } + + /** + * Convert a list of {@link DimensionSchema} into ascending {@link OrderBy}, preserving declared order. + */ + static List declaredOrderAscending(List columns) + { + final List ordering = new ArrayList<>(columns.size()); + for (DimensionSchema d : columns) { + ordering.add(OrderBy.ascending(d.getName())); + } + return Collections.unmodifiableList(ordering); + } } diff --git a/processing/src/main/java/org/apache/druid/data/input/impl/RollupTableProjectionSpec.java b/processing/src/main/java/org/apache/druid/data/input/impl/RollupTableProjectionSpec.java new file mode 100644 index 000000000000..a5a2b81d5f8d --- /dev/null +++ b/processing/src/main/java/org/apache/druid/data/input/impl/RollupTableProjectionSpec.java @@ -0,0 +1,354 @@ +/* + * 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.data.input.impl; + +import com.fasterxml.jackson.annotation.JsonCreator; +import com.fasterxml.jackson.annotation.JsonIgnore; +import com.fasterxml.jackson.annotation.JsonInclude; +import com.fasterxml.jackson.annotation.JsonProperty; +import com.fasterxml.jackson.annotation.JsonTypeName; +import org.apache.druid.error.InvalidInput; +import org.apache.druid.java.util.common.granularity.Granularities; +import org.apache.druid.java.util.common.granularity.Granularity; +import org.apache.druid.query.OrderBy; +import org.apache.druid.query.aggregation.AggregatorFactory; +import org.apache.druid.segment.VirtualColumn; +import org.apache.druid.segment.VirtualColumns; +import org.apache.druid.segment.column.ColumnHolder; + +import javax.annotation.Nullable; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.HashSet; +import java.util.List; +import java.util.Objects; +import java.util.Set; + +/** + * {@link BaseTableProjectionSpec} for a rollup table: the plain {@link TableProjectionSpec} shape plus + * {@link #aggregators}. The declared {@link #groupingColumns} are the grouping columns, rows with identical grouping values + * (after query-granularity flooring of {@code __time}) are aggregated into one at ingest time, with each aggregator + * producing a metric column after the grouping columns. As in the plain layout, the time position is an explicit + * positional entry in {@link #groupingColumns} named {@code __time}, declared column order is the segment storage and sort + * order, and query granularity rides as the {@link Granularities#GRANULARITY_VIRTUAL_COLUMN_NAME} carrier virtual + * column (the only virtual column accepted). + *

+ * Operator-facing counterpart of the segment metadata-side + * {@link org.apache.druid.segment.projections.RollupTableProjectionSchema}. + *

+ * Aggregators may be empty: a rollup table without metrics collapses rows with identical grouping values, mirroring + * the legacy {@code rollup=true} granularity spec without a {@code metricsSpec}. + */ +@JsonTypeName(RollupTableProjectionSpec.TYPE_NAME) +public final class RollupTableProjectionSpec implements BaseTableProjectionSpec +{ + public static final String TYPE_NAME = "rollupTable"; + + private final VirtualColumns virtualColumns; + private final List groupingColumns; + private final AggregatorFactory[] aggregators; + private final DimensionsSpec dimensionsSpec; + private final List ordering; + + public static Builder builder() + { + return new Builder(); + } + + @JsonCreator + public RollupTableProjectionSpec( + @JsonProperty("virtualColumns") @Nullable VirtualColumns virtualColumns, + @JsonProperty("groupingColumns") List groupingColumns, + @JsonProperty("aggregators") @Nullable AggregatorFactory[] aggregators + ) + { + BaseTableProjectionSpec.validateDeclaredColumns(groupingColumns, "groupingColumns", TYPE_NAME); + this.virtualColumns = virtualColumns == null ? VirtualColumns.EMPTY : virtualColumns; + BaseTableProjectionSpec.validateGranularityOnlyVirtualColumns(this.virtualColumns, TYPE_NAME); + this.groupingColumns = Collections.unmodifiableList(new ArrayList<>(groupingColumns)); + this.aggregators = aggregators == null ? new AggregatorFactory[0] : aggregators; + validateAggregators(this.groupingColumns, this.aggregators); + this.dimensionsSpec = DimensionsSpec.builder() + .setDimensions(this.groupingColumns) + .setForceSegmentSortByTime(false) + .build(); + this.ordering = BaseTableProjectionSpec.declaredOrderAscending(this.groupingColumns); + } + + @Override + @JsonProperty + @JsonInclude(JsonInclude.Include.NON_DEFAULT) + public VirtualColumns getVirtualColumns() + { + return virtualColumns; + } + + /** + * The full, ordered list of grouping columns in segment order, with the explicit {@code __time} (or + * query-granularity) marker at its declared position. + */ + @JsonProperty("groupingColumns") + public List getGroupingColumns() + { + return groupingColumns; + } + + /** + * The aggregators computed over the rows collapsed into each grouping tuple, each producing a metric column stored + * after the grouping columns. + */ + @Override + @JsonProperty("aggregators") + @JsonInclude(JsonInclude.Include.NON_EMPTY) + public AggregatorFactory[] getMetrics() + { + return aggregators; + } + + @Override + @JsonIgnore + public List getOrdering() + { + return ordering; + } + + @Override + @JsonIgnore + public DimensionsSpec getDimensionsSpec() + { + return dimensionsSpec; + } + + @Override + public boolean isRollup() + { + return true; + } + + /** + * Returns a copy of this spec with a new {@code queryGranularity}, expressed as a + * {@link Granularities#GRANULARITY_VIRTUAL_COLUMN_NAME} virtual column added to {@link #getVirtualColumns()}. A + * {@code null} or {@code NONE} granularity is a no-op (no flooring), so this returns {@code this} unchanged. + *

+ * {@code ALL} is rejected: it has no granularity virtual column representation, so accepting it would silently + * degrade to {@code NONE} when the spec is read back. + *

+ * Idempotent: if the spec already declares a query-granularity virtual column, that one is authoritative and this is + * a no-op. (The compaction path attaches the virtual column up front; the MSQ generation path then calls this again + * with the query-derived granularity, which must not double-add.) + */ + @Override + public RollupTableProjectionSpec withQueryGranularity(@Nullable Granularity queryGranularity) + { + if (Granularities.ALL.equals(queryGranularity)) { + throw InvalidInput.exception( + "Query granularity[ALL] is not supported for [%s] base tables", + TYPE_NAME + ); + } + if (queryGranularity == null + || Granularities.NONE.equals(queryGranularity) + || virtualColumns.getVirtualColumn(Granularities.GRANULARITY_VIRTUAL_COLUMN_NAME) != null) { + return this; + } + BaseTableProjectionSpec.validateQueryGranularity(queryGranularity, TYPE_NAME); + final VirtualColumn granularityVirtualColumn = + Granularities.toVirtualColumn(queryGranularity, Granularities.GRANULARITY_VIRTUAL_COLUMN_NAME); + final List merged = new ArrayList<>(Arrays.asList(virtualColumns.getVirtualColumns())); + merged.add(granularityVirtualColumn); + return new RollupTableProjectionSpec(VirtualColumns.create(merged), groupingColumns, aggregators); + } + + /** + * Compares rollup-spec state for compaction up-to-date checks: query granularity is compared separately (via its + * carrier virtual column), so it is stripped from both sides; everything else (the grouping columns, the + * aggregators, and any other virtual columns) must match. A spec of a different type is never equivalent. + */ + @Override + public boolean hasEqualCompactionState(BaseTableProjectionSpec other) + { + if (!(other instanceof RollupTableProjectionSpec)) { + return false; + } + return withoutQueryGranularity().equals(((RollupTableProjectionSpec) other).withoutQueryGranularity()); + } + + @Override + public RollupTableProjectionSpec withAdditionalColumns(@Nullable List additionalColumns) + { + if (additionalColumns == null || additionalColumns.isEmpty()) { + return this; + } + final List revised = new ArrayList<>(groupingColumns.size() + additionalColumns.size()); + revised.addAll(groupingColumns); + for (DimensionSchema additionalColumn : additionalColumns) { + if (ColumnHolder.TIME_COLUMN_NAME.equals(additionalColumn.getName())) { + throw InvalidInput.exception( + "Cannot append column [%s] to a [%s] base table; it must be declared at its position in the grouping column list", + ColumnHolder.TIME_COLUMN_NAME, + TYPE_NAME + ); + } + revised.add(additionalColumn); + } + // Duplicates of a declared column or aggregator, and a column named for the query-granularity carrier, are + // rejected by the constructor's validation. + return new RollupTableProjectionSpec(virtualColumns, revised, aggregators); + } + + /** + * Returns a copy of this spec with the {@link Granularities#GRANULARITY_VIRTUAL_COLUMN_NAME} virtual column removed, + * the inverse of {@link #withQueryGranularity(Granularity)}. If no such virtual column is present this returns + * {@code this} unchanged. Used to compare schema independently of query granularity in {@link #hasEqualCompactionState}. + */ + private RollupTableProjectionSpec withoutQueryGranularity() + { + if (virtualColumns.getVirtualColumn(Granularities.GRANULARITY_VIRTUAL_COLUMN_NAME) == null) { + return this; + } + final List remaining = new ArrayList<>(); + for (VirtualColumn vc : virtualColumns.getVirtualColumns()) { + if (!Granularities.GRANULARITY_VIRTUAL_COLUMN_NAME.equals(vc.getOutputName())) { + remaining.add(vc); + } + } + return new RollupTableProjectionSpec(VirtualColumns.create(remaining), groupingColumns, aggregators); + } + + /** + * Aggregators produce the metric columns, so their names share a namespace with the grouping columns — a collision + * means one column with two definitions — and with the query-granularity carrier: virtual columns shadow stored + * columns when selectors resolve, so a metric named for the carrier would be unreadable once a query granularity is + * attached. + */ + private static void validateAggregators(List groupingColumns, AggregatorFactory[] aggregators) + { + final Set names = new HashSet<>(); + for (DimensionSchema column : groupingColumns) { + names.add(column.getName()); + } + for (AggregatorFactory aggregator : aggregators) { + if (aggregator == null) { + throw InvalidInput.exception("aggregators must not contain null entries"); + } + if (Granularities.GRANULARITY_VIRTUAL_COLUMN_NAME.equals(aggregator.getName())) { + throw InvalidInput.exception( + "aggregator cannot be named [%s]; it is the query-granularity virtual column, not a metric column", + Granularities.GRANULARITY_VIRTUAL_COLUMN_NAME + ); + } + if (!names.add(aggregator.getName())) { + throw InvalidInput.exception( + "aggregator [%s] duplicates the name of a column or another aggregator", + aggregator.getName() + ); + } + // The spec's aggregators are applied uniformly by every consumer: ingestion combines whatever arrives under the + // metric column's name, and compaction re-aggregates the rows the table has already stored. Both are only + // correct for an aggregator that combines its own output, so anything else (a COUNT, a sketch build, an input + // field that differs from the output) is rejected rather than silently changing results on re-aggregation. + final AggregatorFactory combining = aggregator.getCombiningFactory().withName(aggregator.getName()); + if (!aggregator.equals(combining)) { + throw InvalidInput.exception( + "aggregator [%s] is not its own combining form: a rollup table re-aggregates the rows it has stored, so" + + " its aggregators must combine their own output. Declare the combining form instead, for example [%s]", + aggregator.getName(), + combining + ); + } + } + } + + @Override + public boolean equals(Object o) + { + if (this == o) { + return true; + } + if (o == null || getClass() != o.getClass()) { + return false; + } + RollupTableProjectionSpec that = (RollupTableProjectionSpec) o; + return Objects.equals(virtualColumns, that.virtualColumns) + && Objects.equals(groupingColumns, that.groupingColumns) + && Arrays.equals(aggregators, that.aggregators); + } + + @Override + public int hashCode() + { + return Objects.hash(virtualColumns, groupingColumns, Arrays.hashCode(aggregators)); + } + + @Override + public String toString() + { + return "RollupTableProjectionSpec{" + + "virtualColumns=" + virtualColumns + + ", groupingColumns=" + groupingColumns + + ", aggregators=" + Arrays.toString(aggregators) + + '}'; + } + + /** + * Fluent builder for {@link RollupTableProjectionSpec}, avoiding the constructor's positional nullable leading + * {@code virtualColumns} argument. {@link #groupingColumns} (the full ordered grouping column list) is required. + */ + public static final class Builder + { + @Nullable + private VirtualColumns virtualColumns; + private List groupingColumns = Collections.emptyList(); + @Nullable + private AggregatorFactory[] aggregators; + + public Builder virtualColumns(@Nullable VirtualColumns virtualColumns) + { + this.virtualColumns = virtualColumns; + return this; + } + + /** + * The full, ordered list of grouping columns in segment order, including the explicit time marker. + */ + public Builder groupingColumns(List groupingColumns) + { + this.groupingColumns = groupingColumns; + return this; + } + + public Builder groupingColumns(DimensionSchema... groupingColumns) + { + return groupingColumns(Arrays.asList(groupingColumns)); + } + + public Builder aggregators(AggregatorFactory... aggregators) + { + this.aggregators = aggregators; + return this; + } + + public RollupTableProjectionSpec build() + { + return new RollupTableProjectionSpec(virtualColumns, groupingColumns, aggregators); + } + } +} diff --git a/processing/src/main/java/org/apache/druid/data/input/impl/TableProjectionSpec.java b/processing/src/main/java/org/apache/druid/data/input/impl/TableProjectionSpec.java index 18545e457366..5f36f13eca3c 100644 --- a/processing/src/main/java/org/apache/druid/data/input/impl/TableProjectionSpec.java +++ b/processing/src/main/java/org/apache/druid/data/input/impl/TableProjectionSpec.java @@ -24,7 +24,6 @@ import com.fasterxml.jackson.annotation.JsonInclude; import com.fasterxml.jackson.annotation.JsonProperty; import com.fasterxml.jackson.annotation.JsonTypeName; -import com.google.common.collect.Sets; import org.apache.druid.error.InvalidInput; import org.apache.druid.java.util.common.granularity.Granularities; import org.apache.druid.java.util.common.granularity.Granularity; @@ -40,7 +39,6 @@ import java.util.Collections; import java.util.List; import java.util.Objects; -import java.util.Set; /** * {@link BaseTableProjectionSpec} for a plain (non-clustered, non-rollup) table: the operator declares a single @@ -86,15 +84,15 @@ public TableProjectionSpec( @JsonProperty("columns") List columns ) { - validate(columns); + BaseTableProjectionSpec.validateDeclaredColumns(columns, "columns", TYPE_NAME); this.virtualColumns = virtualColumns == null ? VirtualColumns.EMPTY : virtualColumns; - validateVirtualColumns(this.virtualColumns); + BaseTableProjectionSpec.validateGranularityOnlyVirtualColumns(this.virtualColumns, TYPE_NAME); this.columns = Collections.unmodifiableList(new ArrayList<>(columns)); this.dimensionsSpec = DimensionsSpec.builder() .setDimensions(this.columns) .setForceSegmentSortByTime(false) .build(); - this.ordering = computeOrdering(this.columns); + this.ordering = BaseTableProjectionSpec.declaredOrderAscending(this.columns); } @Override @@ -225,74 +223,6 @@ private TableProjectionSpec withoutQueryGranularity() return new TableProjectionSpec(VirtualColumns.create(remaining), columns); } - private static void validate(List columns) - { - if (columns == null || columns.isEmpty()) { - throw InvalidInput.exception("columns must be non-empty for [%s] base table", TYPE_NAME); - } - final Set seen = Sets.newHashSetWithExpectedSize(columns.size()); - for (DimensionSchema d : columns) { - if (!seen.add(d.getName())) { - throw InvalidInput.exception("columns contains duplicate name [%s]", d.getName()); - } - } - boolean foundTime = false; - for (DimensionSchema column : columns) { - final String name = column.getName(); - // The query-granularity virtual column is a granularity carrier in virtualColumns (it floors the stored __time - // column); it is not itself a stored column, so it must not be declared in 'columns'. - if (Granularities.GRANULARITY_VIRTUAL_COLUMN_NAME.equals(name)) { - throw InvalidInput.exception( - "[%s] is the query-granularity virtual column, not a stored column; declare it in 'virtualColumns' and use" - + " [%s] as the time column in 'columns'", - Granularities.GRANULARITY_VIRTUAL_COLUMN_NAME, - ColumnHolder.TIME_COLUMN_NAME - ); - } - if (ColumnHolder.TIME_COLUMN_NAME.equals(name)) { - foundTime = true; - } - } - if (!foundTime) { - throw InvalidInput.exception( - "[%s] base table must include [%s] in 'columns' to define the time position", - TYPE_NAME, - ColumnHolder.TIME_COLUMN_NAME - ); - } - } - - /** - * A plain base table accepts only the {@link Granularities#GRANULARITY_VIRTUAL_COLUMN_NAME} carrier: the standard - * segment-generation path does not evaluate spec virtual columns, so any other virtual column would be dead metadata - * whose column never gets materialized. Computed columns belong in a {@code transformSpec}. - */ - private static void validateVirtualColumns(VirtualColumns virtualColumns) - { - for (VirtualColumn virtualColumn : virtualColumns.getVirtualColumns()) { - if (!Granularities.GRANULARITY_VIRTUAL_COLUMN_NAME.equals(virtualColumn.getOutputName())) { - throw InvalidInput.exception( - "virtual column [%s] is not supported: a [%s] base table stores only declared columns and does not" - + " materialize virtual columns; use a transformSpec to compute columns at ingest, or the [%s] virtual" - + " column to carry query granularity", - virtualColumn.getOutputName(), - TYPE_NAME, - Granularities.GRANULARITY_VIRTUAL_COLUMN_NAME - ); - } - BaseTableProjectionSpec.validateGranularity(virtualColumn, TYPE_NAME); - } - } - - private static List computeOrdering(List columns) - { - final List ordering = new ArrayList<>(columns.size()); - for (DimensionSchema d : columns) { - ordering.add(OrderBy.ascending(d.getName())); - } - return Collections.unmodifiableList(ordering); - } - @Override public boolean equals(Object o) { diff --git a/processing/src/test/java/org/apache/druid/data/input/impl/RollupTableProjectionSpecTest.java b/processing/src/test/java/org/apache/druid/data/input/impl/RollupTableProjectionSpecTest.java new file mode 100644 index 000000000000..e37876719eb8 --- /dev/null +++ b/processing/src/test/java/org/apache/druid/data/input/impl/RollupTableProjectionSpecTest.java @@ -0,0 +1,301 @@ +/* + * 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.data.input.impl; + +import com.fasterxml.jackson.databind.ObjectMapper; +import com.google.common.collect.ImmutableList; +import org.apache.druid.error.DruidException; +import org.apache.druid.java.util.common.DateTimes; +import org.apache.druid.java.util.common.granularity.Granularities; +import org.apache.druid.java.util.common.granularity.PeriodGranularity; +import org.apache.druid.query.OrderBy; +import org.apache.druid.query.aggregation.AggregatorFactory; +import org.apache.druid.query.aggregation.CountAggregatorFactory; +import org.apache.druid.query.aggregation.LongSumAggregatorFactory; +import org.apache.druid.segment.TestHelper; +import org.apache.druid.segment.VirtualColumn; +import org.apache.druid.testing.InitializedNullHandlingTest; +import org.joda.time.Period; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; + +import java.util.Collections; +import java.util.stream.Collectors; + +/** + * The column and granularity-carrier rules shared with {@link TableProjectionSpec} live in + * {@link BaseTableProjectionSpec} statics and are covered exhaustively in {@link TableProjectionSpecTest}; this covers + * what the rollup layout adds — the aggregators, and advertising rollup. + */ +class RollupTableProjectionSpecTest extends InitializedNullHandlingTest +{ + private static RollupTableProjectionSpec pagesSpec() + { + return RollupTableProjectionSpec.builder() + .groupingColumns(new StringDimensionSchema("page"), new LongDimensionSchema("__time")) + .aggregators(new LongSumAggregatorFactory("total", "total"), new LongSumAggregatorFactory("rows", "rows")) + .build(); + } + + @Test + void testShapeFollowsDeclaredOrderAndCarriesAggregators() + { + final RollupTableProjectionSpec spec = pagesSpec(); + + // Declared grouping order is both storage and sort order, with __time at its declared position; the aggregators + // become the metric columns and advertise the rollup layout. + Assertions.assertEquals( + ImmutableList.of(OrderBy.ascending("page"), OrderBy.ascending("__time")), + spec.getOrdering() + ); + Assertions.assertEquals(spec.getGroupingColumns(), spec.getDimensionsSpec().getDimensions()); + Assertions.assertFalse(spec.getDimensionsSpec().isForceSegmentSortByTime()); + Assertions.assertArrayEquals( + new AggregatorFactory[]{new LongSumAggregatorFactory("total", "total"), new LongSumAggregatorFactory("rows", "rows")}, + spec.getMetrics() + ); + Assertions.assertTrue(spec.isRollup()); + } + + @Test + void testEmptyAggregatorsAllowed() + { + // A rollup table without metrics collapses rows with identical grouping values, mirroring the legacy rollup=true + // granularity spec without a metricsSpec. + final RollupTableProjectionSpec spec = RollupTableProjectionSpec.builder() + .groupingColumns(new StringDimensionSchema("page"), new LongDimensionSchema("__time")) + .build(); + Assertions.assertEquals(0, spec.getMetrics().length); + Assertions.assertTrue(spec.isRollup()); + } + + @Test + void testAggregatorNameCollidingWithColumnRejected() + { + final DruidException e = Assertions.assertThrows( + DruidException.class, + () -> RollupTableProjectionSpec.builder() + .groupingColumns(new StringDimensionSchema("page"), new LongDimensionSchema("__time")) + .aggregators(new LongSumAggregatorFactory("page", "cnt")) + .build() + ); + Assertions.assertTrue( + e.getMessage().contains("aggregator [page] duplicates the name of a column or another aggregator"), + e.getMessage() + ); + } + + @Test + void testAggregatorNameCollidingWithAggregatorRejected() + { + final DruidException e = Assertions.assertThrows( + DruidException.class, + () -> RollupTableProjectionSpec.builder() + .groupingColumns(new StringDimensionSchema("page"), new LongDimensionSchema("__time")) + .aggregators(new LongSumAggregatorFactory("total", "total"), new CountAggregatorFactory("total")) + .build() + ); + Assertions.assertTrue( + e.getMessage().contains("aggregator [total] duplicates the name of a column or another aggregator"), + e.getMessage() + ); + } + + /** + * A metric named for the query-granularity carrier would be shadowed by the carrier virtual column once a query + * granularity is attached (virtual columns resolve before stored columns), so it is rejected at construction. + */ + @Test + void testAggregatorNamedForGranularityCarrierRejected() + { + final DruidException e = Assertions.assertThrows( + DruidException.class, + () -> RollupTableProjectionSpec.builder() + .groupingColumns(new StringDimensionSchema("page"), new LongDimensionSchema("__time")) + .aggregators(new LongSumAggregatorFactory(Granularities.GRANULARITY_VIRTUAL_COLUMN_NAME, "cnt")) + .build() + ); + Assertions.assertTrue( + e.getMessage().contains( + "aggregator cannot be named [" + Granularities.GRANULARITY_VIRTUAL_COLUMN_NAME + "]" + ), + e.getMessage() + ); + } + + /** + * The spec's aggregators are applied uniformly: ingestion combines what arrives under the metric column's name, and + * compaction re-aggregates stored rows. An aggregator that is not its own combining form (a COUNT, a sketch build, + * an input field differing from the output) would silently change results on re-aggregation, so it is rejected — a + * count is stored by summing a count column. + */ + @Test + void testNonSelfCombiningAggregatorRejected() + { + final DruidException count = Assertions.assertThrows( + DruidException.class, + () -> RollupTableProjectionSpec.builder() + .groupingColumns(new StringDimensionSchema("page"), new LongDimensionSchema("__time")) + .aggregators(new CountAggregatorFactory("rows")) + .build() + ); + Assertions.assertTrue(count.getMessage().contains("aggregator [rows] is not its own combining form"), count.getMessage()); + + final DruidException differentInput = Assertions.assertThrows( + DruidException.class, + () -> RollupTableProjectionSpec.builder() + .groupingColumns(new StringDimensionSchema("page"), new LongDimensionSchema("__time")) + .aggregators(new LongSumAggregatorFactory("total", "cnt")) + .build() + ); + Assertions.assertTrue( + differentInput.getMessage().contains("aggregator [total] is not its own combining form"), + differentInput.getMessage() + ); + } + + @Test + void testNullAggregatorEntryRejected() + { + final DruidException e = Assertions.assertThrows( + DruidException.class, + () -> RollupTableProjectionSpec.builder() + .groupingColumns(new StringDimensionSchema("page"), new LongDimensionSchema("__time")) + .aggregators(new LongSumAggregatorFactory("total", "total"), null) + .build() + ); + Assertions.assertTrue(e.getMessage().contains("aggregators must not contain null entries"), e.getMessage()); + } + + @Test + void testWithQueryGranularityAddsCarrierAndKeepsAggregators() + { + final RollupTableProjectionSpec spec = pagesSpec().withQueryGranularity(Granularities.HOUR); + + final VirtualColumn vc = spec.getVirtualColumns().getVirtualColumn(Granularities.GRANULARITY_VIRTUAL_COLUMN_NAME); + Assertions.assertNotNull(vc); + Assertions.assertEquals(Granularities.HOUR, spec.getQueryGranularity()); + Assertions.assertEquals(pagesSpec().getGroupingColumns(), spec.getGroupingColumns()); + Assertions.assertArrayEquals(pagesSpec().getMetrics(), spec.getMetrics()); + + // Idempotent once a carrier is present. + Assertions.assertSame(spec, spec.withQueryGranularity(Granularities.DAY)); + } + + /** + * The granularity rules are the shared {@link BaseTableProjectionSpec} validators, exercised exhaustively in + * {@link TableProjectionSpecTest}; this just proves the rollup spec is wired to them. + */ + @Test + void testWithQueryGranularityRejectionsAreWired() + { + Assertions.assertTrue( + Assertions.assertThrows(DruidException.class, () -> pagesSpec().withQueryGranularity(Granularities.ALL)) + .getMessage() + .contains("ALL") + ); + Assertions.assertTrue( + Assertions.assertThrows( + DruidException.class, + () -> pagesSpec().withQueryGranularity( + new PeriodGranularity(new Period("P1D"), null, DateTimes.inferTzFromString("America/Los_Angeles")) + ) + ).getMessage().contains("only period granularities in the UTC time zone") + ); + } + + @Test + void testHasEqualCompactionStateIgnoresQueryGranularityAndComparesAggregators() + { + Assertions.assertTrue(pagesSpec().withQueryGranularity(Granularities.HOUR).hasEqualCompactionState(pagesSpec())); + Assertions.assertTrue(pagesSpec().hasEqualCompactionState(pagesSpec().withQueryGranularity(Granularities.HOUR))); + + final RollupTableProjectionSpec differentAggregators = RollupTableProjectionSpec.builder() + .groupingColumns(new StringDimensionSchema("page"), new LongDimensionSchema("__time")) + .aggregators(new LongSumAggregatorFactory("total", "total")) + .build(); + Assertions.assertFalse(pagesSpec().hasEqualCompactionState(differentAggregators)); + } + + @Test + void testHasEqualCompactionStateRejectsDifferentSpecType() + { + // A plain table over the same columns is a different layout: it stores every row, a rollup table aggregates them. + final TableProjectionSpec plain = TableProjectionSpec.builder() + .columns(new StringDimensionSchema("page"), new LongDimensionSchema("__time")) + .build(); + Assertions.assertFalse(pagesSpec().hasEqualCompactionState(plain)); + Assertions.assertFalse(plain.hasEqualCompactionState(pagesSpec())); + } + + @Test + void testWithAdditionalColumnsAppendsGroupingColumnsAndKeepsAggregators() + { + final RollupTableProjectionSpec spec = pagesSpec() + .withQueryGranularity(Granularities.HOUR) + .withAdditionalColumns(ImmutableList.of(new StringDimensionSchema("city"))); + + Assertions.assertEquals( + ImmutableList.of("page", "__time", "city"), + spec.getGroupingColumns().stream().map(DimensionSchema::getName).collect(Collectors.toList()) + ); + Assertions.assertArrayEquals(pagesSpec().getMetrics(), spec.getMetrics()); + Assertions.assertEquals(Granularities.HOUR, spec.getQueryGranularity()); + + Assertions.assertSame(spec, spec.withAdditionalColumns(null)); + Assertions.assertSame(spec, spec.withAdditionalColumns(Collections.emptyList())); + } + + @Test + void testWithAdditionalColumnsRejectsTimeColumnAndAggregatorNames() + { + Assertions.assertTrue( + Assertions.assertThrows( + DruidException.class, + () -> pagesSpec().withAdditionalColumns(ImmutableList.of(new LongDimensionSchema("__time"))) + ).getMessage().contains("[__time]") + ); + // An appended column colliding with an aggregator would be one column with two definitions; the constructor's + // aggregator validation rejects it. + Assertions.assertTrue( + Assertions.assertThrows( + DruidException.class, + () -> pagesSpec().withAdditionalColumns(ImmutableList.of(new LongDimensionSchema("total"))) + ).getMessage().contains("aggregator [total] duplicates the name of a column or another aggregator") + ); + } + + @Test + void testSerdeRoundTripsThroughInterface() throws Exception + { + final ObjectMapper mapper = TestHelper.makeJsonMapper(); + + final BaseTableProjectionSpec bare = pagesSpec(); + final String bareJson = mapper.writeValueAsString(bare); + Assertions.assertTrue(bareJson.contains("\"type\":\"rollupTable\""), bareJson); + Assertions.assertEquals(bare, mapper.readValue(bareJson, BaseTableProjectionSpec.class)); + + final BaseTableProjectionSpec withGranularity = pagesSpec().withQueryGranularity(Granularities.HOUR); + Assertions.assertEquals( + withGranularity, + mapper.readValue(mapper.writeValueAsString(withGranularity), BaseTableProjectionSpec.class) + ); + } +} diff --git a/server/src/main/java/org/apache/druid/catalog/model/DatasourceBaseTableMetadata.java b/server/src/main/java/org/apache/druid/catalog/model/DatasourceBaseTableMetadata.java index 17b2de200206..8587eae97419 100644 --- a/server/src/main/java/org/apache/druid/catalog/model/DatasourceBaseTableMetadata.java +++ b/server/src/main/java/org/apache/druid/catalog/model/DatasourceBaseTableMetadata.java @@ -37,6 +37,7 @@ @JsonTypeInfo(use = JsonTypeInfo.Id.NAME, include = JsonTypeInfo.As.EXISTING_PROPERTY, property = "type") @JsonSubTypes(value = { @JsonSubTypes.Type(name = TableBaseTableMetadata.TYPE_NAME, value = TableBaseTableMetadata.class), + @JsonSubTypes.Type(name = RollupTableBaseTableMetadata.TYPE_NAME, value = RollupTableBaseTableMetadata.class), @JsonSubTypes.Type( name = ClusteredValueGroupsBaseTableMetadata.TYPE_NAME, value = ClusteredValueGroupsBaseTableMetadata.class diff --git a/server/src/main/java/org/apache/druid/catalog/model/RollupTableBaseTableMetadata.java b/server/src/main/java/org/apache/druid/catalog/model/RollupTableBaseTableMetadata.java new file mode 100644 index 000000000000..e7c05fb91844 --- /dev/null +++ b/server/src/main/java/org/apache/druid/catalog/model/RollupTableBaseTableMetadata.java @@ -0,0 +1,232 @@ +/* + * 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.catalog.model; + +import com.fasterxml.jackson.annotation.JsonCreator; +import com.fasterxml.jackson.annotation.JsonInclude; +import com.fasterxml.jackson.annotation.JsonProperty; +import com.fasterxml.jackson.annotation.JsonTypeName; +import org.apache.druid.data.input.impl.DimensionSchema; +import org.apache.druid.data.input.impl.RollupTableProjectionSpec; +import org.apache.druid.error.InvalidInput; +import org.apache.druid.query.aggregation.AggregatorFactory; +import org.apache.druid.segment.VirtualColumns; +import org.apache.druid.segment.column.ColumnType; +import org.apache.druid.utils.CollectionUtils; + +import javax.annotation.Nullable; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.HashMap; +import java.util.HashSet; +import java.util.List; +import java.util.Map; +import java.util.Objects; +import java.util.Set; + +/** + * Catalog layout metadata for {@link RollupTableProjectionSpec} base tables. The catalog column list declares the + * output schema of the rollup: the columns named by an entry of {@link #aggregators} are the metric columns, every + * other declared column is a grouping column, and, mirroring the physical layout, the grouping columns must be + * declared before the metric columns. {@link #createSpec(List)} combines the declared columns with this metadata into + * the physical spec used to generate segments; as with the other layouts, the optional {@link #virtualColumns} carry + * only the query granularity, and the optional {@link #columnSchemas} customize the {@link DimensionSchema} used for a + * declared grouping column (a metric column's physical representation is fixed by its aggregator). + */ +@JsonTypeName(RollupTableBaseTableMetadata.TYPE_NAME) +public class RollupTableBaseTableMetadata implements DatasourceBaseTableMetadata +{ + public static final String TYPE_NAME = RollupTableProjectionSpec.TYPE_NAME; + + private final VirtualColumns virtualColumns; + private final AggregatorFactory[] aggregators; + private final List columnSchemas; + + @JsonCreator + public RollupTableBaseTableMetadata( + @JsonProperty("virtualColumns") @Nullable VirtualColumns virtualColumns, + @JsonProperty("aggregators") @Nullable AggregatorFactory[] aggregators, + @JsonProperty("columnSchemas") @Nullable List columnSchemas + ) + { + this.virtualColumns = virtualColumns == null ? VirtualColumns.EMPTY : virtualColumns; + this.aggregators = aggregators == null ? new AggregatorFactory[0] : aggregators; + this.columnSchemas = columnSchemas == null ? Collections.emptyList() : columnSchemas; + } + + @Override + @JsonProperty("type") + public String getType() + { + return TYPE_NAME; + } + + @Override + @JsonProperty("virtualColumns") + @JsonInclude(JsonInclude.Include.NON_DEFAULT) + public VirtualColumns getVirtualColumns() + { + return virtualColumns; + } + + /** + * The aggregators computed over the rows collapsed into each grouping tuple. Each names the declared metric column + * it fills. + */ + @JsonProperty("aggregators") + @JsonInclude(JsonInclude.Include.NON_EMPTY) + public AggregatorFactory[] getAggregators() + { + return aggregators; + } + + /** + * Per-column customizations of the {@link DimensionSchema} used during segment creation, keyed by + * {@link DimensionSchema#getName()}; empty when every declared grouping column uses the schema derived from its + * declared type. These do not define columns: every entry must customize a declared grouping column. + */ + @JsonProperty("columnSchemas") + @JsonInclude(JsonInclude.Include.NON_EMPTY) + public List getColumnSchemas() + { + return columnSchemas; + } + + /** + * Creates the physical spec from the declared catalog columns: a declared column named by an aggregator is a metric + * column, whose declared type must match what its aggregator produces; every other declared column is a grouping + * column, derived like the other layouts derive theirs. Declared column order is the physical segment order, so the + * grouping columns must be declared before the metric columns; the metric columns' declared order becomes the + * aggregator order of the spec. + */ + @Override + public RollupTableProjectionSpec createSpec(List columns) + { + if (CollectionUtils.isNullOrEmpty(columns)) { + throw InvalidInput.exception( + "Cannot define a [%s] base table without declared columns; the catalog column list defines the table schema", + TYPE_NAME + ); + } + final Map aggregatorsByName = new HashMap<>(); + for (AggregatorFactory aggregator : aggregators) { + if (aggregator == null) { + throw InvalidInput.exception("aggregators must not contain null entries"); + } + if (aggregatorsByName.put(aggregator.getName(), aggregator) != null) { + throw InvalidInput.exception("aggregators contains duplicate entries for column [%s]", aggregator.getName()); + } + } + final Map customSchemas = BaseTableColumns.indexColumnSchemas(columnSchemas); + + final Set declaredNames = new HashSet<>(); + final List groupingColumns = new ArrayList<>(columns.size()); + final List orderedAggregators = new ArrayList<>(aggregators.length); + for (ColumnSpec column : columns) { + declaredNames.add(column.name()); + final AggregatorFactory aggregator = aggregatorsByName.get(column.name()); + if (aggregator == null) { + if (!orderedAggregators.isEmpty()) { + throw InvalidInput.exception( + "grouping column [%s] is declared after metric column [%s]; the declared order is the physical segment" + + " order, so grouping columns must be declared before the metric columns", + column.name(), + orderedAggregators.get(orderedAggregators.size() - 1).getName() + ); + } + groupingColumns.add(BaseTableColumns.toDimensionSchema(column, customSchemas.get(column.name()), TYPE_NAME)); + } else { + final ColumnType declaredType = Columns.druidType(column); + final ColumnType producedType = aggregator.getIntermediateType(); + if (declaredType == null || !declaredType.equals(producedType)) { + throw InvalidInput.exception( + "metric column [%s] is declared as type [%s], but its aggregator produces [%s]; the declared type of a" + + " metric column must match what its aggregator stores", + column.name(), + column.dataType(), + producedType + ); + } + if (customSchemas.containsKey(column.name())) { + throw InvalidInput.exception( + "columnSchemas cannot customize metric column [%s]: the physical representation of a metric column is" + + " fixed by its aggregator", + column.name() + ); + } + orderedAggregators.add(aggregator); + } + } + for (String aggregatorName : aggregatorsByName.keySet()) { + if (!declaredNames.contains(aggregatorName)) { + throw InvalidInput.exception( + "aggregator [%s] does not fill a declared column; every aggregator produces a metric column, so declare" + + " [%s] in the table's column list", + aggregatorName, + aggregatorName + ); + } + } + for (String customized : customSchemas.keySet()) { + if (!declaredNames.contains(customized)) { + throw InvalidInput.exception( + "columnSchemas entry [%s] does not customize a declared column; column schemas do not define columns," + + " declare [%s] in the table's column list", + customized, + customized + ); + } + } + return RollupTableProjectionSpec.builder() + .virtualColumns(virtualColumns) + .groupingColumns(groupingColumns) + .aggregators(orderedAggregators.toArray(new AggregatorFactory[0])) + .build(); + } + + @Override + public boolean equals(Object o) + { + if (o == null || getClass() != o.getClass()) { + return false; + } + RollupTableBaseTableMetadata that = (RollupTableBaseTableMetadata) o; + return Objects.equals(virtualColumns, that.virtualColumns) + && Arrays.equals(aggregators, that.aggregators) + && Objects.equals(columnSchemas, that.columnSchemas); + } + + @Override + public int hashCode() + { + return Objects.hash(virtualColumns, Arrays.hashCode(aggregators), columnSchemas); + } + + @Override + public String toString() + { + return "RollupTableBaseTableMetadata{" + + "virtualColumns=" + virtualColumns + + ", aggregators=" + Arrays.toString(aggregators) + + ", columnSchemas=" + columnSchemas + + '}'; + } +} diff --git a/server/src/main/java/org/apache/druid/segment/indexing/DataSchema.java b/server/src/main/java/org/apache/druid/segment/indexing/DataSchema.java index 51051337a3a5..f2f70bbc379e 100644 --- a/server/src/main/java/org/apache/druid/segment/indexing/DataSchema.java +++ b/server/src/main/java/org/apache/druid/segment/indexing/DataSchema.java @@ -217,7 +217,8 @@ private static void rejectLegacyTopLevelSchemaFields( * {@link GranularitySpec}, which is already the complete source of truth (segment + query granularity + rollup), so * it is returned as-is. In baseTable mode the pieces live apart and are recombined: segment granularity + intervals * from {@code segmentGranularitySpec} (or defaults), query granularity from the spec's - * {@link Granularities#GRANULARITY_VIRTUAL_COLUMN_NAME} virtual column (absent ⇒ {@code NONE}). + * {@link Granularities#GRANULARITY_VIRTUAL_COLUMN_NAME} virtual column (absent ⇒ {@code NONE}), and rollup from the + * spec's layout type ({@link BaseTableProjectionSpec#isRollup()}). */ private static GranularitySpec computeEffectiveGranularitySpec( BaseTableProjectionSpec effectiveSpec, @@ -234,7 +235,7 @@ private static GranularitySpec computeEffectiveGranularitySpec( return new UniformGranularitySpec( segmentGranularity, queryGranularityFromSpec(effectiveSpec), - false, + effectiveSpec.isRollup(), intervals ); } diff --git a/server/src/test/java/org/apache/druid/catalog/model/RollupTableBaseTableMetadataTest.java b/server/src/test/java/org/apache/druid/catalog/model/RollupTableBaseTableMetadataTest.java new file mode 100644 index 000000000000..9af97f9dcc86 --- /dev/null +++ b/server/src/test/java/org/apache/druid/catalog/model/RollupTableBaseTableMetadataTest.java @@ -0,0 +1,267 @@ +/* + * 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.catalog.model; + +import com.fasterxml.jackson.databind.InjectableValues; +import com.fasterxml.jackson.databind.ObjectMapper; +import org.apache.druid.data.input.impl.DimensionSchema; +import org.apache.druid.data.input.impl.RollupTableProjectionSpec; +import org.apache.druid.data.input.impl.StringDimensionSchema; +import org.apache.druid.error.DruidException; +import org.apache.druid.guice.BuiltInTypesModule; +import org.apache.druid.jackson.DefaultObjectMapper; +import org.apache.druid.java.util.common.granularity.Granularities; +import org.apache.druid.math.expr.ExprMacroTable; +import org.apache.druid.query.aggregation.AggregatorFactory; +import org.apache.druid.query.aggregation.LongSumAggregatorFactory; +import org.apache.druid.segment.DefaultColumnFormatConfig; +import org.apache.druid.segment.VirtualColumns; +import org.apache.druid.testing.InitializedNullHandlingTest; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; + +import java.util.Arrays; +import java.util.Collections; +import java.util.List; +import java.util.stream.Collectors; + +/** + * The column-derivation rules shared with the other layouts are covered exhaustively in + * {@link ClusteredValueGroupsBaseTableMetadataTest} and route through the same {@code BaseTableColumns} helper; this + * covers what the rollup layout adds — splitting the declared columns into grouping and metric columns by aggregator + * name. + */ +public class RollupTableBaseTableMetadataTest extends InitializedNullHandlingTest +{ + static { + BuiltInTypesModule.registerHandlersAndSerde(); + } + + // The granularity macro table (not nil) so the carrier's timestamp_floor expression can be parsed back. + private final ObjectMapper mapper = new DefaultObjectMapper().setInjectableValues( + new InjectableValues.Std() + .addValue(ExprMacroTable.class, ExprMacroTable.granularity()) + .addValue(DefaultColumnFormatConfig.class, new DefaultColumnFormatConfig(null, null, null, null)) + ); + + private static final List COLUMNS = Arrays.asList( + new ColumnSpec("tenant", Columns.SQL_VARCHAR, null), + new ColumnSpec(Columns.TIME_COLUMN, Columns.SQL_TIMESTAMP, null), + new ColumnSpec("total", Columns.SQL_BIGINT, null) + ); + + private static VirtualColumns hourCarrier() + { + return VirtualColumns.create( + Granularities.toVirtualColumn(Granularities.HOUR, Granularities.GRANULARITY_VIRTUAL_COLUMN_NAME) + ); + } + + private static AggregatorFactory[] totalSum() + { + return new AggregatorFactory[]{new LongSumAggregatorFactory("total", "total")}; + } + + @Test + public void testSerde() throws Exception + { + final DatasourceBaseTableMetadata metadata = new RollupTableBaseTableMetadata(hourCarrier(), totalSum(), null); + final String json = mapper.writeValueAsString(metadata); + Assertions.assertTrue(json.contains("\"type\":\"rollupTable\""), json); + final DatasourceBaseTableMetadata fromJson = mapper.readValue(json, DatasourceBaseTableMetadata.class); + Assertions.assertEquals(metadata, fromJson); + } + + @Test + public void testSerdeEmpty() throws Exception + { + final DatasourceBaseTableMetadata metadata = new RollupTableBaseTableMetadata(null, null, null); + final String json = mapper.writeValueAsString(metadata); + Assertions.assertFalse(json.contains("virtualColumns")); + Assertions.assertFalse(json.contains("aggregators")); + Assertions.assertFalse(json.contains("columnSchemas")); + final DatasourceBaseTableMetadata fromJson = mapper.readValue(json, DatasourceBaseTableMetadata.class); + Assertions.assertEquals(metadata, fromJson); + } + + @Test + public void testCreateSpec() + { + final RollupTableProjectionSpec spec = + new RollupTableBaseTableMetadata(hourCarrier(), totalSum(), null).createSpec(COLUMNS); + + // Declared columns split by aggregator name: everything else is a grouping column, in declared order; the carrier + // rides through as the query granularity, and the layout advertises rollup. + Assertions.assertEquals( + List.of("tenant", Columns.TIME_COLUMN), + spec.getGroupingColumns().stream().map(DimensionSchema::getName).collect(Collectors.toList()) + ); + Assertions.assertArrayEquals(totalSum(), spec.getMetrics()); + Assertions.assertEquals(Granularities.HOUR, spec.getQueryGranularity()); + Assertions.assertTrue(spec.isRollup()); + } + + @Test + public void testCreateSpecWithoutAggregatorsCollapsesDuplicates() + { + // No aggregators is a dedup table: every declared column is a grouping column. + final RollupTableProjectionSpec spec = new RollupTableBaseTableMetadata(null, null, null).createSpec( + Arrays.asList( + new ColumnSpec("tenant", Columns.SQL_VARCHAR, null), + new ColumnSpec(Columns.TIME_COLUMN, Columns.SQL_TIMESTAMP, null) + ) + ); + Assertions.assertEquals(0, spec.getMetrics().length); + Assertions.assertTrue(spec.isRollup()); + } + + @Test + public void testCreateSpecGroupingColumnAfterMetricColumnFails() + { + final DruidException e = Assertions.assertThrows( + DruidException.class, + () -> new RollupTableBaseTableMetadata(null, totalSum(), null).createSpec( + Arrays.asList( + new ColumnSpec(Columns.TIME_COLUMN, Columns.SQL_TIMESTAMP, null), + new ColumnSpec("total", Columns.SQL_BIGINT, null), + new ColumnSpec("tenant", Columns.SQL_VARCHAR, null) + ) + ) + ); + Assertions.assertTrue( + e.getMessage().contains("grouping column [tenant] is declared after metric column [total]"), + e.getMessage() + ); + } + + /** + * A declared metric column named for the query-granularity carrier slips past the grouping-column carrier-name + * check (it is matched to its aggregator, not treated as a grouping column), so the spec's aggregator validation is + * what rejects it: the carrier virtual column would shadow the metric once a query granularity is attached. + */ + @Test + public void testCreateSpecAggregatorNamedForGranularityCarrierFails() + { + final DruidException e = Assertions.assertThrows( + DruidException.class, + () -> new RollupTableBaseTableMetadata( + null, + new AggregatorFactory[]{ + new LongSumAggregatorFactory(Granularities.GRANULARITY_VIRTUAL_COLUMN_NAME, "cnt") + }, + null + ).createSpec( + Arrays.asList( + new ColumnSpec("tenant", Columns.SQL_VARCHAR, null), + new ColumnSpec(Columns.TIME_COLUMN, Columns.SQL_TIMESTAMP, null), + new ColumnSpec(Granularities.GRANULARITY_VIRTUAL_COLUMN_NAME, Columns.SQL_BIGINT, null) + ) + ) + ); + Assertions.assertTrue( + e.getMessage().contains( + "aggregator cannot be named [" + Granularities.GRANULARITY_VIRTUAL_COLUMN_NAME + "]" + ), + e.getMessage() + ); + } + + @Test + public void testCreateSpecAggregatorWithoutDeclaredColumnFails() + { + final DruidException e = Assertions.assertThrows( + DruidException.class, + () -> new RollupTableBaseTableMetadata( + null, + new AggregatorFactory[]{new LongSumAggregatorFactory("nope", "nope")}, + null + ).createSpec(COLUMNS) + ); + Assertions.assertTrue( + e.getMessage().contains("aggregator [nope] does not fill a declared column"), + e.getMessage() + ); + } + + @Test + public void testCreateSpecMetricTypeMismatchFails() + { + final DruidException e = Assertions.assertThrows( + DruidException.class, + () -> new RollupTableBaseTableMetadata(null, totalSum(), null).createSpec( + Arrays.asList( + new ColumnSpec("tenant", Columns.SQL_VARCHAR, null), + new ColumnSpec(Columns.TIME_COLUMN, Columns.SQL_TIMESTAMP, null), + new ColumnSpec("total", Columns.SQL_DOUBLE, null) + ) + ) + ); + Assertions.assertTrue( + e.getMessage().contains("metric column [total] is declared as type [DOUBLE], but its aggregator produces [LONG]"), + e.getMessage() + ); + } + + @Test + public void testCreateSpecColumnSchemaForMetricColumnFails() + { + final DruidException e = Assertions.assertThrows( + DruidException.class, + () -> new RollupTableBaseTableMetadata( + null, + totalSum(), + Collections.singletonList(new StringDimensionSchema("total")) + ).createSpec(COLUMNS) + ); + Assertions.assertTrue( + e.getMessage().contains("columnSchemas cannot customize metric column [total]"), + e.getMessage() + ); + } + + @Test + public void testCreateSpecColumnSchemaForGroupingColumnApplied() + { + final RollupTableProjectionSpec spec = new RollupTableBaseTableMetadata( + null, + totalSum(), + Collections.singletonList( + new StringDimensionSchema("tenant", DimensionSchema.MultiValueHandling.ARRAY, false) + ) + ).createSpec(COLUMNS); + Assertions.assertEquals( + new StringDimensionSchema("tenant", DimensionSchema.MultiValueHandling.ARRAY, false), + spec.getGroupingColumns().get(0) + ); + } + + @Test + public void testCreateSpecNoDeclaredColumnsFails() + { + final DruidException e = Assertions.assertThrows( + DruidException.class, + () -> new RollupTableBaseTableMetadata(null, totalSum(), null).createSpec(null) + ); + Assertions.assertTrue( + e.getMessage().contains("Cannot define a [rollupTable] base table without declared columns"), + e.getMessage() + ); + } +} diff --git a/server/src/test/java/org/apache/druid/segment/indexing/DataSchemaTest.java b/server/src/test/java/org/apache/druid/segment/indexing/DataSchemaTest.java index 7836f057d326..3e1fa441aa16 100644 --- a/server/src/test/java/org/apache/druid/segment/indexing/DataSchemaTest.java +++ b/server/src/test/java/org/apache/druid/segment/indexing/DataSchemaTest.java @@ -34,6 +34,7 @@ import org.apache.druid.data.input.impl.ClusteredValueGroupsBaseTableProjectionSpec; import org.apache.druid.data.input.impl.DimensionsSpec; import org.apache.druid.data.input.impl.LongDimensionSchema; +import org.apache.druid.data.input.impl.RollupTableProjectionSpec; import org.apache.druid.data.input.impl.StringDimensionSchema; import org.apache.druid.data.input.impl.TableProjectionSpec; import org.apache.druid.data.input.impl.TimestampSpec; @@ -969,6 +970,47 @@ void testBaseTableModeWithTableProjectionSpec() throws IOException Assertions.assertFalse(effectiveGranularity.isRollup()); } + /** + * The rollup layout is the plain shape plus aggregators, so what it adds at this level is: aggregators delegate + * through {@link DataSchema#getAggregators()}, and the recombined granularity carries {@code rollup=true} from the + * spec's layout type (where the plain layout recombines with {@code rollup=false}). + */ + @Test + void testBaseTableModeWithRollupTableProjectionSpec() throws IOException + { + final RollupTableProjectionSpec spec = RollupTableProjectionSpec.builder() + .groupingColumns( + new StringDimensionSchema("page"), + new LongDimensionSchema("__time") + ) + .aggregators(new LongSumAggregatorFactory("total", "total")) + .build() + .withQueryGranularity(Granularities.HOUR); + final DataSchema original = DataSchema.builder() + .withDataSource("datasource") + .withTimestamp(TIMESTAMP_SPEC) + .withSegmentGranularity(new SegmentGranularitySpec(Granularities.DAY, null)) + .withBaseTable(spec) + .build(); + + final String serialized = jsonMapper.writeValueAsString(original); + final JsonNode root = jsonMapper.readTree(serialized); + Assertions.assertEquals("rollupTable", root.get("baseTable").get("type").asText()); + Assertions.assertFalse(root.has("metricsSpec")); + final DataSchema deserialized = jsonMapper.readValue(serialized, DataSchema.class); + Assertions.assertEquals(original, deserialized); + Assertions.assertEquals(spec, deserialized.getBaseTable()); + + Assertions.assertArrayEquals( + new AggregatorFactory[]{new LongSumAggregatorFactory("total", "total")}, + original.getAggregators() + ); + final GranularitySpec effectiveGranularity = original.getGranularitySpec(); + Assertions.assertEquals(Granularities.DAY, effectiveGranularity.getSegmentGranularity()); + Assertions.assertEquals(Granularities.HOUR, effectiveGranularity.getQueryGranularity()); + Assertions.assertTrue(effectiveGranularity.isRollup()); + } + @Test void testBaseTableModeLegacyAccessorsDelegateToSpec() { diff --git a/sql/src/main/java/org/apache/druid/sql/calcite/planner/ProjectionSpecTranslator.java b/sql/src/main/java/org/apache/druid/sql/calcite/planner/ProjectionSpecTranslator.java index dbb48b78f828..85586e10a542 100644 --- a/sql/src/main/java/org/apache/druid/sql/calcite/planner/ProjectionSpecTranslator.java +++ b/sql/src/main/java/org/apache/druid/sql/calcite/planner/ProjectionSpecTranslator.java @@ -31,6 +31,7 @@ import org.apache.druid.catalog.model.ColumnSpec; import org.apache.druid.catalog.model.Columns; import org.apache.druid.catalog.model.DatasourceBaseTableMetadata; +import org.apache.druid.catalog.model.RollupTableBaseTableMetadata; import org.apache.druid.catalog.model.TableBaseTableMetadata; import org.apache.druid.data.input.impl.AggregateProjectionSpec; import org.apache.druid.data.input.impl.DimensionSchema; @@ -65,6 +66,7 @@ import org.apache.druid.sql.calcite.table.DatasourceTable; import org.apache.druid.sql.calcite.table.DatasourceTable.PhysicalDatasourceMetadata; import org.apache.druid.sql.calcite.table.DruidTable; +import org.apache.druid.utils.CollectionUtils; import org.joda.time.Interval; import javax.annotation.Nullable; @@ -151,9 +153,14 @@ public AggregateProjectionSpec translate( * Without {@code CLUSTERED BY} the layout is a plain table, which stores columns as they arrive; the only * expression it accepts is {@code TIME_FLOOR(__time, )} selected as {@code __time}, which becomes the * table's query granularity (carried as the granularity virtual column of the spec). + *

+ * A body with {@code GROUP BY} declares the rollup layout: the body is the query that rolls the table up onto + * itself. A grouping output is a declared column referenced as itself (or the {@code TIME_FLOOR} of {@code __time}), + * and an aggregate output becomes the aggregator that fills the declared metric column it is named for, reading the + * column it fills, since aggregating a stored rollup table is done with the combining form of its aggregators. * * @param clusteredBy the columns segments are clustered on, which must be the leading prefix of the column list; - * null selects the plain-table layout + * null selects the plain-table (or, with GROUP BY, the rollup) layout */ public DatasourceBaseTableMetadata translateBaseTable( final String tableName, @@ -162,27 +169,39 @@ public DatasourceBaseTableMetadata translateBaseTable( @Nullable final SqlNodeList clusteredBy ) { - if (body.getWhere() != null || body.getGroup() != null) { + if (body.getWhere() != null) { throw invalid( BASE_PROJECTION_NAME, - "its body filters or groups. The base table stores every ingested row, so it can do neither" + "its body filters. The base table stores every ingested row, so it cannot filter" ); } rejectSubqueries(BASE_PROJECTION_NAME, body); - final DruidQuery druidQuery = planBody(tableName, columns, BASE_PROJECTION_NAME, body); - final List computed = liftComputedColumns(columns, druidQuery); final DatasourceBaseTableMetadata metadata; - if (clusteredBy == null) { - metadata = new TableBaseTableMetadata(plainTableVirtualColumns(computed), null); + if (body.getGroup() != null) { + // A GROUP BY body declares the rollup layout: rows with identical grouping values aggregate into one. + if (clusteredBy != null) { + throw invalid( + BASE_PROJECTION_NAME, + "its body groups and declares CLUSTERED BY. A clustered base table stores every ingested row, so a base" + + " table is either clustered or rollup, not both" + ); + } + metadata = liftRollupBaseTable(tableName, columns, planBody(tableName, columns, BASE_PROJECTION_NAME, body)); } else { - // liftComputedColumns already verified that the body plans to a scan. - final VirtualColumns planned = druidQuery.getQuery().getVirtualColumns(); - metadata = new ClusteredValueGroupsBaseTableMetadata( - clusteringColumns(clusteredBy), - materializedVirtualColumns(computed, planned, columns), - null - ); + final DruidQuery druidQuery = planBody(tableName, columns, BASE_PROJECTION_NAME, body); + final List computed = liftComputedColumns(columns, druidQuery); + if (clusteredBy == null) { + metadata = new TableBaseTableMetadata(plainTableVirtualColumns(computed), null); + } else { + // liftComputedColumns already verified that the body plans to a scan. + final VirtualColumns planned = druidQuery.getQuery().getVirtualColumns(); + metadata = new ClusteredValueGroupsBaseTableMetadata( + clusteringColumns(clusteredBy), + materializedVirtualColumns(computed, planned, columns), + null + ); + } } // Derive the physical spec now. The catalog does this too when the write lands, but doing it here attributes @@ -288,6 +307,128 @@ private static void collectDependencies( } } + /** + * Lift a GROUP BY body into the rollup layout's metadata. The body plans like an aggregate projection body, and its + * outputs must name every declared column in declared order: an output backed by a grouping dimension must reference + * the declared column itself (or {@code TIME_FLOOR(__time, ) AS __time}, which becomes the query + * granularity carrier) and an output backed by an aggregate becomes the aggregator filling the declared metric + * column it is named for. + */ + private static RollupTableBaseTableMetadata liftRollupBaseTable( + final String tableName, + final List columns, + final DruidQuery druidQuery + ) + { + final PlannedAggregation planned = plannedAggregation(BASE_PROJECTION_NAME, tableName, druidQuery); + if (planned.virtualColumnsAndFilter.filter(BASE_PROJECTION_NAME) != null) { + throw invalid(BASE_PROJECTION_NAME, "its body filters. The base table stores every ingested row"); + } + + final List outputNames = druidQuery.getOutputRowType().getFieldNames(); + if (outputNames.size() != columns.size()) { + throw invalid( + BASE_PROJECTION_NAME, + StringUtils.format( + "it selects %d column(s) but the table declares %d. The body lists the columns in the order segments" + + " store them, so it must name every declared column", + outputNames.size(), + columns.size() + ) + ); + } + + final Map dimensionsByInternalName = new HashMap<>(); + for (DimensionSpec dimension : planned.dimensions) { + dimensionsByInternalName.put(dimension.getOutputName(), dimension); + } + final Map aggregatorsByInternalName = new HashMap<>(); + for (AggregatorFactory aggregator : planned.aggregators) { + aggregatorsByInternalName.put(aggregator.getName(), aggregator); + } + + final RowSignature sources = druidQuery.getOutputRowSignature(); + final List carrier = new ArrayList<>(1); + final List aggregators = new ArrayList<>(planned.aggregators.size()); + for (int i = 0; i < columns.size(); i++) { + final String declared = columns.get(i).name(); + if (!declared.equals(outputNames.get(i))) { + throw invalid( + BASE_PROJECTION_NAME, + StringUtils.format( + "its column %d is [%s] but the table declares [%s] there. The body lists the columns in the order" + + " segments store them", + i + 1, + outputNames.get(i), + declared + ) + ); + } + final String source = sources.getColumnName(i); + final AggregatorFactory aggregator = aggregatorsByInternalName.get(source); + if (aggregator != null) { + aggregators.add(aggregator.withName(declared)); + continue; + } + final DimensionSpec dimension = dimensionsByInternalName.get(source); + if (!(dimension instanceof DefaultDimensionSpec)) { + throw invalid( + BASE_PROJECTION_NAME, + "its grouping column [" + declared + "] is not a plain column reference" + ); + } + final String grouped = dimension.getDimension(); + if (grouped.equals(declared)) { + // A grouping column referenced as itself: derived from the declared column by the catalog, nothing to lift. + continue; + } + final VirtualColumn virtualColumn = + planned.virtualColumnsAndFilter.virtualColumns.getVirtualColumn(grouped); + if (ColumnHolder.TIME_COLUMN_NAME.equals(declared) && virtualColumn != null) { + final Granularity granularity = Granularities.fromTimeVirtualColumn(virtualColumn); + if (granularity == null) { + throw invalid( + BASE_PROJECTION_NAME, + StringUtils.format( + "it computes [%s] with an expression that is not a granularity. A base table computes [%s] only as" + + " TIME_FLOOR(%s, ), which declares the table's query granularity", + ColumnHolder.TIME_COLUMN_NAME, + ColumnHolder.TIME_COLUMN_NAME, + ColumnHolder.TIME_COLUMN_NAME + ) + ); + } + carrier.add(Granularities.toVirtualColumn(granularity, Granularities.GRANULARITY_VIRTUAL_COLUMN_NAME)); + } else if (virtualColumn != null) { + throw invalid( + BASE_PROJECTION_NAME, + StringUtils.format( + "its column [%s] is computed by an expression. Computed columns are not supported; compute the column" + + " at ingestion time instead", + declared + ) + ); + } else { + throw invalid( + BASE_PROJECTION_NAME, + StringUtils.format( + "its column %d groups [%s] but declares it as [%s]. A grouping column is stored under the name it" + + " groups, so nothing would fill [%s]", + i + 1, + grouped, + declared, + declared + ) + ); + } + } + return new RollupTableBaseTableMetadata( + VirtualColumns.create(carrier), + aggregators.toArray(new AggregatorFactory[0]), + null + ); + } + private static List clusteringColumns(@Nullable final SqlNodeList clusteredBy) { if (clusteredBy == null) { @@ -475,6 +616,44 @@ private static AggregateProjectionSpec lift( final String tableName, final DruidQuery druidQuery ) + { + final PlannedAggregation planned = plannedAggregation(projectionName, tableName, druidQuery); + return AggregateProjectionSpec + .builder(projectionName) + .virtualColumns(planned.virtualColumnsAndFilter.virtualColumns) + .filter(planned.virtualColumnsAndFilter.filter(projectionName)) + .groupingColumns(groupingColumns(projectionName, planned.dimensions)) + .aggregators(renameToOutputNames(projectionName, druidQuery, planned.aggregators)) + .build(); + } + + /** + * The pieces of a planned aggregating body that both the aggregate-projection lift and the rollup base-table lift + * consume: the grouping dimensions, the aggregators, and the virtual columns / filter / intervals of the query. + */ + private static final class PlannedAggregation + { + private final List dimensions; + private final List aggregators; + private final VirtualColumnsAndFilter virtualColumnsAndFilter; + + private PlannedAggregation( + List dimensions, + List aggregators, + VirtualColumnsAndFilter virtualColumnsAndFilter + ) + { + this.dimensions = dimensions; + this.aggregators = aggregators; + this.virtualColumnsAndFilter = virtualColumnsAndFilter; + } + } + + private static PlannedAggregation plannedAggregation( + final String projectionName, + final String tableName, + final DruidQuery druidQuery + ) { final DataSource dataSource = druidQuery.getDataSource(); if (!(dataSource instanceof TableDataSource) || !tableName.equals(((TableDataSource) dataSource).getName())) { @@ -501,27 +680,36 @@ private static AggregateProjectionSpec lift( } final Query query = druidQuery.getQuery(); - final List dimensions; - final List aggregators; - final VirtualColumnsAndFilter virtualColumnsAndFilter; if (query instanceof GroupByQuery) { final GroupByQuery groupBy = (GroupByQuery) query; - dimensions = groupBy.getDimensions(); - aggregators = groupBy.getAggregatorSpecs(); - virtualColumnsAndFilter = new VirtualColumnsAndFilter( - groupBy.getVirtualColumns(), - groupBy.getDimFilter(), - groupBy.getIntervals() + if (!CollectionUtils.isNullOrEmpty(groupBy.getSubtotalsSpec())) { + // Subtotal groupings produce rows of several different grouping shapes; a stored grouping has exactly one. + throw invalid( + projectionName, + "its body uses GROUPING SETS (or ROLLUP/CUBE), which a projection cannot store: each grouping is a" + + " different row shape. Store each grouping as its own projection instead" + ); + } + return new PlannedAggregation( + groupBy.getDimensions(), + groupBy.getAggregatorSpecs(), + new VirtualColumnsAndFilter( + groupBy.getVirtualColumns(), + groupBy.getDimFilter(), + groupBy.getIntervals() + ) ); } else if (query instanceof TimeseriesQuery) { // GROUP BY () plans to a timeseries over all time; it has no grouping columns. final TimeseriesQuery timeseries = (TimeseriesQuery) query; - dimensions = Collections.emptyList(); - aggregators = List.of(timeseries.getAggregatorSpecs().toArray(new AggregatorFactory[0])); - virtualColumnsAndFilter = new VirtualColumnsAndFilter( - timeseries.getVirtualColumns(), - timeseries.getFilter(), - timeseries.getIntervals() + return new PlannedAggregation( + Collections.emptyList(), + List.of(timeseries.getAggregatorSpecs().toArray(new AggregatorFactory[0])), + new VirtualColumnsAndFilter( + timeseries.getVirtualColumns(), + timeseries.getFilter(), + timeseries.getIntervals() + ) ); } else { throw invalid( @@ -529,14 +717,6 @@ private static AggregateProjectionSpec lift( "its body did not plan to an aggregation. Add a GROUP BY clause, or use SELECT DISTINCT" ); } - - return AggregateProjectionSpec - .builder(projectionName) - .virtualColumns(virtualColumnsAndFilter.virtualColumns) - .filter(virtualColumnsAndFilter.filter(projectionName)) - .groupingColumns(groupingColumns(projectionName, dimensions)) - .aggregators(renameToOutputNames(projectionName, druidQuery, aggregators)) - .build(); } private static List groupingColumns( diff --git a/sql/src/test/java/org/apache/druid/sql/calcite/CalciteCatalogDdlTest.java b/sql/src/test/java/org/apache/druid/sql/calcite/CalciteCatalogDdlTest.java index f6d65b5f4d76..d4c926bb9435 100644 --- a/sql/src/test/java/org/apache/druid/sql/calcite/CalciteCatalogDdlTest.java +++ b/sql/src/test/java/org/apache/druid/sql/calcite/CalciteCatalogDdlTest.java @@ -31,6 +31,7 @@ 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.RollupTableBaseTableMetadata; import org.apache.druid.catalog.model.TableBaseTableMetadata; import org.apache.druid.catalog.model.TableId; import org.apache.druid.catalog.model.TableMetadata; @@ -42,6 +43,8 @@ import org.apache.druid.guice.BuiltInTypesModule; import org.apache.druid.java.util.common.JodaUtils; import org.apache.druid.java.util.common.granularity.Granularities; +import org.apache.druid.query.aggregation.AggregatorFactory; +import org.apache.druid.query.aggregation.LongSumAggregatorFactory; import org.apache.druid.query.filter.RangeFilter; import org.apache.druid.segment.VirtualColumn; import org.apache.druid.segment.VirtualColumns; @@ -1058,21 +1061,188 @@ public void testBaseProjectionClusteringMustBeLeadingPrefix() } @Test - public void testBaseProjectionRejectsFilterOrGrouping() + public void testBaseProjectionRejectsFilter() { - for (String body : new String[]{ - "SELECT tenant, __time WHERE tenant <> 'x'", - "SELECT tenant, __time GROUP BY tenant, __time" - }) { - final DruidException e = assertThrows( - DruidException.class, - () -> execute( - "CREATE TABLE tbl SEALED (tenant VARCHAR, __time TIMESTAMP, PROJECTION __base AS (" + body + "))" - ), - body - ); - assertTrue(e.getMessage().contains("filters or groups"), e.getMessage()); - } + final DruidException e = assertThrows( + DruidException.class, + () -> execute( + "CREATE TABLE tbl SEALED (tenant VARCHAR, __time TIMESTAMP," + + " PROJECTION __base AS (SELECT tenant, __time WHERE tenant <> 'x'))" + ) + ); + assertTrue(e.getMessage().contains("its body filters"), e.getMessage()); + } + + /** + * A {@code __base} body with {@code GROUP BY} declares the rollup layout: the body is the query that rolls the + * table up onto itself. Grouping columns reference themselves, {@code TIME_FLOOR(__time, )} becomes the + * query granularity carrier, and each aggregate becomes the aggregator filling the declared metric column it is + * named for, reading the column it fills, since aggregating a stored rollup table is done with the combining form + * of its aggregators. + */ + @Test + public void testCreateTableWithRollupBaseProjection() + { + execute( + "CREATE TABLE tbl (tenant VARCHAR, __time TIMESTAMP, total BIGINT," + + " PROJECTION __base AS (SELECT tenant, TIME_FLOOR(__time, 'PT1H') AS __time, SUM(total) AS total" + + " GROUP BY 1, 2))" + ); + assertEquals( + new RollupTableBaseTableMetadata( + VirtualColumns.create( + Granularities.toVirtualColumn(Granularities.HOUR, Granularities.GRANULARITY_VIRTUAL_COLUMN_NAME) + ), + new AggregatorFactory[]{new LongSumAggregatorFactory("total", "total")}, + null + ), + WRITER.calls.get(0).spec.properties().get(DatasourceDefn.BASE_TABLE_PROPERTY) + ); + } + + /** + * Grouping on {@code __time} directly (no TIME_FLOOR) is a rollup at {@code NONE} granularity, and a body with no + * aggregates at all is a dedup table: rows with identical grouping values collapse into one. + */ + @Test + public void testRollupBaseProjectionWithBareTimeAndNoAggregates() + { + execute( + "CREATE TABLE tbl (tenant VARCHAR, __time TIMESTAMP," + + " PROJECTION __base AS (SELECT tenant, __time GROUP BY tenant, __time))" + ); + assertEquals( + new RollupTableBaseTableMetadata(null, null, null), + WRITER.calls.get(0).spec.properties().get(DatasourceDefn.BASE_TABLE_PROPERTY) + ); + } + + @Test + public void testRollupBaseProjectionRejectsClusteredBy() + { + final DruidException e = assertThrows( + DruidException.class, + () -> execute( + "CREATE TABLE tbl (tenant VARCHAR, __time TIMESTAMP, total BIGINT," + + " PROJECTION __base AS (SELECT tenant, __time, SUM(total) AS total GROUP BY 1, 2 CLUSTERED BY tenant))" + ) + ); + assertTrue( + e.getMessage().contains("a base table is either clustered or rollup, not both"), + e.getMessage() + ); + assertTrue(WRITER.calls.isEmpty()); + } + + @Test + public void testRollupBaseProjectionRejectsComputedGroupingColumn() + { + final DruidException e = assertThrows( + DruidException.class, + () -> execute( + "CREATE TABLE tbl (tenant VARCHAR, __time TIMESTAMP, total BIGINT," + + " PROJECTION __base AS (SELECT UPPER(tenant) AS tenant, __time, SUM(total) AS total GROUP BY 1, 2))" + ) + ); + assertTrue( + e.getMessage().contains( + "its column [tenant] is computed by an expression. Computed columns are not supported; compute the column" + + " at ingestion time instead" + ), + e.getMessage() + ); + assertTrue(WRITER.calls.isEmpty()); + } + + /** + * The body plans against the declared columns, so an aggregate reading an undeclared column cannot resolve: the + * rollup body is the table's self-rollup query, so an aggregate reads the metric column it fills. + */ + @Test + public void testRollupBaseProjectionRejectsAggregateOverUndeclaredColumn() + { + final DruidException e = assertThrows( + DruidException.class, + () -> execute( + "CREATE TABLE tbl (tenant VARCHAR, __time TIMESTAMP, total BIGINT," + + " PROJECTION __base AS (SELECT tenant, __time, SUM(cnt) AS total GROUP BY 1, 2))" + ) + ); + assertTrue(e.getMessage().contains("cnt"), e.getMessage()); + assertTrue(WRITER.calls.isEmpty()); + } + + /** + * Subtotal groupings produce rows of several different grouping shapes, while a stored grouping has exactly one, so + * GROUPING SETS (and ROLLUP/CUBE, which plan to the same form) are rejected rather than silently storing only the + * full grouping tuple. The guard is shared with aggregate projections, so both spellings are covered. + */ + @Test + public void testGroupingSetsRejected() + { + final DruidException base = assertThrows( + DruidException.class, + () -> execute( + "CREATE TABLE tbl (tenant VARCHAR, __time TIMESTAMP, total BIGINT," + + " PROJECTION __base AS (SELECT tenant, __time, SUM(total) AS total" + + " GROUP BY GROUPING SETS ((tenant, __time), (__time))))" + ) + ); + assertTrue(base.getMessage().contains("GROUPING SETS"), base.getMessage()); + + final DruidException projection = assertThrows( + DruidException.class, + () -> execute( + "CREATE TABLE tbl (tenant VARCHAR, page VARCHAR, __time TIMESTAMP, cnt BIGINT," + + " PROJECTION subtotals AS (SELECT tenant, page, SUM(cnt) AS total" + + " GROUP BY GROUPING SETS ((tenant, page), (tenant))))" + ) + ); + assertTrue(projection.getMessage().contains("GROUPING SETS"), projection.getMessage()); + assertTrue(WRITER.calls.isEmpty()); + } + + /** + * A rollup table's aggregators must combine their own output (the same aggregators serve ingestion and + * re-aggregation of stored rows), so an aggregate that is not its own combining form is rejected: a count is stored + * by summing a count column, not with {@code COUNT(*)}. + */ + @Test + public void testRollupBaseProjectionRejectsNonSelfCombiningAggregate() + { + final DruidException e = assertThrows( + DruidException.class, + () -> execute( + "CREATE TABLE tbl (tenant VARCHAR, __time TIMESTAMP, cnt BIGINT," + + " PROJECTION __base AS (SELECT tenant, __time, COUNT(*) AS cnt GROUP BY 1, 2))" + ) + ); + assertTrue( + e.getMessage().contains("aggregator [cnt] is not its own combining form"), + e.getMessage() + ); + assertTrue(WRITER.calls.isEmpty()); + } + + /** + * The declared type of a metric column must match what its aggregator stores; the check runs in the metadata's + * createSpec, attributed to the statement. + */ + @Test + public void testRollupBaseProjectionRejectsMetricTypeMismatch() + { + final DruidException e = assertThrows( + DruidException.class, + () -> execute( + "CREATE TABLE tbl (tenant VARCHAR, __time TIMESTAMP, total DOUBLE," + + " PROJECTION __base AS (SELECT tenant, __time, COUNT(*) AS total GROUP BY 1, 2))" + ) + ); + assertTrue( + e.getMessage().contains("metric column [total] is declared as type [DOUBLE], but its aggregator produces [LONG]"), + e.getMessage() + ); + assertTrue(WRITER.calls.isEmpty()); } /**