Skip to content

feat: ddl statements + catalog - #19830

Merged
clintropolis merged 31 commits into
apache:masterfrom
clintropolis:catalog-ddl-statements
Sep 16, 2026
Merged

clintropolis merged 31 commits into
apache:masterfrom
clintropolis:catalog-ddl-statements

Conversation

@clintropolis

@clintropolis clintropolis commented Jul 30, 2026 •

Copy link
Copy Markdown
Member

Description

This PR adds DDL statement support to Druid so that CREATE and ALTER statements can be used to manage the catalog contents. This functionality requires WRITE permissions to the datasource (same as the backing catalog APIs), and is gated behind a new runtime property druid.sql.planner.enableCatalogDdl which is false by default. This PR includes support for defining the logical schema (columns list in catalog), aggregate projections (projections property in catalog), base table projection (baseTable property in catalog), clustering (clusterKeys in catalog), time partitioning (segmentGranularity in catalog), and any other properties (sealed, targetSegmentRows) with generic property setting syntax.

These operations are purely metadata operations, and so DROP TABLE has been omitted since in my mind it has the most room for confusion (a DROP TABLE only cleared metadata and didn't do anything else sounds kind of confusing to me). Dropping columns and projections from a table felt more easily explainable that these do not modify any existing data and instead stage the schema for subsequent ingestion (and future work should wire this catalog stuff better into compaction/reindexing so that we can also frame it such that compaction/reindexing will begin to eventually converge on the updated schema).

These all go through the regular query path, and return 0 rows on successful operation. Some follow-up work is needed to improve web-console syntax highlighting and behaviors, and since there is no drop table statement it would still probably be nice to eventually add a catalog management ui, but the statements at least already work as-is through the current query interface.

Some examples:

create table:

CREATE TABLE "wikipedia" (
  "channel" VARCHAR,
  "__time" TIMESTAMP,
  "page" VARCHAR,
  "namespace" VARCHAR,
  "user" VARCHAR,
  "comment" VARCHAR,
  "added" BIGINT,
  "delta" BIGINT
)
PARTITIONED BY DAY
CLUSTERED BY "channel"

alter table to add column:

ALTER TABLE "wikipedia" ADD COLUMN "deleted" BIGINT

alter table to add projection:

ALTER TABLE "wikipedia" ADD PROJECTION "channel_sums" AS (
  SELECT 
    TIME_FLOOR("__time", 'PT1H'),
    "channel",
    SUM("added") as "sum_added",
    SUM("delta") as "sum_delta",
    SUM("deleted") as "sum_deleted"
  GROUP BY 1,2
)

alter table set property:

ALTER TABLE "wikipedia" SET PROPERTIES (sealed = true)

alter table set 'base table' projection to create clustered segments:

ALTER TABLE "wikipedia" ADD PROJECTION __base AS (
  SELECT
    "channel",
    "__time",
    "page",
    "namespace,"
    "user",
    "comment",
    "added",
    "delta",
    "deleted"
  CLUSTERED BY "channel"
)

create table with base table and projection definitions:

CREATE TABLE "wikipedia" SEALED (
  "channel" VARCHAR,
  "__time" TIMESTAMP,
  "page" VARCHAR,
  "namespace" VARCHAR,
  "user" VARCHAR,
  "comment" VARCHAR,
  "added" BIGINT,
  "delta" BIGINT,
  "deleted" BIGINT,
  PROJECTION __base AS (
    SELECT
      "channel",
      "__time",
      "page",
      "namespace",
      "user",
      "comment",
      "added",
      "delta",
      "deleted"
    CLUSTERED BY "channel"
  ),
  PROJECTION "channel_sums" AS (
    SELECT 
      TIME_FLOOR("__time", 'PT1H'),
      "channel",
      SUM("added") as "sum_added",
      SUM("delta") as "sum_delta",
      SUM("deleted") as "sum_deleted"
    GROUP BY 1,2
  ),
  PROJECTION "channel_page_max" AS (
    SELECT 
      TIME_FLOOR("__time", 'PT1H'),
      "channel",
      "page",
      MAX("added") as "sum_added",
      MAX("delta") as "sum_delta",
      MAX("deleted") as "sum_deleted"
    GROUP BY 1,2,3
  )
)
PARTITIONED BY DAY
CLUSTERED BY channel

Comment thread sql/src/main/java/org/apache/druid/sql/calcite/parser/DruidSqlAlterTable.java Dismissed
Comment thread sql/src/main/java/org/apache/druid/sql/calcite/parser/DruidSqlAlterTable.java Dismissed
Comment thread sql/src/main/java/org/apache/druid/sql/calcite/parser/DruidSqlAlterTable.java Dismissed
Comment thread sql/src/main/java/org/apache/druid/sql/calcite/parser/DruidSqlAlterTable.java Dismissed
Comment thread sql/src/main/java/org/apache/druid/sql/calcite/parser/DruidSqlAlterTable.java Dismissed
Comment thread sql/src/main/java/org/apache/druid/sql/calcite/parser/DruidSqlCreateTable.java Dismissed
Comment thread sql/src/main/java/org/apache/druid/sql/calcite/parser/DruidSqlPropertyAssignment.java Dismissed
Comment thread sql/src/main/java/org/apache/druid/sql/calcite/parser/SqlProjectionSpec.java Dismissed
Comment thread sql/src/main/java/org/apache/druid/sql/calcite/parser/DruidSqlAlterTable.java Dismissed
Comment thread sql/src/main/java/org/apache/druid/sql/calcite/parser/DruidSqlAlterTable.java Dismissed

@FrankChen021 FrankChen021 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Severity Findings
P0 0
P1 2
P2 3
P3 0
Total 5

Reviewed 42 of 42 changed files.


This is an automated review by Codex GPT-5.6-Sol

@Override
protected void execute(CatalogTableWriter writer)
{
writer.updateProperties(tableId, properties);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

[P1] Revalidate the complete table after property edits

This uses the property-only edit endpoint, whose transaction loads and validates properties without columns, so DatasourceDefn.validate(ResolvedTable) and its cross-field checks never run. For example, after defining a DAY projection, SET PROPERTIES can change segmentGranularity to PT1H even though full validation rejects a projection coarser than its segments; it can likewise clear sealed while __base remains. The catalog then contains an invalid specification and subsequent ingestion fails. Load and validate the complete revised TableSpec inside the Coordinator transaction.

@Override
protected void execute(CatalogTableWriter writer)
{
writer.updateColumns(tableId, Collections.singletonList(column));

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

[P1] Enforce column operation predicates atomically

UpdateColumns appends a column when its name is absent, so ALTER TABLE t ALTER COLUMN typo SET DATA TYPE BIGINT silently adds typo instead of rejecting the nonexistent target. ADD COLUMN has the inverse predicate checked by a separate Broker read, allowing concurrent ADDs to both pass and the later merge to overwrite the first type. Add/alter existence semantics need dedicated checks inside the Coordinator's column-update transaction.

engine,
sql,
query,
CONTEXT,

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

[P2] Preserve statement context when planning projections

The nested planner receives only the hard-coded CONTEXT and discards the enclosing statement context, even though SET clauses are explicitly supported before DDL. For example, SET sqlTimeZone = 'America/Los_Angeles' followed by a projection using TIME_FLOOR stores a UTC expression, so the equivalent query under the same context plans differently and cannot match the projection. Merge relevant outer PlannerContext values before applying the deterministic overrides.

)
);
}
final VirtualColumn virtualColumn = planned.getVirtualColumn(selected.get(i));

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

[P2] Handle direct aliases in base projections

ScanQuery deduplicates its column list, while outputNames retains every SELECT item. A valid declared layout [id, copy] with SELECT id, id AS copy therefore has two outputs but only one selected entry, and the second iteration throws IndexOutOfBoundsException; a direct alias also has no virtual column to materialize copy. Preserve the select-to-source mapping or reject this form with a user-facing validation error.

}
name.unparse(writer, leftPrec, rightPrec);

final SqlWriter.Frame frame = writer.startList("(", ")");

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

[P2] Do not unparse omitted columns as empty parentheses

The grammar permits CREATE TABLE tbl PARTITIONED BY DAY with no parenthesized element list, but unparse always emits CREATE TABLE tbl () PARTITIONED BY DAY. Empty parentheses cannot be parsed because AddDruidTableElement is mandatory once '(' is present, so a valid AST does not round-trip. Omit the frame when both lists are empty or teach the grammar to accept ().

@jtuglu1

jtuglu1 commented Aug 6, 2026

Copy link
Copy Markdown
Contributor

cc @maytasm

@FrankChen021 FrankChen021 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Severity Findings
P0 0
P1 1
P2 1
P3 0
Total 2

Reviewed 44 of 44 changed files. Found two current-head correctness issues around projection type compatibility and atomic __base updates.


This is an automated review by Codex GPT-5.6-Luna(max)

}

final String granularity = table.stringProperty(SEGMENT_GRANULARITY_PROPERTY);
DataSchema.validateProjections(

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

[P1] ALTER COLUMN can leave projections type-incompatible

Projection validation checks referenced-column presence but does not reconcile projection dimension and aggregator types with current column definitions. A table with x BIGINT and SUM(x) can accept ALTER COLUMN x SET DATA TYPE VARCHAR, leaving an incompatible projection that later fails during initialization and can break ingestion. Validate compatibility or reject incompatible column edits before commit.


if (BASE_PROJECTION_NAME.equals(projectionName)) {
// The base table is a property of the table, not one of its projections, so it is set rather than appended.
if (existing.spec().properties().get(DatasourceDefn.BASE_TABLE_PROPERTY) != null) {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

[P2] Base projection existence check is not atomic

The __base branch performs the IF NOT EXISTS or duplicate check from a prior table read, then calls updateProperties without repeating that predicate inside the Coordinator transaction. Concurrent ADD PROJECTION __base statements can both observe absence and commit, with the later update overwriting the earlier projection; the non-IF form can also succeed instead of reporting a duplicate. Move the check into the atomic update.

@FrankChen021 FrankChen021 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Severity Findings
P0 0
P1 1
P2 1
P3 0
Total 2

Reviewed 44 of 44 changed files.


This is an automated review by Codex GPT-5.6-Luna(max)

if (baseTable == null) {
throw CatalogException.badRequest("A base table layout is required");
}
return catalog.tables().updateProperties(id, table -> {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

[P1] Base-table edits are not actually serialized

updateProperties performs an unlocked read followed by an update of the properties blob, with no row lock or version predicate. Concurrent ADD [IF NOT EXISTS] PROJECTION __base calls can both observe no layout and let the latter overwrite the former; a concurrent column edit can also leave __base referencing removed or retyped columns. Use row locking or compare-and-set updates and validate against the committed version.

requireSealed(Boolean.TRUE.equals(existing.spec().properties().get(DatasourceDefn.SEALED_PROPERTY)));
writer.setBaseTable(
tableId,
translateBaseTable(handlerContext, tableId.name(), columns, alterTable.getProjection()),

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

[P2] IF NOT EXISTS still plans an existing base layout

The handler translates and plans the supplied base-table query before the Coordinator evaluates ifNotExists. If __base already exists, ALTER TABLE ... ADD IF NOT EXISTS PROJECTION __base AS (...) should be a no-op, but an invalid or missing column in the supplied body can still make it fail. Short-circuit known existing layouts before translation while retaining the server-side race check.

@FrankChen021 FrankChen021 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Severity Findings
P0 0
P1 2
P2 0
P3 0
Total 2

Reviewed 45 of 45 changed files.


This is an automated review by Codex GPT-5.6-Luna(max)

.bind(COLUMNS_COL, JacksonUtils.toBytes(jsonMapper, revised.columns()))
.bind(blobColumn, JacksonUtils.toBytes(jsonMapper, blob.apply(revised)))
.bind(UPDATE_TIME_COL, updateTime)
.bind(OLD_VERSION_PARAM, existing.updateTime())

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

[P1] CAS version can repeat

The new compare-and-set update uses System.currentTimeMillis() as the replacement version, while the predicate compares the row against existing.updateTime(). Two edits in the same millisecond can therefore reuse the same version: both transactions may pass the old-version predicate, overwrite one another, and publish indistinguishable cache versions. Generate a strictly increasing version atomically (or use an equivalent database-side mechanism) before relying on this CAS.

{
String path = tablePath(TABLE_EDIT_PATH, tableId);
FutureUtils.getUnchecked(
serviceClient.asyncRequest(

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

[P1] Retries can replay catalog mutations

This new mutating POST goes through a client configured with StandardRetryPolicy.maxAttempts(6). If the Coordinator commits the DDL but the response is lost, a retryable channel error or 5xx can replay the non-idempotent edit (and the analogous create request), producing conflicts or duplicate effects while the SQL caller sees failure. Disable retries for writes or add request idempotency/deduplication.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

producing conflicts or duplicate effects
I don't think this can really happen with the compare and set stuff, so mainly the risk is a confusing error message? Re: idempotency, the SQL supports IF NOT EXISTS/IF EXISTS, so i don't think we want it to be implicitly idempotent.

Also this retrying isn't really new in this PR, createTable was already going through it prior to this PR.. maybe in a follow-up we could consider tightening up retry on write operations, since I think all it really gives us is a more accurate error message?

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

I agree that compare-and-set prevents a retry from overwriting a concurrent update, and that the retry policy predates this PR for createTable. The remaining issue is the new editTable POST: ADD COLUMN, ADD PROJECTION without IF NOT EXISTS, and base-table edits without existence flags all use the same retrying CatalogClient. If the Coordinator commits the first request but its response is lost, the retry reads the new version and runs the guard again, returning an error such as already exists even though the DDL committed. CAS prevents stale concurrent writes; it does not deduplicate post-commit replays. Please track this explicitly as a write-retry/idempotency follow-up, or disable retries for mutating calls. Reviewed 11 of 45 changed files.

@FrankChen021 FrankChen021 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

I have reviewed the code for correctness, edge cases, concurrency, and integration risks; no new issues found.

Reviewed 45 of 45 changed files.

Validation: git diff --check on the current 45-file PR diff; tests and builds not run.


This is an automated review by Codex GPT-5.6-Luna(max)

@capistrant capistrant left a comment

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.

This is a very cool addition to Druid, thank you for all the work and iteration that has gone into the PR. I reviewed test coverage and docs extensively to understand the features. I support the overall design and impl decisions spec'd out. I appreciate the detailed doc and the coverage that spans UTs and Embedded tests to provide confidence in the fact that the feature does what is spec'd out in the PR description and docs.

The many rounds with @FrankChen021 seem to have also worked through a lot of potential gotchas and have contributed to confidence in bugs being smoked out in this process versus in live testing after merge.

I'm approving with some nits I found reading through all the test files and javadocs. I also took into consideration that this is an experimental feature defaulted to off that is part of an experimental extension. release notes should stress this rawness and solicit bug reports/feedback from any early adopters.

Comment on lines +565 to +575
@SuppressWarnings("unchecked")
private DatasourceProjectionMetadata projection(int index)
{
return ((List<DatasourceProjectionMetadata>) WRITER.calls.get(0).spec.properties().get("projections")).get(index);
}

private String projectionsJson() throws Exception
{
return queryFramework().queryJsonMapper()
.writeValueAsString(WRITER.calls.get(0).spec.properties().get("projections"));
}

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.

nit: don't interleave these in the tests

Comment on lines +336 to +343
private List<DatasourceProjectionMetadata> projectionsOf(String tableName)
{
return TestHelper.JSON_MAPPER.convertValue(
client.readTable(TableId.datasource(tableName))
.spec().properties().get(DatasourceDefn.PROJECTIONS_KEYS_PROPERTY),
new TypeReference<List<DatasourceProjectionMetadata>>() {}
);
}

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.

nit: another interleaved private that could be moved out of test area

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.

### SQL DDL

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.

nice doc, ty for adding it up front. greatly helped set stage for review

Comment on lines +613 to +620
private String columnType(String tableName, String columnName) throws CatalogException
{
return catalog.tables().read(TableId.datasource(tableName)).spec().columns().stream()
.filter(c -> columnName.equals(c.name()))
.findFirst()
.orElseThrow(() -> new AssertionError("No column [" + columnName + "]"))
.dataType();
}

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.

nit: don't interleave with tests here. sorry to be repeating so many of these 😅

}

/**
* The Druid type of a column, which for {@code __time} is always {@link ColumnType#LONG} whatever was declared.

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.

The end reads confusing to me. are you meaning "regardless of what was declared" or something in that direction

@FrankChen021 FrankChen021 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

I have reviewed the code for correctness, edge cases, concurrency, and integration risks; no issues found.

Reviewed 45 of 45 changed files.


This is an automated review by Codex GPT-5.6-Luna(max)

@gianm gianm left a comment

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.

Reviewed only the documentation for now.

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

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.

never create, modify, or delete segments

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?)

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The 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).

ingestion.

These statements are disabled by default. Set `druid.sql.planner.enableCatalogDdl` to `true` on the Broker to enable
them. They require `WRITE` permission on the datasource, the same permission the catalog API requires, so enabling

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.

Can you make these require both READ and WRITE? This is a change I'm proposing to #20075 for INSERT and REPLACE, on the grounds that users without READ can't see a table and therefore can't write to it.


```sql
CREATE [OR REPLACE] TABLE [IF NOT EXISTS] <table>
[ ( { <column> <type> | PROJECTION <name> AS ( <select> ) } [, ...] ) ]

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.

This one line may be clearer if formatted on multiple lines.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

split out into a separate 'element' declaration to try to be a bit easier on the eyes

[ ( { <column> <type> | PROJECTION <name> AS ( <select> ) } [, ...] ) ]
[ PARTITIONED BY <granularity> ]
[ CLUSTERED BY <column> [, ...] ]
[ SEALED ]

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.

Feels like this should go closer to the column list. Maybe this?

CREATE TABLE tbl SEALED ( columns ... )


### SQL DDL

Tables can be defined with SQL instead of by posting a table specification. `CREATE TABLE` and `ALTER TABLE` are

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.

It would be great to have SHOW CREATE TABLE or DESCRIBE TABLE too, otherwise people will find these APIs hard to use without referring back to the JSON. This can be done in a follow up.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The 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

A projection body accepts a select list, an optional `WHERE` and an optional `GROUP BY`. It cannot use `ORDER BY`,
`LIMIT` or `HAVING`: a projection's ordering follows its grouping columns and is not something you choose. It also
cannot use joins, subqueries, or expressions computed over aggregates. Store the aggregates instead: `SUM(x)` and
`COUNT(x)` rather than `AVG(x)`.

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.

This doc point about AVG(x) doesn't really make sense unless you know that Druid internally converts AVG to SUM / COUNT. Not all databases do this, so a better example would be an explicit expression like SUM(x) / COUNT(x). There is also a list of aggregation functions that can be used at ingest time at docs/multi-stage-query/concepts.md#rollup, which you can link to.

```

Both take effect for subsequent ingestion. Segments already built keep whatever projections they were built with, so
dropping a projection does not rewrite data.

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.

But will it trigger compaction?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The 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

`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.

A computed column is written by the expression, not by the ingestion query, so an `INSERT` must supply 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.

"computed column is written by the expression, not by the ingestion query" to me is unclear wording. The example helps though. Maybe try:

Computed columns, like bucket in the example above, are computed based on inputs provided by INSERT or REPLACE. In the example above, the INSERT or REPLACE command should provide user_id, not bucket.


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.

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.

Re: "defining it makes the table a clustered table": Is it an error to provide __base without CLUSTERED BY?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The 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.

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

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.

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?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The 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)

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.

This could probably use pattern matching to be cooler.

}

/**
* Validate the revised spec as a whole before it is written back. Every edit is a read-modify-write of one part of

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 don't think there's really a need to justify the concept of validating the entire table as a unit.

@Override
public PlannerResult plan()
{
execute(handlerContext.plannerContext().getPlannerToolbox().catalogTableWriter());

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.

It's sketchy to execute the operation in plan(). It should happen only when PlannerResult#run is called. Please also include a test that when the DDL operation is unauthorized, no methods are called on the catalog client.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

good catch, fixed

{
try {
final TableMetadata table = client.table(tableId);
cache.updated(

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.

Is this just needed so the Broker can immediately see its own writes? I guess it should eventually update its cache even without this code.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

this is reworked a bit from the resync method added to address other problems using update

final TableMetadata table = client.table(tableId);
cache.updated(
new UpdateEvent(
table == null ? UpdateEvent.EventType.DELETE : UpdateEvent.EventType.UPDATE,

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 think doing every change as an UPDATE is going to confuse the CachedMetadataCatalog. It will cause computeUpdate to run, which expects to see a pre-existing object and does a LOG.error if it doesn't see one. In some cases it won't see one, like CREATE TABLE.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

reworked this a bit into to use resync to be cool for this usage

SqlStatementHandler.HandlerContext handlerContext = new HandlerContextImpl();

if (query instanceof DruidSqlCreateTable || query instanceof DruidSqlAlterTable) {
// The grammar does not admit EXPLAIN of a DDL statement; this guards the case anyway, since a DDL statement

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.

Please include a unit test for this, if it does not already exist. (I haven't checked yet.)

if (explain != null) {
throw InvalidSqlInput.exception("EXPLAIN is not supported for [%s]", query.getKind());
}
return createDdlHandler(handlerContext, query);

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 think structuring it this way will cause DDL to be accepted by all engines. It makes me wonder what we want. I think what we want is for the Task engine, which is meant to be async/offline, to reject DDL. I think we want the interactive ones (native, Dart) to accept it. Consider adding EngineFeature#CAN_DDL (and maybe consolidate CAN_INSERT and CAN_REPLACE into CAN_DML while you're at it) so we can do this.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

done, though didn't refactor CAN_DML, can do that in a follow-up

// The version is also the compare-and-set token, so it must actually change on every write: a commit
// landing in the same millisecond as the write that produced the version it read would otherwise
// leave the token as it was, letting a second writer's predicate match after this one commits.
final long updateTime = Math.max(System.currentTimeMillis(), existing.updateTime() + 1);

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.

Should become a helper function that gets applied to all the UPDATE_TIME_COL usages. There's a few others.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

done, also drive creates through here too instead of using current time inline

* extension and defaults to {@link #NOT_AVAILABLE} when that extension is absent.
* <p>
* The methods are deliberately semantic rather than a generic "apply this edit request", so that the extension's edit
* request types need not be visible here. Each corresponds to exactly one atomic Coordinator operation, which is why

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'm not sure this is actually better than moving the extension's edit types into core. With the approach currently in the PR, the edit types don't need to be here, but we have a 1-1 coupling anyway where there needs to be one method here for each edit type. Seems like a judgment call though. Thoughts?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

left this for now

@github-actions github-actions Bot added Area - Batch Ingestion Area - MSQ For multi stage queries - https://github.com/apache/druid/issues/12262 labels Sep 4, 2026

@FrankChen021 FrankChen021 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Severity Findings
P0 0
P1 0
P2 2
P3 0
Total 2
Severity Findings
P0 0
P1 0
P2 2
P3 0
Total 2

Reviewed 54 of 54 changed files.


This is an automated review by Codex GPT-5.6-Luna(max)

Comment thread sql/src/main/codegen/includes/ddl.ftl Outdated
final SqlProjectionSpec projection;
}
{
LOOKAHEAD(3)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

[P2] Avoid misclassifying projection-named complex columns

CREATE TABLE t (projection TYPE('COMPLEX<json>')) is a valid column declaration: PROJECTION is intentionally non-reserved and TYPE('...') is the supported Druid-native type escape hatch. With LOOKAHEAD(3), the token sequence PROJECTION TYPE ( selects DruidProjectionDefinition (TYPE is also a non-reserved Calcite keyword), which then expects SELECT and rejects the declaration. This makes valid schema syntax depend on the column name; distinguish the complete projection production or special-case TYPE after a projection identifier.

"WHERE schemaName = :schemaName\n" +
" AND name = :name\n";
" AND name = :name\n" +
" AND updateTime = :oldVersion";

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

[P2] Keep CAS edits restricted to active rows

The new compare-and-set update for updateProperties (and the analogous UPDATE_COLUMNS_STMT) checks only schema/name/updateTime; it omits the state = 'A' predicate used by REPLACE_SPEC_STMT. markDeleting changes state to D without a version predicate and assigns nextVersion(0), so if deletion and an edit occur in the same millisecond, the deleting row can retain the version the edit read and the edit still rewrites the DELETING row and emits an update event. TableMetadata.forUpdate reports that event as ACTIVE, violating the existing "can't update when deleting" contract and letting caches observe metadata after deletion. Add AND state = 'A' to both incremental update statements or make deletion participate in the same version guard.

@FrankChen021 FrankChen021 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

No new issues found in the changes since the previous MergeLens review; existing inline findings remain tracked separately.

Reviewed 54 of 54 changed files.


This is an automated review by Codex GPT-5.6-Luna(max)

user_id BIGINT,
payload TYPE('COMPLEX<json>'),
PROJECTION __base AS (
SELECT tenant, ABS(user_id) % 128 AS bucket, __time, user_id, payload

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 don't think ABS(user_id) % 128 is valid Druid SQL. (Should use the MOD function.)

ALTER TABLE <table> ADD COLUMN <column> <type>
ALTER TABLE <table> DROP COLUMN <column>
ALTER TABLE <table> ALTER COLUMN <column> SET DATA TYPE <type>
ALTER TABLE <table> ADD [IF NOT EXISTS] PROJECTION <name> AS ( <select> )

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.

Please change it to ADD PROJECTION [IF NOT EXISTS] <name>, as this is more commonly the way such syntaxes are done in SQL systems. It also better matches DROP PROJECTION [IF EXISTS].

+ " PROJECTION recent AS (SELECT page, SUM(cnt) AS total"
+ " WHERE __time >= TIMESTAMP '2020-01-01 00:00:00' GROUP BY page))"
);
assertNotNull(projection(0).getSpec().getFilter(), "time filter must survive as a filter");

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.

Isn't it bad that time filters are retained as filters? Typically when a real SQL query is actually issued, __time filters will be moved to intervals, so they will not actually show up in the filter. So, does that mean they won't be able to match the stored projection? Or, is there something that causes this to be handled well?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

Good eye, there is not currently logic that can handle matching __time filters in projections against query intervals, so this is currently modeling building an unreachable projection.

Claude came up with this, I left it because i think retaining them as filters makes more sense than trying to give a standalone interval to projections, so i view at as the right form and since this test is just covering projection construction and not query time matching it seemed harmless to leave. That is, assuming I add support for special handling projections with __time filters that are comparable to intervals as a follow-up, which I think it would be pretty straightforward to add, so I have it on my list.

+ " PROJECTION by_tenant AS (SELECT tenant, SUM(v) AS sum_v GROUP BY tenant))"
);
final TableSpec spec = WRITER.calls.get(0).spec;
assertNotNull(spec.properties().get(DatasourceDefn.BASE_TABLE_PROPERTY));

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.

This test class has a lot of assertNotNull that really should be asserting equality to the expected thing. Please tighten them up.

"Projection [%s] references column [%s], which table [%s] does not declare",
spec.getName(),
required,
table.spec().type()

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.

Should be table name, not type, in this position.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

table name isn't actually available here but i think that is probably ok since I think there could be at most 1 table involved in these statements

* The reserved name of the base-table projection, which describes the physical layout of the table itself. Handled
* as a separate catalog property, not as one of the aggregate projections.
*/
public static final String BASE_PROJECTION_NAME = "__base";

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.

Use Projections.BASE_TABLE_PROJECTION_NAME?

// The TYPE('...') escape hatch names a Druid native type. Parse and validate rather than passing the raw
// string downstream, where a malformed type string would silently resolve to a different type, and return the
// canonical form.
if (target == Target.EXTERNAL && !StringUtils.toLowerCase(simpleName).startsWith("complex<")) {

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.

Any reason to not allow non-complex types here? Seems like it would be fine to allow arrays and primitives.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

not sure why it was like this (this was migrated from Externals to here), I don't see any problem allowing any native type in that function, but I have not changed this in this PR yet, I think i'd rather save for a follow-up if we want to do that.

throw unsupportedType(name, dataType);
}
final ColumnType nativeType = ColumnType.fromString(simpleName);
if (nativeType == null) {

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.

Should check that the complex type is valid, or else this will cause failures later when we attempt to apply the table spec, leading to a broken table.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

Initially I was trying to avoid really changing anything here (this logic was basically ported from Externals.java), I would assume the validation wasn't here in part to handle cases where like the coordinator doesn't have some extension loaded that provides the type but presumably whatever is actually using the catalog definition to execute does. However, this PR made it so that in some cases at least we already do this validation from any projection stuff used since that kicks in a validation of those specs which are computed from these definitions, so I went ahead and changed things to just do the validation for datasource columns too. I updated the catalog docs to mention that the coordinator needs all extensions loaded for complex types defined in a catalog schema.

@clintropolis
clintropolis merged commit 9f5363c into apache:master Sep 16, 2026
27 checks passed
@clintropolis
clintropolis deleted the catalog-ddl-statements branch September 16, 2026 23:10
@github-actions github-actions Bot added this to the 39.0.0 milestone Sep 16, 2026
aho135 added a commit to aho135/druid that referenced this pull request Sep 17, 2026
The merge of master (apache#19830, catalog DDL) dropped the closing brace of
isPlanningTimeLimited() and the opening Javadoc delimiter of the following
isEnableCatalogDdl() method, so PlannerConfig.java no longer compiled. Restore
both.
aho135 added a commit that referenced this pull request Sep 27, 2026
* Add a timeout for SQL query planning on the Broker

An extremely complex SQL query (for example one with a very large IN
filter) can spend tens of seconds in the Calcite planning phase on the
Broker. A flood of such queries can occupy all of the Broker's request
threads and effectively freeze the process, degrading the whole cluster.

This adds a wall-clock timeout on the planning phase (the phase measured
by the sqlQuery/planningTimeMs metric):

- New runtime property druid.sql.planner.maxPlanningTimeMs (default 0,
  meaning no timeout), also overridable per query via the
  maxPlanningTimeMs query context key.
- PlannerContext now owns a Calcite CancelFlag, which PlannerFactory
  wires into the framework Context so that Calcite aborts at its next
  cancellation checkpoint when the flag is tripped.
- A new SqlPlanningTimeout watchdog is armed in DirectStatement.plan().
  When the deadline is exceeded it trips the CancelFlag and interrupts
  the planning thread; the failure is surfaced as a QueryTimeoutException
  (HTTP 504). The watchdog clears the interrupt on close so it cannot
  leak to a pooled request thread.

Adds unit tests for the watchdog, the new config (including the query
context override), and an end-to-end DirectStatement planning-timeout
test, plus documentation.

* Make planning-timeout code comments more concise

* Remove cancelled planning-timeout tasks from the scheduler queue

Address review feedback: build the scheduler as a ScheduledThreadPoolExecutor
with setRemoveOnCancelPolicy(true) so that cancelled watchdog tasks (the common
case, when planning finishes before the deadline) are removed from the queue
immediately instead of lingering until their delay elapses.

* Address review: propagate cancel flag to Hep/view planners, tighten timeout

Addresses PR review feedback beyond the scheduler queue fix:

- Propagate the per-query CancelFlag to the Hep-based planning stages
  (pre/reduction/pre-Volcano/cleanup). Calcite's Programs.of builds its
  HepPlanner with a null Context, so it never observed the flag; the new
  CalciteRulesManager#hepProgram helper builds the HepPlanner with a
  Context carrying the flag, otherwise identical to Programs.of. This is
  the stage where RexSimplify (CALCITE-3178) runs, so it must be
  cancellable.

- Share the top-level query's CancelFlag with nested planners created
  during view expansion (DruidViewMacro), via
  PlannerContext#withInheritedCancelFlag, so one planning timeout governs
  the whole session.

- Reject a plan that completes after the watchdog fired (e.g. in a
  non-cancellable section that swallowed the interrupt) instead of
  returning it: check isTimedOut() on the success path too.

- Measure the planning budget from planningStartNanos so time spent
  constructing the planner counts against maxPlanningTimeMs and a query
  cannot get a fresh full budget afterwards.

- Emit maxPlanningTimeMs from PlannerConfig#getNonDefaultAsQueryContext()
  so the defensive config<->context round-trip check passes for a
  non-default value; add a round-trip test.

* Enforce planning deadline across construction; drop deprecated Calcite call

Address review follow-ups:

- Enforce the planning deadline against wall-clock time elapsed since
  planningStartNanos, which includes createPlanner() (planner/schema
  construction). If construction alone exhausts maxPlanningTimeMs, fail
  with QueryTimeoutException before arming the watchdog instead of
  scheduling a 1ms timer and letting a fast follow-on plan slip through.
  Also re-check the wall-clock deadline (not just the watchdog flag) on
  the success path, so a plan that completes after the deadline under
  scheduler jitter is still rejected. Adds a unit test that exhausts the
  budget during planner construction.

- Drop the deprecated HepPlanner#registerMetadataProviders call from the
  cancel-flag-aware hepProgram helper; modern Calcite routes rule
  metadata through the cluster provider (cluster.setMetadataProvider),
  which the helper already sets, matching Programs.of behavior. Verified
  against the decoupled/join/array/select Calcite query tests.

* Close watchdog race, cover subquery rewriting, restore Hep metadata provider

Address the second automated review round (3x P2):

- Watchdog last-millisecond race (DirectStatement): compute the remaining
  budget once and, if the timeout is enabled but already exhausted, fail
  immediately. Previously the deadline check and the budget recompute were
  separate steps; if time crossed the deadline in between, arm() received a
  non-positive budget and returned the disabled (no-op) watchdog, so a query
  past the deadline got no cancel-flag trip or interrupt and was only caught
  after planning returned.

- Subquery rewriting cancellation (CalciteRulesManager): buildPreProgram used
  Programs.subQuery(...), whose HepPlanners are built with a null context and
  never see the query CancelFlag. Replaced with buildSubQueryProgram(), a
  cancel-aware replica of Programs.subQuery's two Hep stages, so a query that
  spends its budget in subquery rewriting is interruptible.

- Restore Hep metadata provider (CalciteRulesManager): the previous change
  dropped HepPlanner#registerMetadataProviders, which also dropped the
  HepRelMetadataProvider that method prepends (needed for metadata lookups on
  HepRelVertex during Hep planning). Restored it and suppressed the
  deprecation warning with a comment explaining why it must match Programs.of.

* docs: clarify that a non-positive maxPlanningTimeMs disables the timeout

A negative value is coalesced to disabled everywhere the setting is consumed
(all guards use <= 0 / > 0), same as 0. Document that non-positive disables the
timeout, in the runtime-property table, the query-context table, and the
PlannerConfig javadoc.

* test: use imports instead of fully-qualified class names in SqlStatementTest

Address review: import DruidPlanner, PlannerResult, CancelFlag, and JUnit
Timeout in the planning-timeout tests rather than referencing them by
fully-qualified name inline.

* Fix merge conflict resolution in PlannerConfig

The merge of master (#19830, catalog DDL) dropped the closing brace of
isPlanningTimeLimited() and the opening Javadoc delimiter of the following
isEnableCatalogDdl() method, so PlannerConfig.java no longer compiled. Restore
both.

* Drop druid.sql.planner.maxPlanningTimeMs in favor of default query context

Per review: since dynamic query-context configuration exists, a per-query
setting should not also be exposed as a dedicated druid.sql.planner.* runtime
property (the existing planner properties predate that feature). Remove the
@JsonProperty binding so maxPlanningTimeMs is configured only via the query
context; a cluster-wide default is set with druid.query.default.context.maxPlanningTimeMs
(or the Broker's dynamic queryContext config). The context key, watchdog, and
per-query override are unchanged.

Docs: remove the runtime-property row from the Broker SQL config table, update
the maxPlanningTimeMs query-context entry to point at "Overriding default query
context values", and add that link to the sql-query-context "Learn more" list.

* Trim verbose javadoc on PlannerConfig#maxPlanningTimeMs

* Trim PlannerConfig#maxPlanningTimeMs javadoc to what/why only

* test: use assertThrows in planning-timeout tests

Replace the try/plan()/fail()/catch pattern with Assertions.assertThrows in
testPlanningTimeout and testPlanningTimeoutDuringPlannerConstruction, matching
the file's existing assertThrows usages. Behavior unchanged; still asserts the
QueryTimeoutException message and closes the statement in a finally.

* Drop Hep cancel-flag duplication; rely on Volcano coverage for the timeout

Thread-dumping a planning-bound large-IN query (forced OR path) showed the cost
is in VolcanoPlanner.findBestExp (Druid-convention conversion + memo digest
computation), not the Hep reduction stage as originally assumed. The Volcano
planner already receives the per-query CancelFlag via the framework Context, so
the timeout covers that path without any custom Hep plumbing.

Revert the hepProgram/buildSubQueryProgram helpers that reimplemented Calcite's
Programs.of/Programs.subQuery just to inject the CancelFlag into the internally
constructed HepPlanners. Those duplicated Calcite internals (and a deprecated
API) to cover a Hep-reduction-bound scenario that could not be reproduced. Back
to plain Programs.of/Programs.subQuery. The framework-Context wiring and the
view-planner flag inheritance (PlannerContext#withInheritedCancelFlag) remain.

* Trim verbose comments across the planning-timeout changes
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

Area - Batch Ingestion Area - Documentation Area - MSQ For multi stage queries - https://github.com/apache/druid/issues/12262 Area - Querying Area - SQL

Projects

None yet

Development

Successfully merging this pull request may close these issues.

6 participants