Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
77 changes: 66 additions & 11 deletions docs/development/extensions-core/catalog.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 `<expr> AS <name>` 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, <period>)` 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, <period>)`, 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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I would put this closer to the bottom, at least underneath rollup tables. Generally I expect rollup tables to be more popular than clustered tables. The order in which we present things should reflect how popular we expect them to be.

leading columns of the table. An item written as `<expr> AS <name>` makes that column computed at ingest time, from the
columns it reads:

```sql
CREATE TABLE "druid"."events" SEALED (
Expand All @@ -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.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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, <period>)} 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, <period>)} declares the table's query
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -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()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -112,6 +112,12 @@ public Granularity getQueryGranularity()
return granularitySpec.getQueryGranularity();
}

@Override
public boolean isRollup()
{
return granularitySpec.isRollup();
}

@Override
public BaseTableProjectionSpec withQueryGranularity(@Nullable Granularity queryGranularity)
{
Expand Down
Loading
Loading