Repository navigation
feat: ddl statements + catalog #19830
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
fbc2cbf
ae9c072
b0617ff
c71e15f
00ccd65
aac737a
2b89172
9fb88d7
19fc09a
62d0d7c
957ad33
1dbc471
93a677c
68d530e
b60a480
7d6155e
3fd5d7b
e21b985
3e8505c
c3fa0f6
5781f20
1390ecb
b6d2c82
74af982
2539ac2
610a2d2
c9eb48c
d6f2129
8e24d46
5128bc0
2579237
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -43,6 +43,212 @@ allowing queries to be more concise, and simpler to write. This also allows the | |
| written into a defined column of the table is consistent with that columns definition, minimizing errors where unexpected | ||
| data is written into a particular column of the table. | ||
|
|
||
| ### Effects of a table definition | ||
|
|
||
| A table definition takes effect when it is read, not when it is written: changing one, whether through SQL DDL or the | ||
| REST API, rewrites nothing. What it affects: | ||
|
|
||
| - **SELECT queries** report each declared column with its declared type, in declared order. Segments that store a | ||
| column as a different physical type are converted to the declared type at query time. Columns present in segments | ||
| but not declared in the catalog remain queryable, and follow the declared columns with their physical types. | ||
| - **INSERT and REPLACE** (the catalog applies only to SQL-based ingestion) validate against the definition rather | ||
| than merely defaulting from it. A query column targeting a declared column must be convertible to the declared type | ||
| without changing it: inserting a `BIGINT` value into a column declared `VARCHAR` stores the value as a string, while | ||
| inserting a `VARCHAR` value into a column declared `BIGINT` is an error. Columns the query produces that the table | ||
| does not declare are rejected when the table is [`sealed`](#table-properties) and ingested normally otherwise. | ||
| Table properties such as `segmentGranularity` and `clusterKeys` act as defaults that an individual statement may | ||
| override, such as with its own `PARTITIONED BY`. | ||
| - **Streaming and native batch ingestion** do not consult the catalog; their specs are unaffected by any table | ||
| definition. | ||
|
|
||
| ### SQL DDL | ||
|
|
||
| Tables can be defined with SQL instead of by posting a table specification. `CREATE TABLE` and `ALTER TABLE` are | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. It would be great to have
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. cool, yea that seems nice, can look into that in a follow-up |
||
| submitted to the Broker like any other SQL statement, and write the same catalog metadata the REST API does. They | ||
| return no rows. | ||
|
|
||
| These statements change catalog metadata only. They never create, modify, or delete segments: defining a table does | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
Well, they don't immediately do these things synchronously, but they may trigger these things to happen in the background as part of compaction. This is worth mentioning. (Or do they? What's the current state of catalog-based compaction?)
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. catalog based compaction doesn't currently handle baseTable spec, but once in place it would apply all of the other things like projections and clustering and stuff. I'm going to omit catalog stuff from the docs for now and address this stuff in a follow-up branch I have in progress that wires that up so that it is consistent with the rest of the catalog and (also make catalog compaction config work with MSQ compaction since that didn't work either). |
||
| not ingest anything, and altering a column does not rewrite existing data. Changes take effect for subsequent queries | ||
| and ingestion, as described in [Effects of a table definition](#effects-of-a-table-definition). | ||
|
|
||
| These statements are disabled by default. Set `druid.sql.planner.enableCatalogDdl` to `true` on the Broker to enable | ||
| them. They require both `READ` and `WRITE` permission on the datasource, the same permissions the catalog API | ||
| requires for its write operations, so enabling them lets anyone who can ingest into a datasource also change its | ||
| catalog definition; leave them disabled if you manage catalog entries with your own tooling. The setting cannot be | ||
| overridden per query. | ||
|
|
||
| The `druid-catalog` extension must be loaded on both the Broker and the Coordinator; without it, these statements | ||
| report that the extension is not available. | ||
|
|
||
| ```sql | ||
| CREATE [OR REPLACE] TABLE [IF NOT EXISTS] <table> | ||
| [ [ SEALED ] ( <table element> [, ...] ) ] | ||
| [ PARTITIONED BY <granularity> ] | ||
| [ CLUSTERED BY <column> [, ...] ] | ||
|
|
||
| <table element> ::= | ||
| <column> <type> | ||
| | PROJECTION <name> AS ( <select> ) | ||
| ``` | ||
|
|
||
| `OR REPLACE` replaces the specification of an existing table; `IF NOT EXISTS` leaves an existing table unchanged. | ||
| The two cannot be combined. `PARTITIONED BY` sets [`segmentGranularity`](#table-properties) and `CLUSTERED BY` sets | ||
| `clusterKeys`, both of which a later `INSERT` or `REPLACE` inherits unless it states its own. `SEALED` sets | ||
| [`sealed`](#table-properties), which requires every ingested column to be declared. It is written just before the | ||
| column list because that is what it describes: the list is the table's whole schema, so the statement does not accept | ||
| `SEALED` without one. (The `sealed` property itself can still be set on any table through `SET PROPERTIES` or the | ||
| REST API.) | ||
|
|
||
| Note that the table-level `CLUSTERED BY` is a sort order applied to each ingestion, which is a different thing from | ||
| the `CLUSTERED BY` inside a [`__base` projection](#the-base-table), which defines how segments physically group rows. | ||
|
|
||
| Column types are written as SQL types, such as `VARCHAR`, `BIGINT`, `DOUBLE`, or `VARCHAR ARRAY`. The `__time` column | ||
| is written as `TIMESTAMP`. Types that have no SQL spelling, such as complex types, use `TYPE('...')` with the Druid | ||
| native type string: | ||
|
|
||
| ```sql | ||
| CREATE TABLE "druid"."visits" ( | ||
| __time TIMESTAMP, | ||
| user_id VARCHAR, | ||
| pages_visited BIGINT, | ||
| sketch TYPE('COMPLEX<thetaSketch>') | ||
| ) | ||
| PARTITIONED BY DAY | ||
| CLUSTERED BY user_id | ||
| ``` | ||
|
|
||
| A complex type must be provided by a loaded extension: the statement fails if the named type is not registered, so a | ||
| typo or a missing extension surfaces when the table is defined rather than at the first ingestion or query that | ||
| touches the column. The extension must be loaded on the Coordinator, which validates catalog writes. | ||
|
|
||
| `ALTER TABLE` supports one change per statement, so that each statement is a single atomic catalog operation: | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. There should be some commentary here about what happens after you do an
As you can tell from the parenthetical questions I am not really sure what happens today. Maybe this is discussed elsewhere in the catalog docs. It probably should be, since it applies to JSON manipulations too. If it is, this section should link to the place where it is discussed.
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. added a section about this, but skipped for compaction which I will address in a follow-up re other comment. |
||
|
|
||
| ```sql | ||
| ALTER TABLE <table> ADD COLUMN <column> <type> | ||
| ALTER TABLE <table> DROP COLUMN <column> | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Presumably, if unsealed,
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. documented that it is just a metadata operation and the implications, there is |
||
| ALTER TABLE <table> ALTER COLUMN <column> SET DATA TYPE <type> | ||
| ALTER TABLE <table> ADD PROJECTION [IF NOT EXISTS] <name> AS ( <select> ) | ||
| ALTER TABLE <table> DROP PROJECTION [IF EXISTS] <name> | ||
| ALTER TABLE <table> SET PROPERTIES ( <property> = <value> [, ...] ) | ||
| ``` | ||
|
|
||
| `ADD COLUMN` fails if the column already exists, and `ALTER COLUMN` fails if it does not, so a misspelled column name | ||
| is reported rather than quietly creating or replacing a column. Each statement is also checked against the rest of the | ||
| table definition, not only the part it changes: adding a column, changing a type, or setting a property is rejected if | ||
| the resulting table would be invalid, such as a segment granularity coarser than a projection the table declares. | ||
|
|
||
| `DROP COLUMN` removes the column's declaration, not its data. A table's SQL schema is the declared columns followed by | ||
| any other columns present in its segments, so a dropped column that still has data remains queryable: it loses its | ||
| declared type and position and appears after the declared columns, like any column the catalog does not know about. | ||
| Whether it can still be ingested then follows the usual rule: rejected if the table is sealed, accepted as an | ||
| undeclared column otherwise. To remove a column from query results without rewriting data, use the catalog's | ||
| `hiddenColumns` property instead, which currently has no SQL spelling and is set through the REST API. | ||
|
|
||
| #### Projections | ||
|
|
||
| A table may declare [projections](../../querying/projections.md), which are pre-aggregated views stored inside each | ||
| segment. A projection is written as a `SELECT` over the table's own columns, with no `FROM` clause: | ||
|
|
||
| ```sql | ||
| CREATE TABLE "druid"."visits" ( | ||
| __time TIMESTAMP, | ||
| user_id VARCHAR, | ||
| user_agent VARCHAR, | ||
| pages_visited BIGINT, | ||
| PROJECTION daily_by_agent AS ( | ||
| SELECT TIME_FLOOR(__time, 'P1D'), user_agent, SUM(pages_visited) AS total_pages | ||
| WHERE user_agent IS NOT NULL | ||
| GROUP BY 1, 2 | ||
| ) | ||
| ) | ||
| PARTITIONED BY DAY | ||
| ``` | ||
|
|
||
| The body is planned exactly as the equivalent query would be, so a projection matches the queries it was written to | ||
| serve. Every aggregate needs an alias, which becomes the name of the stored column. Time granularity is expressed | ||
| with `TIME_FLOOR`, as it would be in a query. | ||
|
|
||
| Because the body is planned like a query, it is planned under the statement's own query context, including any `SET` | ||
| clauses. Only context parameters that affect planning can change the stored definition; parameters that only affect | ||
| query execution have no effect, since the body is planned and stored rather than run. The context itself is not part | ||
| of the definition: what the catalog stores is the projection the body planned to, so nothing from the statement's | ||
| context is carried over to queries that later use it. Note also that a projection is matched to a query by its shape, | ||
| so a definition planned under a context that changes that shape only matches queries run under the same context. | ||
|
|
||
| A projection body accepts a select list, an optional `WHERE` and an optional `GROUP BY`. It cannot use `ORDER BY`, | ||
| `LIMIT` or `HAVING`: a projection's ordering follows its grouping columns and is not something you choose. It also | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Hopefully these will be syntax errors with nice error messages. Do we have tests asserting what the error messages are if you do various incorrect things?
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. there is test coverage for this, the error messages aren't super specific, they are like parser's generic syntax errors instead of 'projections cannot use limit' or whatever (we would have to allow them in the grammar and then reject them later in the translator for that to work) |
||
| cannot use joins, subqueries, or expressions computed over aggregates, such as `SUM(x) / COUNT(x)`: store the two | ||
| aggregates as separate columns and divide at query time. A projection can store the same | ||
| [aggregation functions that are supported for rollup at ingestion time](../../multi-stage-query/concepts.md#rollup). | ||
|
|
||
| Projections may also be added to and removed from an existing table: | ||
|
|
||
| ```sql | ||
| ALTER TABLE "druid"."visits" ADD PROJECTION [IF NOT EXISTS] by_agent AS ( | ||
| SELECT user_agent, SUM(pages_visited) AS total_pages GROUP BY user_agent | ||
| ) | ||
| ALTER TABLE "druid"."visits" DROP PROJECTION [IF EXISTS] by_agent | ||
| ``` | ||
|
|
||
| Both take effect for subsequent ingestion. Segments already built keep whatever projections they were built with, so | ||
| dropping a projection does not rewrite data. | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. But will it trigger compaction?
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. not yet, will address this in follow-up that fixes up catalog compaction to honor baseTable spec and stuff |
||
|
|
||
| #### 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. | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Re: "defining it makes the table a clustered table": Is it an error to provide
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. it is right now, since only the clustered base table spec exists. I do plan on adding regular table and rollup-table specs eventually so that all of the v10 base table metadata is represented and everything can be fully expressed, but this will be done in follow-up PR(s) which can expand on the docs here. |
||
|
|
||
| Its body lists the columns in the order segments store them, so it must name every declared column, in declared | ||
| order. An item written as `<expr> AS <name>` makes that column computed at ingest time, from the columns it reads: | ||
|
|
||
| ```sql | ||
| CREATE TABLE "druid"."events" SEALED ( | ||
| tenant VARCHAR, | ||
| bucket BIGINT, | ||
| __time TIMESTAMP, | ||
| user_id BIGINT, | ||
| payload TYPE('COMPLEX<json>'), | ||
| PROJECTION __base AS ( | ||
| SELECT tenant, MOD(ABS(user_id), 128) AS bucket, __time, user_id, payload | ||
| CLUSTERED BY tenant, bucket | ||
| ) | ||
| ) | ||
| 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`. | ||
|
|
||
| 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. | ||
|
|
||
| `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. | ||
|
|
||
| Every other name beginning with `__` remains reserved. | ||
|
|
||
| #### Setting table properties | ||
|
|
||
| `SET PROPERTIES` merges the given [table properties](#table-properties) into the table. A value of `NULL` removes a | ||
| property. Values must be literals: | ||
|
|
||
| ```sql | ||
| ALTER TABLE "druid"."visits" SET PROPERTIES (targetSegmentRows = 3000000, sealed = TRUE) | ||
| ALTER TABLE "druid"."visits" SET PROPERTIES (sealed = NULL) | ||
| ``` | ||
|
|
||
| Both statements require `READ` and `WRITE` permission on the datasource, the same permissions the REST API checks. Table | ||
| names may be unqualified or qualified with the `druid` schema; other schemas are rejected. Names are case-sensitive. | ||
|
|
||
| There is no `DROP TABLE`. Deleting a table's catalog entry without deleting its data would be a surprising meaning | ||
| for the statement, so removing a specification is left to the [delete API](#delete-a-table) until the semantics are | ||
| settled. | ||
|
|
||
| ### API Objects | ||
|
|
||
| #### TableSpec | ||
|
|
@@ -90,10 +296,11 @@ The endpoint supports a set of optional query parameters to enforce optimistic l | |
| is meant to update a table rather than create a new one. In the default case, with no query parameters set, this request | ||
| will return an error if a table of the same name already exists in the schema specified. | ||
|
|
||
| | Parameter | Type | Description | | ||
| |-------------|---------|-------------------------------------------------------------------------------------------------------------------------------| | ||
| | `version` | Long | the expected version of an existing table. The version must match. If not (or if the table does not exist), returns an error. | | ||
| | `overwrite` | boolean | if true, then overwrites any existing table. Otherwise, the operation fails if the table already exists. | | ||
| | Parameter | Type | Description | | ||
| |---------------|---------|-------------------------------------------------------------------------------------------------------------------------------| | ||
| | `version` | Long | the expected version of an existing table. The version must match. If not (or if the table does not exist), returns an error. | | ||
| | `overwrite` | boolean | if true, then overwrites any existing table. Otherwise, the operation fails if the table already exists. | | ||
| | `ifNotExists` | boolean | if true, then leaves an existing table unchanged and reports a version of 0. Otherwise, the operation fails if the table already exists. | | ||
|
|
||
| ##### Responses | ||
|
|
||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
nice doc, ty for adding it up front. greatly helped set stage for review