Skip to content

[SPARK-34586][SQL] Support declaring a write distribution and ordering in CREATE/REPLACE TABLE - #58153

Open
anuragmantri wants to merge 8 commits into
apache:masterfrom
anuragmantri:SPARK-34586-write-distribution-ordering-create-table
Open

anuragmantri wants to merge 8 commits into
apache:masterfrom
anuragmantri:SPARK-34586-write-distribution-ordering-create-table

Conversation

@anuragmantri

@anuragmantri anuragmantri commented Aug 20, 2026 •

Copy link
Copy Markdown
Contributor

This PR is based on the previous works of @aokolnychyi and @RussellSpitzer. Both are co-authors on the commit.

What changes were proposed in this pull request?

Lets CREATE/REPLACE TABLE, and their AS SELECT forms, carry the write distribution and sort order the table should be written with:

[DISTRIBUTED BY PARTITION]
[[LOCALLY] ORDERED BY { (write_order_field, ...) | write_order_field, ... } | UNORDERED]

write_order_field: { col | transform({ col | constant }, ...) } [ASC|DESC] [NULLS {FIRST|LAST}]

The request travels to the catalog on TableInfo, via two new builder methods (withWriteDistributionMode, withWriteOrdering) and two accessors. A plain column in ORDERED BY reaches the catalog as a NamedReference, any other key as a Transform, and the ordering is never null. Table gains writeDistributionMode() / writeOrdering() so a catalog can report back what a table declares, DelegatingTable forwards both, and SHOW CREATE TABLE / DESCRIBE TABLE EXTENDED read them.

A catalog has to advertise the new TableCatalogCapability.SUPPORTS_CREATE_TABLE_WITH_WRITE_DISTRIBUTION_AND_ORDERING, otherwise the statement fails with UNSUPPORTED_FEATURE.TABLE_OPERATION before anything is created: during planning for a v2 catalog, and during analysis when the session catalog would create a v1 table. The request reaches the catalog only through the four TableInfo overloads (createTable and the three StagingTableCatalog.stage*), so a catalog that reports the capability must override each one it can be reached through. DelegatingCatalogExtension does not forward this capability from its delegate, because it does not forward those overloads.

Layers:

  • Grammar (4 new non-reserved keywords),
  • AstBuilder (which also rejects DISTRIBUTED BY PARTITION on an unpartitioned table),
  • Four CreateTable/ReplaceTable(AsSelect) plans, PreprocessTableCreation (normalizes the ordering) and CheckAnalysis (rejects unknown columns)
  • The create/replace execs, which build the TableInfo, and DataSourceV2Strategy, which checks the capability

A sort key has to resolve against the table's columns, so ORDERED BY needs a statement that defines a schema (a column list, typed partition columns, or AS SELECT); on a v2 catalog a schemaless statement is rejected with SPECIFY_WRITE_ORDERING_IS_NOT_ALLOWED, the position PARTITIONED BY already takes there. DISTRIBUTED BY PARTITION needs the statement to declare a partitioning (SPECIFY_DISTRIBUTED_BY_PARTITION_WITHOUT_PARTITIONING_IS_NOT_ALLOWED) and cannot be combined with CLUSTER BY (SPECIFY_CLUSTER_BY_WITH_DISTRIBUTED_BY_PARTITION_IS_NOT_ALLOWED); a cluster_by(...) transform in PARTITIONED BY counts as clustering, not partitioning.

writeDistributionMode is a new WriteDistributionMode enum (HASH/RANGE/NONE). null means the statement said nothing, which is distinct from NONE. null leaves the choice to the catalog's own default; NONE is an explicit request not to distribute. See the design decisions below.

Why are the changes needed?

Spark can already enforce a write layout: a connector reports one from RequiresDistributionAndOrdering on its Write, and DistributionAndOrderingUtils inserts the shuffle and sort. What is missing is a way for the user to author it, so it is persisted as table metadata and honored by the first write and by every later one, including writes from other engines.

Without it, a user who wants a sorted table has to run three statements: CREATE, then a connector-specific ALTER TABLE, then INSERT. That reaches the same physical layout but is not atomic: the table is visible unsorted in between, a concurrent writer can land unsorted data, and for REPLACE ... AS SELECT the table is temporarily empty. Tools that generate CTAS have no hook between create and load, so for them the first load can never be sorted.

Iceberg's community has asked for this in the engine's DDL twice apache/iceberg#3547 and apache/iceberg#14612

Does this PR introduce any user-facing change?

Yes, additive optional clauses on CREATE/REPLACE TABLE, and four new keywords (DISTRIBUTED, LOCALLY, ORDERED, UNORDERED), all non-reserved in every mode, so existing identifiers with those names keep working. A statement that does not use the clauses builds the same TableInfo as before.

CREATE TEMPORARY TABLE ... USING, CREATE MATERIALIZED VIEW and CREATE STREAMING TABLE reject the clauses with INVALID_STATEMENT_OR_CLAUSE.

New error conditions: WRITE_ORDERING_WITH_UNKNOWN_COLUMN (42703), SPECIFY_WRITE_ORDERING_IS_NOT_ALLOWED (42601), SPECIFY_DISTRIBUTED_BY_PARTITION_WITHOUT_PARTITIONING_IS_NOT_ALLOWED (42908) and SPECIFY_CLUSTER_BY_WITH_DISTRIBUTED_BY_PARTITION_IS_NOT_ALLOWED (42908). The JDBC getSQLKeywords and Thrift CLI_ODBC_KEYWORDS lists include the four new keywords.

Source compatibility for extensions and connectors:

  • Extensions that pattern-match on the CreateTable, ReplaceTable, CreateTableAsSelect and ReplaceTableAsSelect plans need to add the two new fields (writeDistributionMode, writeOrdering) to their patterns. The fields are defaulted, so constructing the plans by name still compiles, but extensions compiled against an earlier release must be recompiled. Moving the pair onto TableSpec to restore the signatures is SPARK-59943.
  • The create/replace exec nodes gain the same two parameters without defaults, so extensions that match or construct them positionally (for example a custom CTAS strategy) must add them. On CreateTableAsSelectExec and ReplaceTableAsSelectExec they come before the defaulted transaction.
  • Implementations of V2CreateTablePlan need to add its two new abstract members (writeOrdering, withWriteOrdering).
  • DelegatingCatalogExtension.capabilities() now returns the delegate's capabilities without SUPPORTS_CREATE_TABLE_WITH_WRITE_DISTRIBUTION_AND_ORDERING. A subclass that overrides the TableInfo overloads adds it back by overriding capabilities().
  • TableInfo rejects a null write ordering at construction.

Design decisions

Four choices here are worth spelling out, with what they cost.

1. A mode enum rather than a typed Distribution. Spark already has Distributions.unspecified() / clustered(...) / ordered(...), used by RequiresDistributionAndOrdering. But a mode is a policy for all future writes while a Distribution describes one write, so a typed Distribution would have to freeze the partition expressions at create time. The table therefore declares a WriteDistributionMode (HASH/RANGE/NONE), a closed set like SortDirection and NullOrdering, which maps directly to connector settings such as Iceberg's write.distribution-mode. TableInfo is a builder, so a typed withDistribution(Distribution) can be added later without breaking callers.

2. Table exposes what the table declares, and both display paths read it. Table.writeDistributionMode() / writeOrdering() default to null / empty, mirroring Table.constraints(). The contract, documented on the methods, is that these are a declared default for future writes and nothing more: an individual write may override it, a connector may narrow it (a hash distribution is meaningless on an unpartitioned table), RequiresDistributionAndOrdering on the Write stays authoritative for what a given write actually requires, and none of it claims anything about how the data already in the table is laid out. A scan reports that itself, via outputPartitioning/outputOrdering. SHOW CREATE TABLE reproduces the clauses from them and DESCRIBE TABLE EXTENDED reports them, which is what makes a table created with these clauses recreatable. Iceberg, for example, downgrades a requested hash to none on an unpartitioned table at write time, so a consumer of these accessors must not read them as what the next write will do. There is also no ALTER TABLE version yet. A connector that has one of its own (Iceberg's ALTER TABLE ... WRITE) is unaffected, but changing the declared default through Spark is follow-up work.

A connector can report more. For example, hash on a table with no partitioning (the parser rejects DISTRIBUTED BY PARTITION there), a range distribution with no ordering, or a sort key the syntax cannot spell. SHOW CREATE TABLE therefore emits the clauses only when the catalog accepts them, the statement would not create a v1 table (the session catalog with a provider that is not a v2 source, decided by the same rule CREATE TABLE uses), and parsing the rendered clauses with the session's parser gives back the same mode and keys, with every column they reference in the table. It retries once with every name quoted, for reserved keywords. Otherwise it omits the pair rather than emitting a clause that would mean something else or would not parse. Literals keep their type: a FLOAT renders as <v>F, a TIMESTAMP or nanosecond TIMESTAMP_LTZ in UTC with an explicit offset, and a TIME with its precision. DESCRIBE TABLE EXTENDED prints both values in every case, renders connector expressions as their describe does, and does not fail on values Catalyst cannot represent or on expressions Spark's SQL builder does not know.

CREATE TABLE ... LIKE deliberately does not copy the declared layout, and is not gated on the capability, because it is not a V2CreateTablePlan, so gating it would mean extending the check to a fifth plan. It does hand the source Table to the connector, which can carry the layout across itself, as the TableCatalog.createTableLike Javadoc says. Happy to fold LIKE in if reviewers would rather have it here.

3. ORDERED BY implies a distribution, and the mode is what says how far the order reaches.

clause distribution ordering
(none) unset, catalog's own default none
ORDERED BY (...) range as written
LOCALLY ORDERED BY (...) none as written
UNORDERED none none
DISTRIBUTED BY PARTITION hash none
DISTRIBUTED BY PARTITION [LOCALLY] ORDERED BY (...) hash as written
DISTRIBUTED BY PARTITION UNORDERED hash none

There is no separate global-vs-local flag on the recorded pair, because the distribution already is one: range means the order holds across the table, hash and none mean it holds within a write task. A bare ORDERED BY has to imply range: sorting within each task does not make the table sorted, so on its own it would record an order the writes cannot achieve. That is also why LOCALLY is the escape hatch for the within-task case, the same word Iceberg's ALTER TABLE ... WRITE LOCALLY ORDERED BY uses.

It also means the implication only applies when DISTRIBUTED BY PARTITION is absent. Beside it the distribution is already fixed and already local, so LOCALLY adds nothing and UNORDERED contributes only "no sort keys". Both are accepted and both produce the same pair as leaving them out. Spark's usual answer for a clause with no effect in a combination is to reject it (see SPECIFY_CLUSTER_BY_WITH_PARTITIONED_BY_IS_NOT_ALLOWED), and that would be defensible here too; they are accepted because both spellings are already valid in Iceberg's ALTER TABLE ... WRITE, and a word that means the same thing on CREATE TABLE as it does on ALTER TABLE seems worth more than the extra strictness.

The cost of the coupling: there is no way to say "range-distribute but record no ordering", nor "record this ordering and leave the distribution unset". Decoupling needs more syntax; There is no spelling for a range distribution on its own (DISTRIBUTED BY PARTITION is the hash one), so one would have to be invented, and the common cases would then take two clauses instead of one.

4. CLUSTER BY is independent, and it is not CLUSTERED BY ... INTO ... BUCKETS. DISTRIBUTED BY PARTITION requires the table to be partitioned, by PARTITIONED BY or by CLUSTERED BY ... INTO n BUCKETS (bucketing is a partition transform). CLUSTER BY (...) is a different clause: it records clustering columns for the data source to interpret, it is not a partition transform, and the grammar already forbids combining it with PARTITIONED BY or CLUSTERED BY ... INTO ... BUCKETS. So a table using CLUSTER BY has no partitioning at all, and DISTRIBUTED BY PARTITION on it is unsatisfiable by construction rather than by a policy choice here. The cost: CLUSTER BY users cannot ask for a per-partition write distribution without moving to PARTITIONED BY. If CLUSTER BY should count as a partitioning for this, that is a change to CLUSTER BY and belongs in its own PR.

CLUSTER BY can still be combined with ORDERED BY, LOCALLY ORDERED BY or UNORDERED: Spark passes the clustering columns and the declared distribution and ordering to the catalog without reconciling them, so the catalog decides how they interact and may reject a combination it does not support.

DISTRIBUTED BY PARTITION on a CLUSTER BY table gets its own error, SPECIFY_CLUSTER_BY_WITH_DISTRIBUTED_BY_PARTITION_IS_NOT_ALLOWED, which says to replace CLUSTER BY with PARTITIONED BY or bucketing. The parser and SHOW CREATE TABLE share one rule for what counts as partitioning: any transform not named cluster_by.

One pre-existing hole to be aware of: TableCatalog.createTable(ident, TableInfo)'s default implementation forwards to the deprecated 4-arg createTable(ident, columns, partitions, properties), so every TableInfo-only field is dropped for a catalog that implements only that overload. constraints already is, today. The capability check is what keeps that from becoming a silent wrong result here: without the capability the statement fails, so a catalog that never looks at TableInfo cannot quietly create a table lacking the requested layout. Fixing the default itself is out of scope.

How was this patch tested?

  • New CreateTableWriteOrderSuite (48 tests): parsing, analysis, what reaches the catalog (createTable and the stage* methods), the capability checks (including DelegatingCatalogExtension), a catalog rejecting a combination, the first write of a CTAS/RTAS planning the declared shuffle and sort, SHOW CREATE TABLE / DESCRIBE TABLE EXTENDED, and every new error condition with checkError and its SQLSTATE.
  • New WriteDistributionAndOrderingUtilsSuite (14 tests): table-driven, over connector-reported sort keys of every literal type and shape, asserting SHOW CREATE TABLE emits a key exactly when parsing it back gives the same key with every column it references in the schema, including under the parser confs that change literals and keywords; each mode's clause form; and DESCRIBE's rendering of connector expressions.
  • One table-driven test in PlanResolutionSuite over the resolved CREATE/CTAS/REPLACE/RTAS plans, and one in CreatePipelineDatasetAsSelectParserSuiteBase.
  • Updated keyword lists in keywords*.sql.out, SparkConnectDatabaseMetaDataSuite and ThriftServerWithSparkContextSuite.

Was this patch authored or co-authored using generative AI tooling?

Generated-by: Claude Code (Opus 5)

Co-authored-by: Peter Toth peter.toth@gmail.com
Co-authored-by: Anton Okolnychyi aokolnychyi@apache.org
Co-authored-by: Russell Spitzer russell.spitzer@gmail.com

@anuragmantri

Copy link
Copy Markdown
Contributor Author

@peter-toth @szehon-ho could you please review this? I would especially like feedback on the the design design decisions mentioned in the PR description. Thanks.

@peter-toth peter-toth 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.

Thanks for picking this up, @anuragmantri!

The shape reads right to me. The request rides on TableInfo, Table reports it back, and the capability check rejects the statement while planning, so a catalog that ignores TableInfo cannot hand back a table quietly missing the layout. On the design decisions you asked about: I'm a co-author, so my read of those isn't independent - all four still look right to me, and the value is in the findings below. Two things block. CI is red on a keyword-list test this misses, and SHOW CREATE TABLE can still emit DDL that does not parse - the ordering half of the hazard the hash-without-partitioning guard already covers. Findings 2, 3 and 5 are measured on this branch, not read off the diff.

Blocking

  • 1. Connect JDBC keyword list not updated: SparkConnectDatabaseMetaDataSuite."getSQLKeywords" asserts its own hardcoded keyword list and is the only failing test on this head. The Thrift-server sibling was updated; this one needs the same four keywords. [inline: sql/hive-thriftserver/src/test/scala/org/apache/spark/sql/hive/thriftserver/ThriftServerWithSparkContextSuite.scala:217]
  • 2. SHOW CREATE TABLE can emit DDL that does not parse: the guard covers the distribution mode but never checks that the sort expressions are spellable, and SortOrder.expression() is typed Expression. A table reporting sort(a + 1, ASC, NULLS FIRST) yields ORDERED BY (id + 1 ASC NULLS FIRST), which fails with PARSE_SYNTAX_ERROR on replay. [inline: sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/ShowCreateTableExec.scala:146]
  • 3. New error condition is named for cases it cannot report: nested struct columns are supported (the suite asserts ORDERED BY p.x succeeds), and the "or is in a map or array" clause is unreachable - those paths raise INVALID_FIELD_NAME. A released condition name can't be renamed. [inline: common/utils/src/main/resources/error/error-conditions.json:9155]

Non-blocking

  • 4. Docs credit the "data source" with a catalog capability: the gate is TableCatalog.capabilities(), so a reader following this paragraph will inspect the USING provider and find nothing there. [inline: docs/sql-ref-syntax-ddl-create-table-datasource.md:153]
  • 5. Docs promise SHOW CREATE TABLE reproduces the clauses, with no omission caveat: the caveat is in the PR body and in the scaladoc, but not on the page users read, and the omitted case is the likely one. [inline: docs/sql-ref-syntax-ddl-create-table-datasource.md:160]

Minor

  • 6. "written file" vs "write task" three lines apart: a task can write several files, so the two are different claims; the second is the accurate one. [inline: docs/sql-ref-syntax-ddl-create-table-datasource.md:133]

val infoValue = client.getInfo(sessionHandle, GetInfoType.CLI_ODBC_KEYWORDS)
// scalastyle:off line.size.limit
assert(infoValue.getStringValue == "ADD,AFTER,AGGREGATE,ALIGN,ALL,ALTER,ALWAYS,ANALYZE,AND,ANTI,ANY,ANY_VALUE,APPLY,APPROX,ARCHIVE,ARRAY,AS,ASC,ASENSITIVE,ASOF,AT,ATOMIC,AUTHORIZATION,AUTO,BEGIN,BERNOULLI,BETWEEN,BIGINT,BIN,BINARY,BINDING,BIN_DISTRIBUTE_RATIO,BIN_END,BIN_START,BOOLEAN,BOTH,BUCKET,BUCKETS,BY,BYTE,CACHE,CALL,CALLED,CASCADE,CASE,CAST,CATALOG,CATALOGS,CDC,CHANGE,CHANGES,CHAR,CHARACTER,CHECK,CLEAR,CLOSE,CLUSTER,CLUSTERED,CODEGEN,COLLATE,COLLATION,COLLATIONS,COLLECTION,COLUMN,COLUMNS,COMMENT,COMMIT,COMPACT,COMPACTIONS,COMPENSATION,COMPUTE,CONCATENATE,CONDITION,CONSTRAINT,CONTAINS,CONTINUE,COST,CREATE,CROSS,CUBE,CURRENT,CURRENT_DATABASE,CURRENT_DATE,CURRENT_PATH,CURRENT_SCHEMA,CURRENT_TIME,CURRENT_TIMESTAMP,CURRENT_USER,CURSOR,DATA,DATABASE,DATABASES,DATE,DATEADD,DATEDIFF,DATE_ADD,DATE_DIFF,DAY,DAYOFYEAR,DAYS,DBPROPERTIES,DEC,DECIMAL,DECLARE,DEFAULT,DEFAULT_PATH,DEFINED,DEFINER,DELAY,DELETE,DELIMITED,DESC,DESCRIBE,DETERMINISTIC,DFS,DIRECTORIES,DIRECTORY,DISTANCE,DISTINCT,DISTRIBUTE,DIV,DO,DOUBLE,DROP,ELSE,ELSEIF,EMPTY,END,ENFORCED,ERROR,ESCAPE,ESCAPED,EVOLUTION,EXACT,EXCEPT,EXCHANGE,EXCLUDE,EXCLUSIVE,EXECUTE,EXISTS,EXIT,EXPLAIN,EXPORT,EXTEND,EXTENDED,EXTERNAL,EXTRACT,FALSE,FETCH,FIELDS,FILEFORMAT,FILTER,FIRST,FLOAT,FLOW,FOLLOWING,FOR,FOREIGN,FORMAT,FORMATTED,FOUND,FROM,FULL,FUNCTION,FUNCTIONS,GENERATED,GEOGRAPHY,GEOMETRY,GLOBAL,GRANT,GROUP,GROUPING,HANDLER,HAVING,HISTORY,HOUR,HOURS,IDENTIFIED,IDENTIFIER,IDENTITY,IF,IGNORE,ILIKE,IMMEDIATE,IMPORT,IN,INCLUDE,INCLUSIVE,INCREMENT,INDEX,INDEXES,INNER,INPATH,INPUT,INPUTFORMAT,INSENSITIVE,INSERT,INT,INTEGER,INTERSECT,INTERVAL,INTO,INVOKER,IS,ITEMS,ITERATE,JOIN,JSON,JSON_EXISTS,JSON_TABLE,JSON_VALUE,KEY,KEYS,LANGUAGE,LAST,LATERAL,LAZY,LEADING,LEAVE,LEFT,LEVEL,LIKE,LIMIT,LINES,LIST,LOAD,LOCAL,LOCALTIME,LOCATION,LOCK,LOCKS,LOGICAL,LONG,LOOP,MACRO,MAP,MATCHED,MATCH_CONDITION,MATERIALIZED,MAX,MEASURE,MERGE,METRICS,MICROSECOND,MICROSECONDS,MILLISECOND,MILLISECONDS,MINUS,MINUTE,MINUTES,MODIFIES,MONTH,MONTHS,MSCK,NAME,NAMESPACE,NAMESPACES,NANOSECOND,NANOSECONDS,NATURAL,NEAREST,NEXT,NO,NONE,NORELY,NOT,NULL,NULLS,NUMERIC,OF,OFFSET,ON,ONLY,OPEN,OPTION,OPTIONS,OR,ORDER,ORDINALITY,OUT,OUTER,OUTPUTFORMAT,OVER,OVERLAPS,OVERLAY,OVERWRITE,PARTITION,PARTITIONED,PARTITIONS,PATH,PERCENT,PIVOT,PLACING,POSITION,PRECEDING,PRIMARY,PRINCIPALS,PROCEDURE,PROCEDURES,PROPERTIES,PURGE,QUALIFY,QUARTER,QUERY,RANGE,READ,READS,REAL,RECORDREADER,RECORDWRITER,RECOVER,RECURSION,RECURSIVE,REDUCE,REFERENCES,REFRESH,RELY,RENAME,REPAIR,REPEAT,REPEATABLE,REPLACE,RESET,RESPECT,RESTRICT,RETURN,RETURNING,RETURNS,REVOKE,RIGHT,ROLE,ROLES,ROLLBACK,ROLLUP,ROW,ROWS,SCD,SCHEMA,SCHEMAS,SECOND,SECONDS,SECURITY,SELECT,SEMI,SEPARATED,SEQUENCE,SERDE,SERDEPROPERTIES,SESSION_USER,SET,SETS,SHORT,SHOW,SIMILARITY,SINGLE,SKEWED,SMALLINT,SOME,SORT,SORTED,SOURCE,SPECIFIC,SQL,SQLEXCEPTION,SQLSTATE,START,STATISTICS,STORED,STRATIFY,STREAM,STREAMING,STRING,STRUCT,SUBSTR,SUBSTRING,SYNC,SYSTEM,SYSTEM_PATH,SYSTEM_TIME,SYSTEM_VERSION,TABLE,TABLES,TABLESAMPLE,TARGET,TBLPROPERTIES,TERMINATED,THEN,TIME,TIMEDIFF,TIMESTAMP,TIMESTAMPADD,TIMESTAMPDIFF,TIMESTAMP_LTZ,TIMESTAMP_NTZ,TINYINT,TO,TOUCH,TRACK,TRAILING,TRANSACTION,TRANSACTIONS,TRANSFORM,TRIM,TRUE,TRUNCATE,TRY_CAST,TYPE,UNARCHIVE,UNBOUNDED,UNCACHE,UNIFORM,UNION,UNIQUE,UNKNOWN,UNLOCK,UNNEST,UNPIVOT,UNSET,UNTIL,UPDATE,USE,USER,USING,VALUE,VALUES,VAR,VARCHAR,VARIABLE,VARIANT,VERSION,VIEW,VIEWS,VOID,WATERMARK,WEEK,WEEKS,WHEN,WHERE,WHILE,WIDTH,WINDOW,WITH,WITHIN,WITHOUT,X,YEAR,YEARS,ZONE")
assert(infoValue.getStringValue == "ADD,AFTER,AGGREGATE,ALIGN,ALL,ALTER,ALWAYS,ANALYZE,AND,ANTI,ANY,ANY_VALUE,APPLY,APPROX,ARCHIVE,ARRAY,AS,ASC,ASENSITIVE,ASOF,AT,ATOMIC,AUTHORIZATION,AUTO,BEGIN,BERNOULLI,BETWEEN,BIGINT,BIN,BINARY,BINDING,BIN_DISTRIBUTE_RATIO,BIN_END,BIN_START,BOOLEAN,BOTH,BUCKET,BUCKETS,BY,BYTE,CACHE,CALL,CALLED,CASCADE,CASE,CAST,CATALOG,CATALOGS,CDC,CHANGE,CHANGES,CHAR,CHARACTER,CHECK,CLEAR,CLOSE,CLUSTER,CLUSTERED,CODEGEN,COLLATE,COLLATION,COLLATIONS,COLLECTION,COLUMN,COLUMNS,COMMENT,COMMIT,COMPACT,COMPACTIONS,COMPENSATION,COMPUTE,CONCATENATE,CONDITION,CONSTRAINT,CONTAINS,CONTINUE,COST,CREATE,CROSS,CUBE,CURRENT,CURRENT_DATABASE,CURRENT_DATE,CURRENT_PATH,CURRENT_SCHEMA,CURRENT_TIME,CURRENT_TIMESTAMP,CURRENT_USER,CURSOR,DATA,DATABASE,DATABASES,DATE,DATEADD,DATEDIFF,DATE_ADD,DATE_DIFF,DAY,DAYOFYEAR,DAYS,DBPROPERTIES,DEC,DECIMAL,DECLARE,DEFAULT,DEFAULT_PATH,DEFINED,DEFINER,DELAY,DELETE,DELIMITED,DESC,DESCRIBE,DETERMINISTIC,DFS,DIRECTORIES,DIRECTORY,DISTANCE,DISTINCT,DISTRIBUTE,DISTRIBUTED,DIV,DO,DOUBLE,DROP,ELSE,ELSEIF,EMPTY,END,ENFORCED,ERROR,ESCAPE,ESCAPED,EVOLUTION,EXACT,EXCEPT,EXCHANGE,EXCLUDE,EXCLUSIVE,EXECUTE,EXISTS,EXIT,EXPLAIN,EXPORT,EXTEND,EXTENDED,EXTERNAL,EXTRACT,FALSE,FETCH,FIELDS,FILEFORMAT,FILTER,FIRST,FLOAT,FLOW,FOLLOWING,FOR,FOREIGN,FORMAT,FORMATTED,FOUND,FROM,FULL,FUNCTION,FUNCTIONS,GENERATED,GEOGRAPHY,GEOMETRY,GLOBAL,GRANT,GROUP,GROUPING,HANDLER,HAVING,HISTORY,HOUR,HOURS,IDENTIFIED,IDENTIFIER,IDENTITY,IF,IGNORE,ILIKE,IMMEDIATE,IMPORT,IN,INCLUDE,INCLUSIVE,INCREMENT,INDEX,INDEXES,INNER,INPATH,INPUT,INPUTFORMAT,INSENSITIVE,INSERT,INT,INTEGER,INTERSECT,INTERVAL,INTO,INVOKER,IS,ITEMS,ITERATE,JOIN,JSON,JSON_EXISTS,JSON_TABLE,JSON_VALUE,KEY,KEYS,LANGUAGE,LAST,LATERAL,LAZY,LEADING,LEAVE,LEFT,LEVEL,LIKE,LIMIT,LINES,LIST,LOAD,LOCAL,LOCALLY,LOCALTIME,LOCATION,LOCK,LOCKS,LOGICAL,LONG,LOOP,MACRO,MAP,MATCHED,MATCH_CONDITION,MATERIALIZED,MAX,MEASURE,MERGE,METRICS,MICROSECOND,MICROSECONDS,MILLISECOND,MILLISECONDS,MINUS,MINUTE,MINUTES,MODIFIES,MONTH,MONTHS,MSCK,NAME,NAMESPACE,NAMESPACES,NANOSECOND,NANOSECONDS,NATURAL,NEAREST,NEXT,NO,NONE,NORELY,NOT,NULL,NULLS,NUMERIC,OF,OFFSET,ON,ONLY,OPEN,OPTION,OPTIONS,OR,ORDER,ORDERED,ORDINALITY,OUT,OUTER,OUTPUTFORMAT,OVER,OVERLAPS,OVERLAY,OVERWRITE,PARTITION,PARTITIONED,PARTITIONS,PATH,PERCENT,PIVOT,PLACING,POSITION,PRECEDING,PRIMARY,PRINCIPALS,PROCEDURE,PROCEDURES,PROPERTIES,PURGE,QUALIFY,QUARTER,QUERY,RANGE,READ,READS,REAL,RECORDREADER,RECORDWRITER,RECOVER,RECURSION,RECURSIVE,REDUCE,REFERENCES,REFRESH,RELY,RENAME,REPAIR,REPEAT,REPEATABLE,REPLACE,RESET,RESPECT,RESTRICT,RETURN,RETURNING,RETURNS,REVOKE,RIGHT,ROLE,ROLES,ROLLBACK,ROLLUP,ROW,ROWS,SCD,SCHEMA,SCHEMAS,SECOND,SECONDS,SECURITY,SELECT,SEMI,SEPARATED,SEQUENCE,SERDE,SERDEPROPERTIES,SESSION_USER,SET,SETS,SHORT,SHOW,SIMILARITY,SINGLE,SKEWED,SMALLINT,SOME,SORT,SORTED,SOURCE,SPECIFIC,SQL,SQLEXCEPTION,SQLSTATE,START,STATISTICS,STORED,STRATIFY,STREAM,STREAMING,STRING,STRUCT,SUBSTR,SUBSTRING,SYNC,SYSTEM,SYSTEM_PATH,SYSTEM_TIME,SYSTEM_VERSION,TABLE,TABLES,TABLESAMPLE,TARGET,TBLPROPERTIES,TERMINATED,THEN,TIME,TIMEDIFF,TIMESTAMP,TIMESTAMPADD,TIMESTAMPDIFF,TIMESTAMP_LTZ,TIMESTAMP_NTZ,TINYINT,TO,TOUCH,TRACK,TRAILING,TRANSACTION,TRANSACTIONS,TRANSFORM,TRIM,TRUE,TRUNCATE,TRY_CAST,TYPE,UNARCHIVE,UNBOUNDED,UNCACHE,UNIFORM,UNION,UNIQUE,UNKNOWN,UNLOCK,UNNEST,UNORDERED,UNPIVOT,UNSET,UNTIL,UPDATE,USE,USER,USING,VALUE,VALUES,VAR,VARCHAR,VARIABLE,VARIANT,VERSION,VIEW,VIEWS,VOID,WATERMARK,WEEK,WEEKS,WHEN,WHERE,WHILE,WIDTH,WINDOW,WITH,WITHIN,WITHOUT,X,YEAR,YEARS,ZONE")

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.

Finding 1. This list got the four new keywords, but its Spark Connect JDBC sibling did not, and CI is red on it.

SparkConnectDatabaseMetaDataSuite."SparkConnectDatabaseMetaData getSQLKeywords" asserts its own hardcoded list at sql/connect/client/jdbc/src/test/scala/org/apache/spark/sql/connect/client/jdbc/SparkConnectDatabaseMetaDataSuite.scala:213. It is the only failing test on c6adea5c405 - one annotation on the "Report test results" check run, and the failing "Build modules: ... connect ..." job is the same test.

getSQLKeywords drops SQL:2003 reserved words, and all four new keywords are non-reserved, so all four need adding:

  • ...,DISTRIBUTE,DISTRIBUTED,DIV,...
  • ...,LOAD,LOCALLY,LOCATION,...
  • ...,OPTIONS,ORDERED,ORDINALITY,...
  • ...,UNLOCK,UNORDERED,UNPIVOT,...

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done.

val orderBy = if (table.writeOrdering().nonEmpty) {
Some(table.writeOrdering()
.map(WriteDistributionAndOrdering.describeSortOrder)
.mkString("ORDERED BY (", ", ", ")"))

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.

Finding 2. The hasPartitioning guard below covers the distribution side of this method's own doc - "emitting a clause that means something else -- or one that does not parse at all -- would be worse than emitting none" - but nothing checks that the sort expressions are spellable. SortOrder.expression() is typed Expression and Table.writeOrdering()'s new contract does not narrow it, so a connector may report an expression the writeOrderField : transform ... rule cannot represent.

Measured on this branch, with a table reporting mode range and sort(a + 1, ASCENDING, NULLS_FIRST):

CREATE TABLE p.t (
  id INT)
USING foo
ORDERED BY (id + 1 ASC NULLS FIRST)

Replaying that gives [PARSE_SYNTAX_ERROR] Syntax error at or near '+' at line 4, pos 15. The whole statement is unrunnable, not one clause lost, which is the worse failure the doc calls out.

Same fix shape as the hash case: drop the pair when it cannot be spelled. Two things to get right. CreateTableWriteOrderSuite builds a bare FieldReference, not a Transform, and that one does round-trip, so the predicate has to admit it. And nulling out orderBy alone is not enough - (none, None) would then emit UNORDERED, declaring no ordering on a table that has one - so the whole method has to bail:

  private def isSpellable(e: V2Expression): Boolean = e match {
    case _: NamedReference => true
    // Mirrors `transformArgument : qualifiedName | constant`.
    case t: Transform =>
      t.arguments().forall(a => a.isInstanceOf[NamedReference] || a.isInstanceOf[Literal[_]])
    case _ => false
  }

then wrap the existing body in if (table.writeOrdering().forall(o => isSpellable(o.expression()))) { ... }. DESCRIBE TABLE EXTENDED still reports both values verbatim, so nothing is hidden.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done. I added isSpellable() and wrapped the rest with it.

"Write for the binary file data source."
]
},
"WRITE_ORDERING_WITH_NESTED_COLUMN_IS_UNSUPPORTED" : {

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.

Finding 3. Both the name and the message describe cases this condition cannot report.

A nested struct column is supported, and the suite asserts it - CREATE TABLE testcat.t (p STRUCT<x: INT>) USING foo ORDERED BY p.x is accepted. So WITH_NESTED_COLUMN_IS_UNSUPPORTED says the opposite of the tested behaviour.

"or is in a map or array" is unreachable. StructType.findNestedField is called with the default includeCollections = false, which throws INVALID_FIELD_NAME on a path through a non-struct rather than returning None. Measured on this branch:

ORDERED BY (m.key) on MAP<STRING, INT>
  -> [INVALID_FIELD_NAME] Field name `m`.`key` is invalid: `m` is not a struct

ORDERED BY (a.element.x) on ARRAY<STRUCT<x: INT>>
  -> [INVALID_FIELD_NAME] Field name `a`.`element`.`x` is invalid: `a` is not a struct

ORDERED BY (truncate(4, m.key)) on MAP<STRING, INT>
  -> [INVALID_FIELD_NAME] Field name `m`.`key` is invalid: `m` is not a struct

The third one is raised from the CheckAnalysis check itself: truncate(...) is an ApplyTransform, so it is not rewritable and never reaches PreprocessTableCreation.

What this condition actually reports is a reference that is not a column of the table. I know it is a faithful copy of PARTITION_WITH_NESTED_COLUMN_IS_UNSUPPORTED, where both faults are equally present, and consistency with the sibling is a real argument - it just does not carry to a name being introduced now. The sibling can be left alone; this one cannot be renamed after a release. Suggest UNSUPPORTED_FEATURE.WRITE_ORDERING_WITH_UNKNOWN_COLUMN with "Invalid write ordering: <cols> is not a column of the table.", or a non-UNSUPPORTED_FEATURE parent, since a missing column is not a feature gap.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Changed it to WRITE_ORDERING_WITH_UNKNOWN_COLUMN.

DISTRIBUTED BY PARTITION UNORDERED;
```

Both clauses are passed to the data source, which has to support them: a data source that 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.

Finding 4. The advertiser is the catalog, not the data source. What gates this is TableCatalog.capabilities() returning TableCatalogCapability.SUPPORTS_CREATE_TABLE_WITH_WRITE_DISTRIBUTION_AND_ORDERING, checked in WriteDistributionAndOrdering.validateCatalogForWriteDistributionAndOrdering. A user who hits UNSUPPORTED_FEATURE.TABLE_OPERATION and follows this paragraph will inspect the USING provider, where there is nothing to inspect.

It also makes the last sentence hard to act on. "The built-in data sources do not support them" is true, but the reason is that V2SessionCatalog does not advertise the capability and the v1 conversion in ResolveSessionCatalog rejects the request outright - nothing to do with parquet or ORC.

Suggest saying catalog throughout: "Both clauses are passed to the catalog, which has to support them: a catalog that does not advertise support ... The built-in catalogs do not."

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done.


What the data source records is a *default* for later writes, not a statement about the data
already in the table: an individual write may override it, and rewriting existing data to match
a newly requested layout is a separate operation. `SHOW CREATE TABLE` reproduces the clauses and

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.

Finding 5. SHOW CREATE TABLE reproduces the clauses only for pairs the syntax can spell, and emits nothing at all for the rest. That caveat is in the PR description and in ShowCreateTableExec's scaladoc, but not here, and this page is what users read.

It is not a corner case. A connector that records a sort order without touching its distribution mode reports (null, non-empty), which has no clause form. Measured on this branch, with a table reporting mode null and ordering id DESC NULLS LAST:

CREATE TABLE n.t (
  id INT)
USING foo

while DESCRIBE TABLE EXTENDED on the same table shows Ordering = id DESC NULLS LAST. So that DDL runs and creates a table without the ordering.

One sentence covers it: SHOW CREATE TABLE reproduces the clauses when the recorded pair has a clause form, and DESCRIBE TABLE EXTENDED reports both values in every case.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done.

The distribution decides how far the order reaches, and this clause picks one when
`DISTRIBUTED BY PARTITION` is absent: a bare `ORDERED BY` range-partitions each write, so the
order holds across the whole table, while `LOCALLY ORDERED BY` asks for it to hold within each
written file only, without a shuffle. `UNORDERED` on its own asks for no distribution either.

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.

Finding 6. "within each written file only" here, "within each write task" three lines down at :136. A task can roll over several files, so these are different claims, and the second is the accurate one - it also matches TableInfo.DISTRIBUTION_MODE_NONE's javadoc ("any ordering holds within a write task only"). DISTRIBUTION_MODE_RANGE's javadoc has the same drift the other way ("across files, not only within one").

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done.

@anuragmantri
anuragmantri force-pushed the SPARK-34586-write-distribution-ordering-create-table branch from c6adea5 to 92cc539 Compare August 24, 2026 23:46

@peter-toth peter-toth 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.

Re-checked through 92cc5395424 - findings 1, 3, 5 and 6 resolved (Connect JDBC keyword list, the condition rename, the omission caveat, "write task"). Finding 2's fix narrows the hole rather than closing it, and finding 4 has one occurrence left; both measured on this head.

On CREATE TABLE ... LIKE: leave it out. CreateTableLikeExec hands sourceTable to TableCatalog.createTableLike alongside the TableInfo, so a connector already has what it needs to carry the layout across, and that exec's own doc names Iceberg sort order as the example. Adding the clauses to LIKE syntax is separate work. Finding 9 is only about saying so in the doc.

Blocking

  • 2. SHOW CREATE TABLE can still emit DDL that does not parse (round 1): isSpellable validates a Transform's arguments but never the transform itself. Expressions.apply("+", column("id"), literal(1)) is public API and yields ORDERED BY (+(id, 1) ASC NULLS FIRST), which fails replay with PARSE_SYNTAX_ERROR. [inline: sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/ShowCreateTableExec.scala:137]
  • 7. The new guard has no test (new): replacing the if at ShowCreateTableExec.scala:154 with if (true) leaves all 30 CreateTableWriteOrderSuite tests green, so nothing pins finding 2's fix - which is why its remaining hole went unnoticed. [inline: sql/core/src/test/scala/org/apache/spark/sql/connector/CreateTableWriteOrderSuite.scala:757]

Non-blocking

  • 4. One "data source" left where the gate is the catalog (round 1): :125 still says omitting the clause "leaves the choice to the data source". The other three are fixed; :117 is about CLUSTER BY and is right as it stands. Following up on the existing thread rather than opening a new one.
  • 8. Docs claim a range ordering "holds across the whole table" (late catch): it holds across one write's tasks, and a later INSERT INTO overlaps it. This now contradicts DISTRIBUTION_MODE_RANGE's javadoc, which this round narrowed to "across write tasks". [inline: docs/sql-ref-syntax-ddl-create-table-datasource.md:131]
  • 9. CreateTableLikeExec's doc enumerates what the TableInfo carries and now omits two fields (new): not copying them is right, but a connector author comparing against constraints will expect otherwise. [inline: sql/catalyst/src/main/java/org/apache/spark/sql/connector/catalog/TableInfo.java:159]

private def isSpellable(e: V2Expression): Boolean = e match {
case _: NamedReference => true
case t: Transform =>
t.arguments().forall(a => a.isInstanceOf[NamedReference] || a.isInstanceOf[Literal[_]])

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.

Finding 2. This closes the case on the parent thread (GeneralScalarExpression is neither a NamedReference nor a Transform), but it validates a Transform's arguments and never the transform itself. Two ways through:

  • The name. Expressions.apply(String name, Expression... args) is public and takes any name, and ApplyTransform.describe() renders it unquoted. applyTransform : transformName=identifier LEFT_PAREN ... needs an identifier there.
  • No arguments. forall on an empty arguments() is true, and the same rule requires at least one transformArgument.

Measured on 92cc5395424, with a fixture reporting mode range and Expressions.sort(Expressions.apply("+", Expressions.column("id"), Expressions.literal(1)), ASCENDING, NULLS_FIRST):

CREATE TABLE reportcat.t (
  id INT)
USING foo
ORDERED BY (+(id, 1) ASC NULLS FIRST)

Replaying that gives [PARSE_SYNTAX_ERROR] Syntax error at or near '+'. SQLSTATE: 42601 (line 4, pos 12) - the same whole-statement failure as before, through a narrower door.

Suggested change
t.arguments().forall(a => a.isInstanceOf[NamedReference] || a.isInstanceOf[Literal[_]])
t.name().matches("[a-zA-Z_][a-zA-Z0-9_]*") && t.arguments().nonEmpty &&
t.arguments().forall(a => a.isInstanceOf[NamedReference] || a.isInstanceOf[Literal[_]])

The regex approximates identifier rather than matching it: a name that is a reserved word (select) passes and still fails to parse under ANSI. That hole is far smaller than the current one and I would not chase it, but it is worth a word in the comment so the next reader knows this is an approximation rather than a claim.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Took your suggestion and also added a comment in the code about the approximation.

.getOrElse(tableInfo.writeDistributionMode())
}

override def writeOrdering(): Array[SortOrder] = tableInfo.writeOrdering()

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.

Finding 7. Nothing pins the guard finding 2 asked for. Measured in a review worktree on this head: replace the if at ShowCreateTableExec.scala:154 with if (true) and sql/testOnly *CreateTableWriteOrderSuite still reports Tests: succeeded 30, failed 0.

This line is why. Every ordering the suite can build arrives through tableInfo.writeOrdering(), so it comes from the parser, and every parser-produced sort key is a NamedReference or a Transform over references and literals - isSpellable is true in all 30 tests. MODE_OVERRIDE gave the distribution side a way to fabricate a value no statement could ask for; the ordering side has no equivalent, which is also why finding 2's remaining hole went unnoticed.

Same shape as MODE_OVERRIDE:

  override def writeOrdering(): Array[SortOrder] = {
    if (tableInfo.properties().containsKey(ReportingInMemoryTable.ORDERING_OVERRIDE)) {
      // A sort key `writeOrderField : transform ...` cannot represent.
      Array(Expressions.sort(
        Expressions.apply("+", Expressions.column("id"), Expressions.literal(1)),
        SortDirection.ASCENDING,
        NullOrdering.NULLS_FIRST))
    } else {
      tableInfo.writeOrdering()
    }
  }

Two assertions are worth having, one per trap named on the finding 2 thread:

  • mode range: no ORDERED BY in the output, and the emitted DDL still runs;
  • mode none: no UNORDERED either. That is the case where bailing out of the orderBy half alone would declare "no ordering" on a table that has one, and it is the half a test written only against range would miss.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done.

clause.

The distribution decides how far the order reaches, and this clause picks one when
`DISTRIBUTED BY PARTITION` is absent: a bare `ORDERED BY` range-partitions each write, so 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.

Finding 8. A range distribution orders one write's output across that write's tasks. It says nothing about the table over time: a later INSERT INTO range-partitions its own rows, so its files overlap the ranges already there and the table is not sorted end to end.

This round narrowed TableInfo.DISTRIBUTION_MODE_RANGE's javadoc from "across files, not only within one" to "across write tasks, not only within one", so the two now disagree and the javadoc is the accurate one. The same claim is in the SQL comment at :140.

The paragraph is about what the clause requests for every write ("recorded on the table so that later writes honor it too"), which is exactly the scope where the stronger claim fails. "so the order holds across the tasks of a write, not only within one" matches the javadoc and is what the mechanism delivers.

I raised the file-vs-task drift on this bullet last round and missed this one, which is the bigger of the two: it promises a table-level layout guarantee.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

I changed the docs, let me know if that is what you meant.

*
* @since 4.4.0
*/
public Builder withWriteOrdering(SortOrder[] writeOrdering) {

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.

Finding 9. TableInfo now carries two more fields that CREATE TABLE ... LIKE could copy from its source, and CreateTableLikeExec does not. That is the right call - it passes sourceTable to TableCatalog.createTableLike, so the connector has everything, and its own doc already names Iceberg sort order as the example - but the exec's class doc enumerates what the TableInfo carries:

columns and partitioning copied from the source, constraints copied from
the source, user-specified TBLPROPERTIES / LOCATION / USING provider ...

so the omission now reads as an oversight (sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/CreateTableLikeExec.scala:43-48). A connector author comparing against constraints, which is copied, will expect the same treatment. One sentence there - the declared write distribution and ordering are deliberately not copied, read them from sourceTable - saves the next reader deriving it from this PR.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Added it in the docs

@anuragmantri
anuragmantri force-pushed the SPARK-34586-write-distribution-ordering-create-table branch from 92cc539 to 19365b5 Compare September 11, 2026 05:34

@anuragmantri anuragmantri left a comment

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Sorry for my late response. I have updated the PR with your comments. Ready for another round. @peter-toth.

private def isSpellable(e: V2Expression): Boolean = e match {
case _: NamedReference => true
case t: Transform =>
t.arguments().forall(a => a.isInstanceOf[NamedReference] || a.isInstanceOf[Literal[_]])

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Took your suggestion and also added a comment in the code about the approximation.

*
* @since 4.4.0
*/
public Builder withWriteOrdering(SortOrder[] writeOrdering) {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Added it in the docs

.getOrElse(tableInfo.writeDistributionMode())
}

override def writeOrdering(): Array[SortOrder] = tableInfo.writeOrdering()

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done.

clause.

The distribution decides how far the order reaches, and this clause picks one when
`DISTRIBUTED BY PARTITION` is absent: a bare `ORDERED BY` range-partitions each write, so the

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

I changed the docs, let me know if that is what you meant.

@peter-toth peter-toth 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.

Re-checked through 19365b5946a — findings 2, 4, 7, 8 and 9 resolved, nothing new.

Finding 7 re-measured rather than read: forcing the gate at ShowCreateTableExec.scala:157 to true now fails SHOW CREATE TABLE omits a pair the syntax cannot spell, and stays runnable on ORDERED BY (+(id, 1) ASC NULLS FIRST), where the same ablation left 30/30 green last round. That also pins finding 2's fix.

Thanks for working through all of these, @anuragmantri — nothing left open from my side.

@peter-toth

Copy link
Copy Markdown
Contributor

@aokolnychyi, @szehon-ho, can you please take a look when you have some time?

@peter-toth

Copy link
Copy Markdown
Contributor

Any comments or suggestions @aokolnychyi, @szehon-ho? I would like to merge this PR this week if there isn't any.

@szehon-ho szehon-ho 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.

How should the new write-layout clauses interact with existing table CLUSTER BY support? For example, CLUSTER BY (a) ORDERED BY (b) and CLUSTER BY (a) UNORDERED are accepted and send both declarations to the connector. Is the connector responsible for reconciling them or rejecting incompatible combinations? Could we document that contract and add coverage for these combinations?

Anton (@aokolnychyi) may have more thoughts on this than me, given his work on the DSv2 write distribution and ordering APIs.

* would read back as a transform *named* `identity` rather than as a plain column reference.
*/
def describeSortOrder(sortOrder: SortOrder): String = {
s"${sortOrder.expression().describe()} ${sortOrder.direction()} ${sortOrder.nullOrdering()}"

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] describe() does not preserve literal types. For example, the parser stores DATE '1970-01-01' as LiteralValue(0, DateType), whose description is 0. Consequently, an ordering such as f(id, DATE '1970-01-01') passes isSpellable but becomes f(id, 0) in SHOW CREATE TABLE. Replaying that DDL changes the argument to IntegerType, potentially changing or invalidating the ordering. Could we render type-preserving SQL literals and add a round-trip test comparing the reconstructed ordering expressions?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Good catch, thanks. describeSortOrder now renders literal arguments through Catalyst's Literal.sql, so their types survive, e.g. DATE '1970-01-01' instead of 0. I added a test that replays the SHOW CREATE TABLE output, and checks that the reconstructed ordering equals the original, literal data types included.

@dongjoon-hyun

Copy link
Copy Markdown
Member

Thank you for working on this, @anuragmantri. I have a few additional comments.

  1. Error class of WRITE_ORDERING_WITH_UNKNOWN_COLUMN: It is placed under UNSUPPORTED_FEATURE (SQLSTATE 0A000), but the error is about a reference to a non-existent column rather than an unsupported feature. A 42703-family condition looks more appropriate. It seems to follow PARTITION_WITH_NESTED_COLUMN_IS_UNSUPPORTED, but that condition was originally about nested columns, so its semantics differ. Since error classes are hard to change after release, could we decide this before merging?

  2. Inconsistent case handling: PreprocessTableCreation normalizes references only for RewritableTransform, so an ApplyTransform such as truncate(4, ID) reaches the case-sensitive findNestedField check in CheckAnalysis as-is. As a result, even under the default case-insensitive analysis, ORDERED BY truncate(4, ID) fails against a column id while ORDERED BY ID succeeds. PARTITIONED BY has the same pre-existing limitation, so this is not a regression, but it would be good to document it or pin it with a test.

  3. Source compatibility for downstream projects: New fields are added to the CreateTable, ReplaceTable, CreateTableAsSelect, and ReplaceTableAsSelect case classes and several Exec case classes. External extensions that pattern-match on these nodes (e.g. Delta, Iceberg Spark extensions) will fail to compile. This is acceptable because they are internal APIs, but could you mention it in the PR description?

  4. Code comments: Some comments describe the review history, e.g. "Reusing checkTransformDuplication here was wrong" and the 13-line normalization explanation in rules.scala. Some are also much more verbose than the surrounding code, e.g. "This is the only thing standing between the user and that silent drop" in WriteDistributionAndOrdering. Could you trim them to briefly describe only the current behavior?

@dongjoon-hyun

Copy link
Copy Markdown
Member

Gentle ping, @anuragmantri .

@dongjoon-hyun

Copy link
Copy Markdown
Member

Gentle ping, @anuragmantri . Please resolve the conflicts and check my review comments.

@anuragmantri
anuragmantri force-pushed the SPARK-34586-write-distribution-ordering-create-table branch from 19365b5 to d2aa867 Compare September 28, 2026 17:15

@anuragmantri anuragmantri left a comment

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Thanks for the reviews @szehon-ho and @dongjoon-hyun.

I rebased onto the latest master to resolve the conflicts and addressed your points.

@dongjoon-hyun for your comments:

  1. WRITE_ORDERING_WITH_UNKNOWN_COLUMN is now a top-level condition with SQLSTATE 42703 instead of a subclass of UNSUPPORTED_FEATURE.
  2. Yes, I have verified this finding. For this PR, I added a test for case-insensitive session does not normalize an ApplyTransform's references to CreateTableWriteOrderSuite. It shows ORDERED BY ID is normalized against a column id while ORDERED BY truncate(4, ID) is rejected, and that PARTITIONED BY behaves the same way.
  3. Added a note to the PR description that extensions pattern-matching on the create/replace plans and exec nodes need to add the two new fields.
  4. Trimmed the comments to describe only the current behavior, including the ones you pointed out in rules.scala and WriteDistributionAndOrdering.

@szehon-ho

How should the new write-layout clauses interact with existing table CLUSTER BY support? For example, CLUSTER BY (a) ORDERED BY (b) and CLUSTER BY (a) UNORDERED are accepted and send both declarations to the connector. Is the connector responsible for reconciling them or rejecting incompatible combinations? Could we document that contract and add coverage for these combinations?

Spark treats them as independent declarations. CLUSTER BY records clustering columns for the catalog to interpret, and the write clauses record a declared distribution and ordering. Spark passes both to the catalog without reconciling them, the same way it passes the write layout through without enforcing it. So yes, the connector is responsible for interpreting the combination, or rejecting it if it can't support it. I documented this contract in sql-ref-syntax-ddl-create-table-datasource.md and added a test showing that CLUSTER BY (a) ORDERED BY (b) and CLUSTER BY (a) UNORDERED both reach the catalog with both declarations intact.

I'm happy to adjust if you think Spark should validate any of these combinations.

* would read back as a transform *named* `identity` rather than as a plain column reference.
*/
def describeSortOrder(sortOrder: SortOrder): String = {
s"${sortOrder.expression().describe()} ${sortOrder.direction()} ${sortOrder.nullOrdering()}"

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Good catch, thanks. describeSortOrder now renders literal arguments through Catalyst's Literal.sql, so their types survive, e.g. DATE '1970-01-01' instead of 0. I added a test that replays the SHOW CREATE TABLE output, and checks that the reconstructed ordering equals the original, literal data types included.

@dongjoon-hyun

dongjoon-hyun commented Sep 29, 2026 •

Copy link
Copy Markdown
Member

Thank you for addressing the previous comments, @anuragmantri. I confirmed that my four points are resolved in d2aa8674fdc. I have a few more comments, mostly on the last round of changes.

  1. SHOW CREATE TABLE can still emit DDL that does not parse for a FLOAT literal. toSQL now renders literals through Catalyst Literal.sql (WriteDistributionAndOrdering.scala:65), but Literal.sql renders a FloatType value as CAST('1.5' AS FLOAT), and transformArgument only accepts qualifiedName | constant. So ORDERED BY f(id, 1.5F) is accepted (isSpellable only checks isInstanceOf[Literal[_]] at ShowCreateTableExec.scala:140), but SHOW CREATE TABLE prints ORDERED BY (f(id, CAST('1.5' AS FLOAT)) ASC NULLS FIRST), which cannot be replayed. The same applies to non-finite DOUBLE values reported by a connector (CAST('NaN' AS DOUBLE)). Could you render FloatType as <v>F (FLOAT_LITERAL), reject non-finite float/double literals in isSpellable, and add 1.5F to the round-trip test at CreateTableWriteOrderSuite.scala:525?

  2. DESCRIBE TABLE EXTENDED / SHOW CREATE TABLE can throw for a connector-reported literal. CatalystLiteral(l.value, l.dataType) requires the value to be in the internal representation already, so a LiteralValue holding a java.lang.String (which Expressions.literal("x") produces) fails with IllegalArgumentException: requirement failed: Literal must have a corresponding value to string .... The connector would be violating the Literal Javadoc, but DESCRIBE is the place where we promise to report the declared layout as-is. CatalystLiteral.create(l.value, l.dataType) accepts both representations.

  3. The CLUSTER BY contract is documented only in the SQL reference. The answer to Szehon's question (Spark passes both declarations without reconciling them, and the catalog interprets or rejects the combination) is meant for connector authors, so could we also state it in the Javadoc of TableCatalogCapability.SUPPORTS_CREATE_TABLE_WITH_WRITE_DISTRIBUTION_AND_ORDERING (TableCatalogCapability.java:98) or of TableInfo#writeDistributionMode()?

  4. Question on DelegatingTable. DelegatingTable now reports the declared layout (DelegatingTable.java:89-96), but Spark reads and writes a DelegatingTable through the v1 path (RelationResolution.scala:467), which never looks at it. If such a catalog advertises the capability, DESCRIBE / SHOW CREATE TABLE show an ORDERED BY that Spark's own writes never apply. Is this intended? If so, could we mention it in the capability Javadoc?

  5. Question on the public API shape. Table#writeDistributionMode() and the new TableInfo methods become public API in 4.4.0 and are hard to change afterwards. Did you consider a small Java enum (like SortDirection / NullOrdering) instead of String plus the TableInfo.DISTRIBUTION_MODE_* constants? It would make the closed set explicit, and SHOW CREATE TABLE would no longer need the unknown-mode (zigzag) case. If we keep String, please add @since 4.4.0 to the three constants (TableInfo.java:40,46,52).

  6. Test suggestion. The main motivation is that the first load of a CTAS/RTAS can already follow the declared layout. Since CreateTableAsSelectExec writes through AppendData on the newly created table, a single test would pin this: a test catalog that builds InMemoryTable(..., distribution, ordering) from the TableInfo, plus withQueryExecutionsCaptured, to check that the inner write plans the required shuffle/sort.

Minor:

  • QueryParsingErrors.distributedByPartitionWithoutPartitioning -> ...Error, like the other methods there.
  • SparkSqlParser.scala:1695 adds a new use of _LEGACY_ERROR_TEMP_0035. It matches the surrounding checks, but invalidStatement (INVALID_STATEMENT_OR_CLAUSE), as used for CREATE TEMPORARY TABLE at :614, would avoid a new legacy-condition usage.
  • The operation string ... DISTRIBUTED BY/ORDERED BY (WriteDistributionAndOrdering.scala:53) is shown even when the statement only has UNORDERED.
  • docs/sql-ref-syntax-ddl-create-table-datasource.md:154-156: Spark rejects the statement, not the catalog.
  • PR description: the suite now has 33 tests, and V2CreateTablePlan gained two abstract members (writeOrdering, withWriteOrdering) that implementors in extensions also need to add.

The current CI failures look unrelated: OracleIntegrationSuite (ORA-12516, cannot connect to the database) and the K8s DepsTestsSuite "SPARK-33748: Launcher python client respecting PYSPARK_PYTHON" timeout. Could you re-trigger CI?

@szehon-ho szehon-ho 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.

Thanks for clarifying the CLUSTER BY interaction. I'm fine keeping CLUSTER BY with ORDERED BY, LOCALLY ORDERED BY, or UNORDERED legal and leaving compatibility to the catalog. The existing rejection of CLUSTER BY with DISTRIBUTED BY PARTITION makes sense. The inline comments cover the rejection-path tests, a connector-transform case in the partitioning guard, and a wording correction.

None
}
// Bucketing counts as partitioning here; CLUSTER BY does not.
val hasPartitioning = table.partitioning.exists(!_.isInstanceOf[ClusterByTransform])

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] Could we use the ClusterByTransform(_) extractor here instead of an implementation-class check? A connector can return Expressions.apply("cluster_by", Expressions.column("a")) (or its own Transform implementation). The existing extractor recognizes that as clustering, but isInstanceOf[ClusterByTransform] is false, so this guard treats it as actual partitioning. If the table reports declared mode hash, SHOW CREATE TABLE then emits DISTRIBUTED BY PARTITION for a clustering-only table, although this method intends to omit that pair.

A focused regression test could stub Table.partitioning() with this generic transform and report mode hash, with both empty and non-empty write ordering. Assert that neither DISTRIBUTED BY PARTITION nor a standalone ORDERED BY is emitted for that unrepresentable pair. Keep an actual partition or bucket transform as a positive control. TransformExtractorSuite already covers recognition of a connector-defined cluster_by transform.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Good catch. The guard now uses the ClusterByTransform(_) extractor, so a connector's generic cluster_by transform counts as clustering, not partitioning. I added "SHOW CREATE TABLE treats a connector's cluster_by transform as clustering". It uses a test catalog whose tables report CLUSTER BY as Expressions.apply("cluster_by", ...). It returns a DelegatingTable, because InMemoryTable rejects a generic cluster_by at creation. With mode hash and both an empty ordering and ORDERED BY (b), the test asserts that neither DISTRIBUTED BY PARTITION nor ORDERED BY is emitted, for both the native and the generic transform. PARTITIONED BY (a) and CLUSTERED BY (a) INTO 4 BUCKETS are the positive controls and still emit DISTRIBUTED BY PARTITION ORDERED BY (b ASC NULLS FIRST). I checked that the generic case fails with the old isInstanceOf check.

}
}

test("CLUSTER BY and the write clauses both reach the catalog") {

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.

Could we complement this successful pass-through test with a catalog that advertises SUPPORTS_CREATE_TABLE_WITH_WRITE_DISTRIBUTION_AND_ORDERING but deliberately rejects one combination, for example CLUSTER BY (a) UNORDERED? This would exercise the connector-rejection contract separately from the existing tests where the catalog lacks the entire capability.

For CREATE, assert that the catalog's error reaches the caller and no table is published. For staged REPLACE/RTAS, start with a populated table and assert that its rows, schema, and clustering metadata survive the rejection. The unchanged-table assertion should be scoped to staging: non-staging ReplaceTableExec drops the original before calling createTable, so late catalog rejection there retains the existing non-atomic replacement limitation.

It would also be useful to include CLUSTER BY (a) LOCALLY ORDERED BY (b) in this positive test, asserting mode none with the ordering retained.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Added "a catalog with the capability can still reject a combination it does not support". It uses a staging catalog that advertises the capability but rejects CLUSTER BY with UNORDERED in createTable and in all three stage* methods. For CREATE TABLE (the createTable path) and CTAS (the stageCreate path), the catalog's error reaches the caller and no table is published. The test then populates a CLUSTER BY (a) table and runs REPLACE TABLE, CREATE OR REPLACE TABLE, RTAS and CREATE OR REPLACE ... AS SELECT, each with CLUSTER BY (a) UNORDERED. After each rejection, the rows, the schema and the clustering are unchanged. As you suggested, it covers only the staging path, because the non-staging ReplaceTableExec drops the original first. I also added CLUSTER BY (a) LOCALLY ORDERED BY (b) to this test, which asserts mode NONE with the ordering b ASC NULLS FIRST kept and the clustering intact.

Both clauses are passed to the catalog, which has to support them: a catalog that does not
advertise support for a write distribution and ordering rejects the statement rather than
creating a table that silently lacks the requested layout. The built-in catalogs do not
support them. Both may also be combined with `CLUSTER BY`. Spark passes the clustering columns

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.

Both may also be combined with CLUSTER BY appears to include DISTRIBUTED BY PARTITION, but that combination is rejected by the parser and by the test at CreateTableWriteOrderSuite:311. Could we name the allowed forms explicitly here: ORDERED BY, LOCALLY ORDERED BY, and UNORDERED may be combined with CLUSTER BY, subject to the catalog accepting the combination?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Thanks, that sentence was too broad. It now reads: "ORDERED BY, LOCALLY ORDERED BY, and UNORDERED may also be combined with CLUSTER BY, subject to the catalog accepting the combination ... DISTRIBUTED BY PARTITION cannot be combined with CLUSTER BY, as described above."

@dongjoon-hyun

Copy link
Copy Markdown
Member

Thank you for the update, @anuragmantri. I took one more pass over d2aa8674fdc and have some additional comments. I continue the numbering from my previous comment so that the points can be referred to unambiguously. None of these is a regression of existing behavior, and they come from reading the code rather than from running it.

  1. SHOW CREATE TABLE silently drops every write clause when a transform name needs quoting. transformName=identifier accepts a quoted name and the parser strips the backticks (parsers.scala:310), so CREATE TABLE cat.t (id INT, p INT) USING foo PARTITIONED BY (p) DISTRIBUTED BY PARTITION ORDERED BY `z-order`(id) is a valid statement. isSpellable rejects the name with its regex (ShowCreateTableExec.scala:139), and because the forall(isSpellable) guard at :151 wraps the whole block, SHOW CREATE TABLE prints neither ORDERED BY nor DISTRIBUTED BY PARTITION. Replaying the output creates a table without the declared layout, with no error. The earlier discussion treated a name like + as having no spelling, but `+`(id, 1) does parse as an applyTransform. Could we render the name with quoteIfNeeded(t.name) in toSQL (WriteDistributionAndOrdering.scala:67) and drop the name check? The regex is the same as QuotingUtils.validIdentPattern, and column references are already quoted this way. The +(id, 1) fixture would then round-trip, so the "unspellable" test needs a key that really has no spelling, e.g. a transform whose argument is another transform.

  2. DISTRIBUTED BY PARTITION on a schemaless CREATE/REPLACE TABLE gives advice that cannot be followed. The parser decides "no partitioning" from the absence of PARTITIONED BY / CLUSTERED BY in the statement (AstBuilder.scala:6518, :6619). For CREATE TABLE cat.t USING foo LOCATION '/x' DISTRIBUTED BY PARTITION, the error says "Please add PARTITIONED BY or CLUSTERED BY ... INTO ... BUCKETS, or drop the clause and use only ORDERED BY." However, adding PARTITIONED BY (p) then fails in PreprocessTableCreation with _LEGACY_ERROR_TEMP_1165 ("It is not allowed to specify partitioning when the table schema is not defined", rules.scala:330-332), and using only ORDERED BY fails with SPECIFY_WRITE_ORDERING_IS_NOT_ALLOWED. So the hash mode cannot be declared at all when the schema and partitioning come from the catalog, while UNORDERED is accepted there. The docs mention the schemaless restriction for ORDERED BY (sql-ref-syntax-ddl-create-table-datasource.md:127) but not for DISTRIBUTED BY PARTITION (:114). Could we at least fix the message and add the sentence to the docs? Alternatively, the check could move next to the schemaless check in PreprocessTableCreation and be skipped when the schema is empty, leaving that case to the catalog.

  3. Built-in transform names are matched case-sensitively, so the pinned limitation is wider than the test shows. visitTransform dispatches on the raw text (applyCtx.identifier.getText match { case "bucket" ... case "days" ... }, AstBuilder.scala:5869), so DAYS(ts) or BUCKET(4, id) becomes an ApplyTransform rather than a DaysTransform / BucketTransform. In the default case-insensitive session, ORDERED BY days(TS) and ORDERED BY bucket(4, ID) succeed, but ORDERED BY DAYS(TS) and ORDERED BY BUCKET(4, ID) fail with WRITE_ORDERING_WITH_UNKNOWN_COLUMN, and BUCKET(4, id) hands the catalog a generic transform named BUCKET. The dispatch is pre-existing and PARTITIONED BY behaves the same, so this is not a regression, but it is another way to hit the limitation pinned by "a case-insensitive session does not normalize an ApplyTransform's references". Could you add it to that test? As a possible follow-up, RewritableTransform is a sealed trait in the same file and its only user is PreprocessTableCreation, so making ApplyTransform a RewritableTransform would let both ORDERED BY and PARTITIONED BY normalize these references and remove the limitation.

  4. The capability Javadoc names only createTable, but the request can be dropped per method. TableCatalog.createTable(Identifier, TableInfo) and StagingTableCatalog.stageCreate / stageReplace / stageCreateOrReplace(Identifier, TableInfo) all forward only columns, partitions and properties by default. A staging catalog that advertises the capability and overrides the three stage* overloads still loses the request silently for CREATE TABLE cat.t (id INT) USING x ORDERED BY id, because a CREATE TABLE with a column list is not staged and goes to createTable (the fixture comment at CreateTableWriteOrderSuite.scala:706 describes this trap). Could the Javadoc of SUPPORTS_CREATE_TABLE_WITH_WRITE_DISTRIBUTION_AND_ORDERING and of TableInfo#writeDistributionMode() list the four TableInfo overloads and say that a catalog advertising the capability must override every one it can be reached through? This pairs with point 3.

  5. PARTITIONED BY and ORDERED BY now render the same transform differently in one output. The literal fix went into a renderer private to the ordering (WriteDistributionAndOrdering.toSQL), while showTablePartitioning still uses t.describe() (ShowCreateTableExec.scala:112), and so does the Part N row in DescribeTableExec.scala:161. For PARTITIONED BY (truncate(s, 10L)) ORDERED BY truncate(s, 10L), SHOW CREATE TABLE prints PARTITIONED BY (truncate(s, 10)) next to ORDERED BY (truncate(s, 10L) ASC NULLS FIRST), and replaying it turns the partition literal into an INT. The partition half is pre-existing, so I'm fine with a follow-up, but could the comment at WriteDistributionAndOrdering.scala:62 say that PARTITIONED BY still goes through describe? Sharing one renderer would be the real fix once points 1 and 2 are settled.

  6. Question on the shape of a sort key in the new API. ORDERED BY id reaches the catalog as sort(identity(id), ...) because visitWriteOrderField wraps whatever visitTransform returns (AstBuilder.scala:6368). Other producers of v2 SortOrder in Spark use a bare NamedReference for a column (e.g. DataSourceStrategy.translateSortOrders), which is also the usual form for requiredOrdering(), and V2ExpressionSQLBuilder does not accept a Transform inside a SortOrder. Spark's own write path resolves both forms, so nothing is broken. However, TableInfo#writeOrdering() and Table#writeOrdering() become public API in 4.4.0, so could we either unwrap the identity transform or document the shape in the Javadoc (each key is a Transform, and a plain column comes as identity)? It would also be good to state there that the array is never null, and whether checking that a key's type is orderable is left to the catalog (e.g. ORDERED BY m on a MAP column is accepted today).

  7. Question on the plan and exec signatures. This is a follow-up to my earlier point on source compatibility. I said the break was acceptable, but there is a precedent going the other way: SPARK-43529 added a default parameter to the same four plans, and its follow-up (7e94f2a5433) moved it into UnresolvedTableSpec to "Restore the signatures of class CreateTable, CreateTableAsSelect, ReplaceTable and ReplaceTableAsSelect". Constraints and the default collation also travel on TableSpec, which already reaches every consumer here, so carrying the pair there would leave the 22 positional patterns (ApplyDefaultCollation, ResolveCatalogs and TypeCoercionBase change only for arity) and the 7 exec signatures untouched. For example, Apache Paimon's PaimonCreateTableAsSelectStrategy matches CreateTableAsSelect and constructs CreateTableAsSelectExec positionally, so it stops compiling with this PR. Did you consider this? I'm fine with a follow-up if it is too much for this PR.

  8. Tests.

    • CreateTableWriteOrderSuite.scala:550 asserts with ddl.contains(expected), and ORDERED BY (id DESC NULLS LAST) is also a substring of LOCALLY ORDERED BY (id DESC NULLS LAST), while the replay test at :532 compares only ordering(). So emitting LOCALLY ORDERED BY for a range table would not fail any test. Each clause is printed on its own line, so ddl.split("\n").contains(expected) would pin the mode. Comparing writeDistributionMode() in the replay test and adding a PARTITIONED BY (c) ORDERED BY (id) case would help too.
    • "the new keywords stay usable as identifiers" toggles only spark.sql.ansi.enabled, but the parser switches to ansiNonReserved only when spark.sql.ansi.enforceReservedKeywords is also true (SQLConf.scala:9779), so both iterations take the nonReserved path. Could you set ENFORCE_RESERVED_KEYWORDS in the ANSI iteration?

Minor:

  • rules.scala:362-365: findNestedField already returns the normalized path, so the body of normalizeResolvableReferences can be schema.findNestedField(fieldNames, resolver = resolver).map { case (path, field) => FieldReference(path :+ field.name) }.getOrElse(ref) instead of three schema traversals through two lookup functions with different rules (DescribeTableExec.scala:117-119 uses the same idiom). normalizeReferences at :343 has a single caller and could stay inline as before.
  • ShowCreateTableExec.scala:131-134: the scaladoc describes the transform rule, not transformArgument. Also, select fails to parse only when spark.sql.ansi.enforceReservedKeywords is true, not under ANSI in general.
  • CreateTableWriteOrderSuite.scala:706: a REPLACE TABLE with a column list is staged (the test at :381 asserts stageReplace); only CREATE TABLE is not.
  • Related to point 1: a connector-reported typed NULL literal renders as CAST(NULL AS INT) (literals.scala:654), so an allow-list of literal types in isSpellable may be simpler than special-casing FLOAT and non-finite values.
  • The hash-without-partitioning check is duplicated verbatim in visitCreateTable and visitReplaceTable. It could live in writeSpecsFrom.

Pre-existing in the shared visitTransform and now reachable from ORDERED BY (fine as follow-ups):

  • bucket(4294967300, id) is silently truncated to bucket(4, id) by .toInt (AstBuilder.scala:5877), while a DECIMAL or TINYINT count is rejected.
  • ORDERED BY IDENTIFIER('id') is not resolved to id because visitQualifiedName uses _.getText (AstBuilder.scala:5806) instead of getIdentifierParts. PARTITIONED BY has the same gap.

@dongjoon-hyun

Copy link
Copy Markdown
Member

Gentle ping, @anuragmantri .

@anuragmantri anuragmantri left a comment

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

@szehon-ho thanks for confirming the CLUSTER BY behavior. It stays combinable with ORDERED BY, LOCALLY ORDERED BY and UNORDERED, and it is still rejected with DISTRIBUTED BY PARTITION. Your three inline points are addressed in 5569461.

Thank you for the two detailed passes, @dongjoon-hyun. Keeping your numbering:

  1. Done. A finite FLOAT now renders as <v>F, isSpellable rejects non-finite FLOAT/DOUBLE, and 1.5F is in the round-trip test. Two related cases came up: Float.MaxValue renders as 3.4028235E38F, which is outside the parser's FLOAT range on replay, so it is treated as unspellable, and DESCRIBE now prints a non-finite FLOAT as CAST('NaN' AS FLOAT) instead of NaNF.

  2. Good catch. Switched to CatalystLiteral.create. A connector-reported java.lang.String literal is now tested in both DESCRIBE TABLE EXTENDED and SHOW CREATE TABLE.

  3. Done. The Javadoc of SUPPORTS_CREATE_TABLE_WITH_WRITE_DISTRIBUTION_AND_ORDERING now says that Spark passes the clustering columns and the requested distribution and ordering without reconciling them, and the catalog interprets or rejects the combination.

  4. Yes, it is intended. A DelegatingTable exposes what the catalog stored, so reporting the layout keeps DESCRIBE and SHOW CREATE TABLE faithful, and another engine reading the same catalog can apply it. Spark's own reads and writes go through the v1 path, which does not apply it, and the capability Javadoc now says so.

  5. Agreed. It is now an enum, WriteDistributionMode (HASH, RANGE, NONE), @since 4.4.0, and the String constants are gone. SHOW CREATE TABLE no longer needs the unknown-mode case. DESCRIBE still prints hash / range / none.

  6. Added "CTAS and RTAS write their first load with the declared distribution and ordering". A test catalog builds InMemoryTable with the distribution and ordering from the TableInfo, and the test checks the inner write: a range shuffle plus sort for ORDERED BY, a sort with no shuffle for LOCALLY ORDERED BY, and a hash shuffle plus sort for CREATE OR REPLACE ... DISTRIBUTED BY PARTITION ... AS SELECT.

  7. Good catch. toSQL now quotes the name with quoteIfNeeded, and the name check in isSpellable is gone. A new test round-trips ORDERED BY `z-order`(id), `+`(id, 1), and the unspellable fixture is now a transform with a transform argument, f(g(id)).

  8. I kept the parse-time check and fixed the message and the docs. Both now say that a statement that defines no schema (no column list, no typed partition columns, and no AS SELECT) cannot declare partitioning, so it cannot use DISTRIBUTED BY PARTITION. I rephrased the ORDERED BY docs sentence the same way, since PARTITIONED BY (p INT) defines a schema without a column list. I preferred this to moving the check because it keeps a simple guarantee for catalogs, that HASH always comes with non-empty partitions(). Relaxing that later stays compatible, but tightening it after 4.4.0 would not.

  9. Added DAYS(TS) and BUCKET(4, ID) to that test, next to the lowercase forms, which do normalize. Filed SPARK-59944 for making ApplyTransform a RewritableTransform. That alone would not stop BUCKET(4, id) from reaching the catalog as a generic transform. That needs case-insensitive dispatch in visitTransform, which also changes PARTITIONED BY, so the JIRA covers both.

  10. Done. The capability Javadoc and TableInfo#writeDistributionMode() now list TableCatalog#createTable(Identifier, TableInfo) and the three StagingTableCatalog stage*(Identifier, TableInfo) overloads, and say that a catalog reporting the capability must override each one it can be reached through. The capability Javadoc also notes that a staging catalog needs all four, because a CREATE TABLE without AS SELECT is not staged.

  11. The comment on toSQL now says that PARTITIONED BY and the Part N rows of DESCRIBE render through describe. Sharing one renderer is SPARK-59946.

  12. I went with unwrapping. ORDERED BY id now reaches the catalog as a bare NamedReference, like the other v2 SortOrder producers, and PreprocessTableCreation normalizes a bare reference the same way it normalizes a transform's references. The TableInfo#writeOrdering() Javadoc now documents the shape (a column is a NamedReference, anything else a Transform), that the array is never null, and that Spark only checks that the referenced columns exist, so orderability and argument types are left to the catalog.

  13. Thanks for the SPARK-43529 precedent. I'd like to do this as a follow-up, SPARK-59943, and land it before 4.4.0. One thing to handle there: once the pair travels on TableSpec, a custom strategy like Paimon's compiles again but skips validateCatalogForWriteDistributionAndOrdering, which runs in DataSourceV2Strategy. So the follow-up should move that check into the execs or into CheckAnalysis.

  14. Done. The SHOW CREATE TABLE assertions now use ddl.split("\n").contains(expected), the replay test compares writeDistributionMode() as well as the ordering, and there is a PARTITIONED BY (c) ORDERED BY (id) case. The keyword test now sets ENFORCE_RESERVED_KEYWORDS in its ANSI iteration.

None
}
// Bucketing counts as partitioning here; CLUSTER BY does not.
val hasPartitioning = table.partitioning.exists(!_.isInstanceOf[ClusterByTransform])

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Good catch. The guard now uses the ClusterByTransform(_) extractor, so a connector's generic cluster_by transform counts as clustering, not partitioning. I added "SHOW CREATE TABLE treats a connector's cluster_by transform as clustering". It uses a test catalog whose tables report CLUSTER BY as Expressions.apply("cluster_by", ...). It returns a DelegatingTable, because InMemoryTable rejects a generic cluster_by at creation. With mode hash and both an empty ordering and ORDERED BY (b), the test asserts that neither DISTRIBUTED BY PARTITION nor ORDERED BY is emitted, for both the native and the generic transform. PARTITIONED BY (a) and CLUSTERED BY (a) INTO 4 BUCKETS are the positive controls and still emit DISTRIBUTED BY PARTITION ORDERED BY (b ASC NULLS FIRST). I checked that the generic case fails with the old isInstanceOf check.

}
}

test("CLUSTER BY and the write clauses both reach the catalog") {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Added "a catalog with the capability can still reject a combination it does not support". It uses a staging catalog that advertises the capability but rejects CLUSTER BY with UNORDERED in createTable and in all three stage* methods. For CREATE TABLE (the createTable path) and CTAS (the stageCreate path), the catalog's error reaches the caller and no table is published. The test then populates a CLUSTER BY (a) table and runs REPLACE TABLE, CREATE OR REPLACE TABLE, RTAS and CREATE OR REPLACE ... AS SELECT, each with CLUSTER BY (a) UNORDERED. After each rejection, the rows, the schema and the clustering are unchanged. As you suggested, it covers only the staging path, because the non-staging ReplaceTableExec drops the original first. I also added CLUSTER BY (a) LOCALLY ORDERED BY (b) to this test, which asserts mode NONE with the ordering b ASC NULLS FIRST kept and the clustering intact.

Both clauses are passed to the catalog, which has to support them: a catalog that does not
advertise support for a write distribution and ordering rejects the statement rather than
creating a table that silently lacks the requested layout. The built-in catalogs do not
support them. Both may also be combined with `CLUSTER BY`. Spark passes the clustering columns

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Thanks, that sentence was too broad. It now reads: "ORDERED BY, LOCALLY ORDERED BY, and UNORDERED may also be combined with CLUSTER BY, subject to the catalog accepting the combination ... DISTRIBUTED BY PARTITION cannot be combined with CLUSTER BY, as described above."

@dongjoon-hyun

Copy link
Copy Markdown
Member

Thank you for updating, @anuragmantri .

@dongjoon-hyun dongjoon-hyun 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.

Thank you for addressing points 1-14, @anuragmantri. I reviewed the latest commit, 5569461, and left 15 inline comments, continuing the numbering from points 1-14. Three of them (19, 21 and 29) cover cases that the fixes for points 8, 2 and 12 do not reach yet. Here is a summary, ordered by priority.

Correctness

  1. ShowCreateTableExec L156: TimeType is missing from the allow-list, so a TIME literal in a sort key drops every write clause from SHOW CREATE TABLE. This is a regression in this commit and the silent drop from point 7.
  2. ShowCreateTableExec L154: true, false and NULL re-parse as column references inside a transform, so BooleanType and NullType do not round-trip.
  3. ShowCreateTableExec L155: TIMESTAMP literals are printed as session-local wall time without an offset, so a replay can shift the value (DST overlap, another time zone) or turn it into TIMESTAMP_NTZ.
  4. ShowCreateTableExec L177: the extractor-based hasPartitioning throws ClassCastException for a table with a cluster_by transform that has a literal argument (a regression vs. master), and disagrees with the parser on PARTITIONED BY (cluster_by(a)) DISTRIBUTED BY PARTITION, so a replay loses HASH.
  5. error-conditions.json L7345: for a CLUSTER BY table, the advice to add PARTITIONED BY or CLUSTERED BY ... INTO ... BUCKETS cannot be followed.
  6. ShowCreateTableExec L146: the "same type" claim fails for connector-declared DECIMAL(p,s) and CHAR/VARCHAR literals and under two legacy confs.
  7. WriteDistributionAndOrdering L73: a literal conversion that throws makes DESCRIBE TABLE EXTENDED fail entirely, and a connector's own NamedReference is printed via describe().
  8. ShowCreateTableExec L137: connector keys named bucket, years, months, days or hours with other argument shapes print DDL that the parser rejects.
  9. AstBuilder L6375: a parameter marker in a transform argument fails with a raw ClassCastException. This is pre-existing in PARTITIONED BY, so a follow-up is fine.

Tests

  1. CreateTableWriteOrderSuite L244: the error checks are hand-rolled, one of them checks no condition, and nothing pins the SQLSTATEs of the three new conditions.
  2. PlanResolutionSuite L3696: the six new tests copy unrelated assertions and only re-pin parser output covered elsewhere.

Docs and comments

  1. CreateTableLikeExec L50: the "not copied" caveat is only in this scaladoc; the public TableCatalog.createTableLike and TableInfo Javadocs still describe tableInfo as complete.
  2. PR description: it still describes the mode as a String with DISTRIBUTION_MODE_* constants (also in design decision 1), says the suite has 30 tests (now 38), and does not mention the two new abstract members of V2CreateTablePlan (writeOrdering, withWriteOrdering) from the minor items under points 1-6.

API and design

  1. TableCatalogCapability L108: DelegatingCatalogExtension forwards capabilities() but not createTable(Identifier, TableInfo), which leaves a latent hole in this Javadoc's promise.
  2. TableInfo L139: the "never null" ordering is not enforced at build().

Cleanup

  1. CheckAnalysis L995: the comment gives a stale reason, and the block copies the partitioning check with workarounds.

The most important are 15-18. 15 is a regression in this commit, and 16-18 make SHOW CREATE TABLE emit DDL that fails, or silently changes the declared layout, on replay.

case (_, s: StringType) => DataTypeUtils.isDefaultStringCharOrVarcharType(s)
case (_, BooleanType | ByteType | ShortType | IntegerType | LongType | _: DecimalType |
BinaryType | DateType | TimestampType | TimestampNTZType | _: DayTimeIntervalType |
_: YearMonthIntervalType) => true

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.

TimeType is missing from this list, so a TIME literal in a sort key drops every write clause from SHOW CREATE TABLE. The parser accepts ORDERED BY f(id, TIME '12:00:00') without any flag (spark.sql.timeType.enabled gates schemas and casts, not literals), and Literal.sql spells it back as TIME '12:00:00', which re-parses as TIME(6). But the key falls to case _ => false at L157, the forall at L167 fails, and ORDERED BY (or LOCALLY ORDERED BY, or DISTRIBUTED BY PARTITION ...) disappears from the output. Replaying it creates a table without the declared layout, which is the silent drop from point 7. This is a regression in this commit: d2aa867 accepted every Literal and printed the clause. The nanosecond timestamp types (preview flag) and the legacy CalendarIntervalType are dropped the same way.

Could we add case (_, TimeType(TimeType.DEFAULT_PRECISION)) => true, and add TIME '12:00:00' to the round-trip test at CreateTableWriteOrderSuite.scala:627? Precision 7-9 does not round-trip, since TIME '12:00:00.1000000' prints as TIME '12:00:00.1' and re-parses as TIME(6).

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Good catch, and thanks for spotting that it was a regression. With the parse-back check there is no type list to miss. TIME '12:00:00' is in the round-trip test, and a TIME with 7 to 9 fractional digits is omitted, because it parses back as TIME(6). The nanosecond timestamp types and CalendarIntervalType are handled by the same check.

case (f: Float, FloatType) => java.lang.Float.isFinite(f) && math.abs(f) < Float.MaxValue
case (d: Double, DoubleType) => java.lang.Double.isFinite(d)
case (_, s: StringType) => DataTypeUtils.isDefaultStringCharOrVarcharType(s)
case (_, BooleanType | ByteType | ShortType | IntegerType | LongType | _: DecimalType |

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.

true, false and NULL do not re-parse as constants inside a transform, so BooleanType and NullType do not belong in this list. In transformArgument : qualifiedName | constant, these non-reserved keywords match both alternatives, and ANTLR resolves the ambiguity to the first one, qualifiedName. (primaryExpression lists constant before identifier, which is why SELECT true is a literal.) If a connector reports RANGE with the key f(id, literal(true)), SHOW CREATE TABLE prints ORDERED BY (f(id, true) ASC NULLS FIRST), and replaying it reads true as a column: the statement fails with WRITE_ORDERING_WITH_UNKNOWN_COLUMN, or binds to a column named true. TRUE is also in ansiNonReserved, so this happens in every mode. FALSE and NULL behave the same unless spark.sql.ansi.enforceReservedKeywords is true, so a table created with ORDERED BY f(id, NULL) in an enforcing session does not replay in a default one.

Could we drop BooleanType here and change L150 to case (null, _) => false? Deleting L150 would let a typed NULL fall through to the type list. The comment at L146 and the constant in the new docs could mention this too.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Thanks. A TRUE, FALSE or NULL argument now parses back as a column reference when the keyword is not reserved, so the comparison fails and the key is omitted. Under spark.sql.ansi.enforceReservedKeywords they parse back as constants and are emitted. The docs now say so.

case (d: Double, DoubleType) => java.lang.Double.isFinite(d)
case (_, s: StringType) => DataTypeUtils.isDefaultStringCharOrVarcharType(s)
case (_, BooleanType | ByteType | ShortType | IntegerType | LongType | _: DecimalType |
BinaryType | DateType | TimestampType | TimestampNTZType | _: DayTimeIntervalType |

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.

TIMESTAMP literals are printed as session-local wall time without an offset, so they can replay as a different value or type. Literal.sql prints TIMESTAMP '<wall time in spark.sql.session.timeZone>'. With the default confs and America/Los_Angeles, ORDERED BY f(id, TIMESTAMP '2020-11-01 01:30:00-08:00') stores 09:30Z, but SHOW CREATE TABLE prints TIMESTAMP '2020-11-01 01:30:00'. Replaying that in the same session resolves the ambiguous fall-back time to the earlier offset (PDT), so the key silently becomes 08:30Z. Replaying in another session time zone shifts it by the offset difference. With spark.sql.timestampType=TIMESTAMP_NTZ, a key written as TIMESTAMP_LTZ '2020-01-01 10:00:00' is printed as TIMESTAMP '...' and re-parsed as TIMESTAMP_NTZ. The round-trip test passes because it replays in the same session, outside a DST overlap, with the default timestamp type.

Could we render TimestampType with an explicit offset in toSQL (WriteDistributionAndOrdering.scala:72), e.g. TIMESTAMP_LTZ '<UTC wall time>Z', like the FLOAT special case? That re-parses as TimestampType in both timestamp modes.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done. toSQL renders a TIMESTAMP as TIMESTAMP_LTZ '<UTC wall time>Z'. A new test creates the table in America/Los_Angeles with TIMESTAMP '2020-11-01 01:30:00-08:00' and replays the DDL in Los Angeles, in Asia/Tokyo and with spark.sql.timestampType=TIMESTAMP_NTZ. Each replay declares the same key.

}
// Bucketing counts as partitioning here; CLUSTER BY does not.
val hasPartitioning = table.partitioning.exists {
case ClusterByTransform(_) => false

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.

This extractor-based check introduced two problems in this commit.

(a) ClusterByTransform.unapply casts every argument with arguments.map(_.asInstanceOf[NamedReference]) (expressions.scala:187), and hasPartitioning is a strict val evaluated for every table whose ordering passes L167, including tables that declare nothing, since an empty ordering passes forall. So SHOW CREATE TABLE now throws a raw ClassCastException for a table whose first partition transform is a cluster_by with a literal argument, e.g. a connector-reported Expressions.apply("cluster_by", Expressions.literal(4)). Master and d2aa867 printed PARTITIONED BY (cluster_by(4)) there.

(b) The parser still counts such a transform as partitioning. PARTITIONED BY (cluster_by(a)) DISTRIBUTED BY PARTITION becomes ApplyTransform("cluster_by", [a]), passes the check at AstBuilder.scala:6320 and sends HASH, and a DelegatingTable-based catalog (like GenericClusterByTableCatalog in the new suite) reports both back. SHOW CREATE TABLE then prints PARTITIONED BY (cluster_by(a)) but treats it as clustering here and omits DISTRIBUTED BY PARTITION, so a replay silently loses HASH. The notation is not exotic: v2 SHOW CREATE TABLE already prints every CLUSTER BY table as PARTITIONED BY (cluster_by(c)).

Could we share one name-based predicate without a cast between writeSpecsFrom and this method, e.g. WriteDistributionAndOrdering.hasPartitioning(p) = p.exists(_.name != "cluster_by"), and evaluate it only in the HASH cases? That keeps szehon-ho's request that a connector's generic cluster_by counts as clustering.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Good catch on both. writeSpecsFrom and SHOW CREATE TABLE now share WriteDistributionAndOrdering.hasPartitioning, which matches by name (any transform not named cluster_by), casts nothing, and is evaluated only in the HASH cases. A connector's cluster_by(4) prints as PARTITIONED BY (cluster_by(4)) again. The parser now rejects PARTITIONED BY (cluster_by(a)) DISTRIBUTED BY PARTITION, so a replay can no longer lose HASH, and @szehon-ho 's case still holds.

"SPECIFY_DISTRIBUTED_BY_PARTITION_WITHOUT_PARTITIONING_IS_NOT_ALLOWED" : {
"message" : [
"Cannot specify DISTRIBUTED BY PARTITION for a table that has no partitioning.",
"Please add PARTITIONED BY or CLUSTERED BY ... INTO ... BUCKETS, or drop the clause. A statement that defines no schema (no column list, no typed partition columns, and no AS SELECT) cannot declare partitioning, so it cannot use DISTRIBUTED BY PARTITION."

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.

For a CLUSTER BY table, this advice cannot be followed either. CREATE TABLE t (id INT, c STRING) USING foo CLUSTER BY (c) DISTRIBUTED BY PARTITION, pinned at CreateTableWriteOrderSuite.scala:343, gets this message. Adding PARTITIONED BY (id) as suggested then fails with SPECIFY_CLUSTER_BY_WITH_PARTITIONED_BY_IS_NOT_ALLOWED, and adding CLUSTERED BY (id) INTO 4 BUCKETS fails with SPECIFY_CLUSTER_BY_WITH_BUCKETING_IS_NOT_ALLOWED, because those checks (AstBuilder.scala:6290-6296) run before writeSpecsFrom. Only the SQL reference explains that CLUSTER BY has to be replaced. This is the CLUSTER BY counterpart of point 8.

Could writeSpecsFrom raise a dedicated condition when ctx.clusterBySpec is present, e.g. SPECIFY_CLUSTER_BY_WITH_DISTRIBUTED_BY_PARTITION_IS_NOT_ALLOWED next to the existing SPECIFY_CLUSTER_BY_WITH_* conditions, or at least add a sentence about CLUSTER BY here, as for the schemaless case? Either way the rejection stays at parse time, so HASH still always comes with non-empty partitions().

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done. writeSpecsFrom raises a dedicated SPECIFY_CLUSTER_BY_WITH_DISTRIBUTED_BY_PARTITION_IS_NOT_ALLOWED (42908) when CLUSTER BY is present. The message says to replace CLUSTER BY with PARTITIONED BY or bucketing, or to drop DISTRIBUTED BY PARTITION. The missing-partitioning message also says now that a cluster_by(...) transform in PARTITIONED BY is clustering. Both stay at parse time.

}
}

test("SPARK-34586: v2 table creation (global writeOrdering)") {

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.

Nit: these six tests (about 230 lines) copy assertions from "Test v2 CreateTable with default catalog" (L620) and "Test v2 CTAS with known catalog in identifier" (L683) that are unrelated to this feature: catalog name, schema, properties and ignoreIfExists. parseAndResolve (L291) runs neither PreprocessTableCreation nor checkAnalysis by default, so these tests only re-pin parser output that the parse tests at CreateTableWriteOrderSuite.scala:80-190 already pin, and they can break for reasons unrelated to the write clauses.

Could we drop them, or fold them into one table-driven test over (sql, mode, ordering) that asserts only the write spec?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Folded the six into one table-driven test over CREATE, CTAS, REPLACE and RTAS with the default and a known catalog. It asserts only the plan type, the mode, the ordering and, for DISTRIBUTED BY PARTITION, the partitioning.

* [[TableCatalog.PROP_OWNER]] set to the current user. Source table properties are intentionally
* excluded so that connectors can decide which custom properties to clone via [[sourceTable]].
*
* The source's declared write distribution and ordering are likewise not copied onto the

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.

Minor: this caveat is only in the exec's scaladoc, while connector authors read the public TableCatalog.createTableLike Javadoc. That Javadoc still says tableInfo contains "columns and partitioning copied from the source, any constraints copied from the source, ..." (TableCatalog.java:368-371) and calls it the "complete description of the new table: columns, partitioning, constraints, ..." (TableCatalog.java:379-381). A connector written against that contract expects tableInfo.writeOrdering() to carry the source's layout like the constraints, but it is always empty, so CREATE TABLE ... LIKE silently creates the target without it unless the connector reads sourceTable. The TableInfo class Javadoc (TableInfo.java:28) also lists only "columns, properties, partitioning and constraints".

Could we add the same sentence to those Javadocs? That would complete peter-toth's finding 9.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done. The TableCatalog.createTableLike Javadoc and its @param tableInfo say the source's write distribution and ordering are not copied (writeDistributionMode() is null, writeOrdering() is empty) and that a connector can read them from sourceTable. The TableInfo class summary lists them too.

* <li>{@link StagingTableCatalog#stageReplace(Identifier, TableInfo)}</li>
* <li>{@link StagingTableCatalog#stageCreateOrReplace(Identifier, TableInfo)}</li>
* </ul>
* Their default implementations drop the request, so a catalog that reports this capability must

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.

Minor: Spark's own DelegatingCatalogExtension cannot keep this promise. It forwards the delegate's capabilities() (DelegatingCatalogExtension.java:57-58) but does not override createTable(Identifier, TableInfo), so it uses the default at TableCatalog.java:359-361, which drops the request. Today the only in-tree delegate, V2SessionCatalog, does not report the capability, so the statement is rejected. Once a delegate reports it, e.g. if V2SessionCatalog gains support later, subclasses that inherit super.capabilities(), such as Delta's AbstractDeltaCatalog and Hudi's HoodieCatalog, pass the check and silently lose the layout. That is the case the PR description says the capability prevents, and SUPPORT_TABLE_CONSTRAINT has had the same gap since 4.1.0. Forwarding createTable(Identifier, TableInfo) to the delegate is not a fix, because it would bypass those subclasses' column-based overrides.

Could DelegatingCatalogExtension.capabilities() remove this new capability from the forwarded set, so that only a subclass overriding all four TableInfo overloads adds it back, or could this Javadoc at least mention it?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Good catch. DelegatingCatalogExtension.capabilities() now removes SUPPORTS_CREATE_TABLE_WITH_WRITE_DISTRIBUTION_AND_ORDERING from the delegate's set, so a subclass reports it only if it overrides the TableInfo overloads and adds it back. createTable(Identifier, TableInfo) is still not forwarded. Both Javadocs say so, and a new test checks that such an extension rejects ORDERED BY before anything is created. I left SUPPORT_TABLE_CONSTRAINT as it is, since that gap predates this PR.

*
* @since 4.4.0
*/
public Builder withWriteOrdering(SortOrder[] writeOrdering) {

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.

Nit: the Javadoc says this must not be null and that writeOrdering() is never null, but neither this setter nor build() (L156) enforces it. A connector that passes null and returns new DelegatingTable(info, name) gets an NPE in DESCRIBE TABLE EXTENDED (DescribeTableExec.scala:232) or SHOW CREATE TABLE (ShowCreateTableExec.scala:167) instead of at build(). withPartitions and withConstraints are unchecked too, but only this field documents "never null".

Could build() add Objects.requireNonNull(writeOrdering, ...) next to the columns check?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done. TableInfo now rejects a null write ordering in its constructor, so both build() and a subclass are covered, and a new test checks it.

"cols" -> badReferences.map(r => toSQLId(r)).mkString(", ")))
}

// PreprocessTableCreation only normalizes column and RewritableTransform references,

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.

Nit: this comment gives a reason that no longer holds. PreprocessTableCreation deliberately keeps an unresolvable ordering reference for this check (rules.scala:354-355), so this is the only check rather than an additional one, and analyzers without that sql/core rule need it as well. The block also copies the partitioning check above with an explicit (ref: NamedReference) => ascription (the import at L33 exists only for it), a trailing .toSeq, and a quote-then-reparse through column.quoted and toSQLId(String). With five or more references, cols follows the hash set order rather than the declaration order.

Could we write it as create.writeOrdering.flatMap(_.expression().references().map(_.fieldNames().toImmutableArraySeq)).distinct.filter(create.tableSchema.findNestedField(_).isEmpty) with toSQLId(parts: Seq[String]), drop the import, and say why the check is here?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done as you suggested. The comment now says this is the only check for an unknown ordering column, because PreprocessTableCreation keeps such a reference for it and analyzers without that rule need it. The NamedReference import is gone, and a new test pins declaration order with six unknown references.

@dongjoon-hyun

Copy link
Copy Markdown
Member

Gentle ping, @anuragmantri ~

@anuragmantri anuragmantri left a comment

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Thank you for another careful pass, @dongjoon-hyun. Points 15-30 are addressed with replies on each thread.

One design change is worth calling out, because several of your points (15, 16, 17, 20, 21, 22) were the same problem: an allow-list could not cover every literal type, name and conf. SHOW CREATE TABLE now renders the keys, parses the clauses back with the session's parser, and emits them only when the parsed mode and keys equal the declared ones. It retries once with every name quoted, for reserved keywords. It also omits the clauses when the catalog does not report the capability or a key references a column the table does not have, since a replay would fail in both cases. WriteDistributionAndOrderingUtilsSuite checks "emitted iff it parses back to the same key" over connector-reported keys of every literal type and shape, including under the confs that change parsing.

case (_, s: StringType) => DataTypeUtils.isDefaultStringCharOrVarcharType(s)
case (_, BooleanType | ByteType | ShortType | IntegerType | LongType | _: DecimalType |
BinaryType | DateType | TimestampType | TimestampNTZType | _: DayTimeIntervalType |
_: YearMonthIntervalType) => true

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Good catch, and thanks for spotting that it was a regression. With the parse-back check there is no type list to miss. TIME '12:00:00' is in the round-trip test, and a TIME with 7 to 9 fractional digits is omitted, because it parses back as TIME(6). The nanosecond timestamp types and CalendarIntervalType are handled by the same check.

case (f: Float, FloatType) => java.lang.Float.isFinite(f) && math.abs(f) < Float.MaxValue
case (d: Double, DoubleType) => java.lang.Double.isFinite(d)
case (_, s: StringType) => DataTypeUtils.isDefaultStringCharOrVarcharType(s)
case (_, BooleanType | ByteType | ShortType | IntegerType | LongType | _: DecimalType |

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Thanks. A TRUE, FALSE or NULL argument now parses back as a column reference when the keyword is not reserved, so the comparison fails and the key is omitted. Under spark.sql.ansi.enforceReservedKeywords they parse back as constants and are emitted. The docs now say so.

case (d: Double, DoubleType) => java.lang.Double.isFinite(d)
case (_, s: StringType) => DataTypeUtils.isDefaultStringCharOrVarcharType(s)
case (_, BooleanType | ByteType | ShortType | IntegerType | LongType | _: DecimalType |
BinaryType | DateType | TimestampType | TimestampNTZType | _: DayTimeIntervalType |

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done. toSQL renders a TIMESTAMP as TIMESTAMP_LTZ '<UTC wall time>Z'. A new test creates the table in America/Los_Angeles with TIMESTAMP '2020-11-01 01:30:00-08:00' and replays the DDL in Los Angeles, in Asia/Tokyo and with spark.sql.timestampType=TIMESTAMP_NTZ. Each replay declares the same key.

}
// Bucketing counts as partitioning here; CLUSTER BY does not.
val hasPartitioning = table.partitioning.exists {
case ClusterByTransform(_) => false

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Good catch on both. writeSpecsFrom and SHOW CREATE TABLE now share WriteDistributionAndOrdering.hasPartitioning, which matches by name (any transform not named cluster_by), casts nothing, and is evaluated only in the HASH cases. A connector's cluster_by(4) prints as PARTITIONED BY (cluster_by(4)) again. The parser now rejects PARTITIONED BY (cluster_by(a)) DISTRIBUTED BY PARTITION, so a replay can no longer lose HASH, and @szehon-ho 's case still holds.

"SPECIFY_DISTRIBUTED_BY_PARTITION_WITHOUT_PARTITIONING_IS_NOT_ALLOWED" : {
"message" : [
"Cannot specify DISTRIBUTED BY PARTITION for a table that has no partitioning.",
"Please add PARTITIONED BY or CLUSTERED BY ... INTO ... BUCKETS, or drop the clause. A statement that defines no schema (no column list, no typed partition columns, and no AS SELECT) cannot declare partitioning, so it cannot use DISTRIBUTED BY PARTITION."

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done. writeSpecsFrom raises a dedicated SPECIFY_CLUSTER_BY_WITH_DISTRIBUTED_BY_PARTITION_IS_NOT_ALLOWED (42908) when CLUSTER BY is present. The message says to replace CLUSTER BY with PARTITIONED BY or bucketing, or to drop DISTRIBUTED BY PARTITION. The missing-partitioning message also says now that a cluster_by(...) transform in PARTITIONED BY is clustering. Both stay at parse time.

* [[TableCatalog.PROP_OWNER]] set to the current user. Source table properties are intentionally
* excluded so that connectors can decide which custom properties to clone via [[sourceTable]].
*
* The source's declared write distribution and ordering are likewise not copied onto the

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done. The TableCatalog.createTableLike Javadoc and its @param tableInfo say the source's write distribution and ordering are not copied (writeDistributionMode() is null, writeOrdering() is empty) and that a connector can read them from sourceTable. The TableInfo class summary lists them too.

* <li>{@link StagingTableCatalog#stageReplace(Identifier, TableInfo)}</li>
* <li>{@link StagingTableCatalog#stageCreateOrReplace(Identifier, TableInfo)}</li>
* </ul>
* Their default implementations drop the request, so a catalog that reports this capability must

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Good catch. DelegatingCatalogExtension.capabilities() now removes SUPPORTS_CREATE_TABLE_WITH_WRITE_DISTRIBUTION_AND_ORDERING from the delegate's set, so a subclass reports it only if it overrides the TableInfo overloads and adds it back. createTable(Identifier, TableInfo) is still not forwarded. Both Javadocs say so, and a new test checks that such an extension rejects ORDERED BY before anything is created. I left SUPPORT_TABLE_CONSTRAINT as it is, since that gap predates this PR.

"cols" -> badReferences.map(r => toSQLId(r)).mkString(", ")))
}

// PreprocessTableCreation only normalizes column and RewritableTransform references,

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done as you suggested. The comment now says this is the only check for an unknown ordering column, because PreprocessTableCreation keeps such a reference for it and analyzers without that rule need it. The NamedReference import is gone, and a new test pins declaration order with six unknown references.

*
* @since 4.4.0
*/
public Builder withWriteOrdering(SortOrder[] writeOrdering) {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done. TableInfo now rejects a null write ordering in its constructor, so both build() and a subclass are covered, and a new test checks it.

}

// A plain column is passed as a bare reference rather than as `identity(col)`.
val key = visitTransform(ctx.transform) match {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Thanks, filed SPARK-59977 for this since it seems like an existing bug.

@dongjoon-hyun dongjoon-hyun 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.

Thank you for addressing points 15-30, @anuragmantri. I reviewed the latest commit, 60a4ea7, and left 8 inline comments, continuing the numbering from points 1-30. The parse-back check works well: I could not find a table created in SQL whose SHOW CREATE TABLE output fails or declares a different layout when replayed in the same session. 33 and 34 are the remaining cases of points 15 and 17, and 39 is point 27, which is still open. Here is a summary, ordered by priority.

Correctness

  1. WriteDistributionAndOrdering L122: referencesExist does not see a column inside a nested IdentityTransform, so SHOW CREATE TABLE can emit an ORDERED BY that fails on replay. This is new in this commit.
  2. WriteDistributionAndOrdering L173: DESCRIBE now prints an Extract or UDF key without its field or function name.
  3. WriteDistributionAndOrdering L188: a TIME(7-9) value with trailing zeros still drops the whole clause (point 15).
  4. WriteDistributionAndOrdering L185: the nanosecond TIMESTAMP_LTZ type is still printed in session-local time without an offset (point 17).

Tests

  1. WriteDistributionAndOrderingUtilsSuite L156: the HASH rows cannot fail, because this suite's replay has no PARTITIONED BY, and no test in the suite emits a HASH or NONE pair.
  2. PreprocessTableCreation L334: SPECIFY_WRITE_ORDERING_IS_NOT_ALLOWED is the only new condition without a query context, and the converted checkError now pins that.

Docs and comments

  1. TableInfo L100: the new sentence that the parser checks the arguments of bucket, years, months, days and hours promises more than the parser does.
  2. sql-ref-syntax-ddl-create-table-datasource.md L41: the syntax block still shows the parentheses of ORDERED BY as required.
  3. PR description (point 27): it still describes a String mode with DISTRIBUTION_MODE_* constants, including design decision 1, says the suite has 30 tests, and does not mention the two new abstract members of V2CreateTablePlan (writeOrdering, withWriteOrdering). It should now also describe the parse-back check and the DelegatingCatalogExtension change. dev/merge_spark_pr.py puts this text into the commit message.

Other notes

  • The GeneralScalarExpression.toString() recursion you found is pre-existing: V2ExpressionSQLBuilder.visitUnexpectedExpr builds its error message with String.valueOf(expr), which calls toString() again. It can be fixed separately from this PR.
  • In the reply to point 16, TRUE is not emitted under spark.sql.ansi.enforceReservedKeywords either. It is also in ansiNonReserved, so it always parses back as a column; only FALSE and NULL are reproduced in that mode. The new docs sentence is right as written.
  • The CI failures look unrelated. SQLAppStatusListenerWithInMemoryStoreSuite "driver side SQL metrics" is a racy test that SPARK-59784 fixed on master after this branch's base, and KafkaRealTimeModeWindowSuite "tumbling window count" timed out.

The most important are 31 and 39: 31 makes SHOW CREATE TABLE emit DDL that fails on replay, and 39 becomes the commit message.

* matches the ordering of a CREATE/REPLACE TABLE statement.
*/
def referencesExist(schema: StructType, writeOrdering: Seq[SortOrder]): Boolean = {
writeOrdering.flatMap(_.expression().references()).forall { ref =>

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.

referencesExist checks the declared key's references(), but ApplyTransform.references collects only top-level references, so a column inside a nested IdentityTransform is never checked. For a table t(id INT), a connector key Expressions.apply("f", Expressions.identity("missing")) passes this check vacuously. toSQL and keyOf both flatten the nested identity, so the parse-back comparison succeeds and SHOW CREATE TABLE prints ORDERED BY (f(missing) ASC NULLS FIRST). Replaying it fails with WRITE_ORDERING_WITH_UNKNOWN_COLUMN (CheckAnalysis.scala:996-1003). identity("ID") for a column id fails the same way, since an ApplyTransform is not normalized (SPARK-59944). This is new in this commit: the previous isSpellable rejected an argument that was neither a reference nor a literal.

Could we check the replayed ordering instead, which is what CheckAnalysis will see, e.g. referencesExist(schema, ordering) inside the find or in ShowCreateTableExec.replay?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Good catch. The column check now runs on the replayed ordering, which is what CheckAnalysis sees, so f(identity(missing)) is omitted. A nested identity(...) is also no longer compared as a plain column: only a key's top level counts as identity(col), so f(identity(id)) is omitted too. Both are in the unit suite and in ORDERING_OVERRIDE with omit-and-replay rows.

case g: GeneralScalarExpression =>
g.children().map(toSQL(_, quote)).mkString(s"${g.name}(", ", ", ")")
case other =>
other.children().map(toSQL(_, quote)).mkString(s"${other.getClass.getSimpleName}(", ", ", ")")

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.

Minor: this fallback now applies to every connector expression, not only to the GeneralScalarExpression names that V2ExpressionSQLBuilder does not know. An Extract key prints as Extract(ts) instead of EXTRACT(YEAR FROM ts), because Extract.children() is just the source (Extract.java:60), and a UserDefinedScalarFunc prints as UserDefinedScalarFunc(id) instead of my_udf(id). So DESCRIBE TABLE EXTENDED loses the field and the function name that 5569461 printed with describe. The GeneralScalarExpression case at L170-171 has a similar effect on known names, e.g. +(id, 1) instead of id + 1.

Could describe stay the default here, with the children-based rendering only for a GeneralScalarExpression name that V2ExpressionSQLBuilder does not know? Once the pre-existing recursion is fixed separately, that special case could go away too.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Thanks. DESCRIBE renders connector expressions with ToStringSQLBuilder again, so EXTRACT(YEAR FROM ts), my_udf(id) and id + 1 print as describe printed them. The builder is a subclass that renders literals and references with their types and quoting, and renders GetArrayItem and VariantGet through their children rather than through toString. It renders an expression it does not know as name(children) instead of throwing, so a known parent keeps its name, e.g. my_udf(DATE_TRUNC(id)). Tests cover these cases, plus a sort order and a null child nested in a key.

case (micros: Long, TimestampType) =>
val utc = TimestampFormatter.getFractionFormatter(ZoneOffset.UTC).format(micros)
s"TIMESTAMP_LTZ '${utc}Z'"
case _ => l.sql

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.

Minor: point 15 is not complete for a TIME(7-9) value with trailing zeros. TIME '12:00:00.100000000' is a TIME(9) literal, but Literal.sql prints TIME '12:00:00.1' (literals.scala:691, whose formatter drops trailing zeros), which parses back as TIME(6). The comparison fails and the whole clause is dropped, although TIME '12:00:00.100000000' would replay correctly. A value without trailing zeros, such as TIME '12:00:00.123456789', does round-trip, so the reply's "a TIME with 7 to 9 fractional digits is omitted" holds only for these values.

Could literalToSQL pad the fraction to the type's precision for TimeType, as Literal.sql already does with padToNanosPrecision for the nanosecond timestamp types (literals.scala:700-702)?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done. literalToSQL pads a TIME with more than microsecond precision to its precision, so TIME '12:00:00.100000000' round-trips as TIME(9). The unit suite has TIME(7) and TIME(9) rows with and without trailing zeros.


private def literalToSQL(l: CatalystLiteral): String = (l.value, l.dataType) match {
case (f: Float, FloatType) if java.lang.Float.isFinite(f) => s"${f}F"
case (micros: Long, TimestampType) =>

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.

Minor (preview): point 17 is not complete for the nanosecond timestamps. With spark.sql.timestampNanosTypes.enabled, a TIMESTAMP literal with 7-9 fractional digits is a TimestampLTZNanosType, which this case does not match, so it falls to l.sql and prints TIMESTAMP_LTZ '<session-local wall time>' without an offset (literals.scala:702). In America/Los_Angeles, ORDERED BY f(id, TIMESTAMP '2020-11-01 01:30:00.123456789-08:00') then fails the comparison in the DST overlap and drops the clause, and a value outside the overlap replays shifted in another session time zone.

Could TimestampLTZNanosType also be printed in UTC with a Z suffix, padded to its precision?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done. A TimestampLTZNanosType literal now renders in UTC with a Z suffix, padded to its precision. A test checks TIMESTAMP_LTZ '2020-11-01 01:30:00.123456789-08:00' in America/Los_Angeles with spark.sql.timestampNanosTypes.enabled, and a nanosecond TIMESTAMP_NTZ row was added too.

test("the pairs with no clause form are not emitted") {
val keys = Seq(key(id))
Seq(
(WriteDistributionMode.HASH, keys, Seq.empty[Transform]),

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.

Minor: these two HASH rows cannot fail. This suite's replay (L51) parses CREATE TABLE t (id INT) USING foo $clauses, which has no PARTITIONED BY, so the parser rejects every DISTRIBUTED BY PARTITION ... rendering and writeClausesSQL returns None whatever it emits. Production replays with PARTITIONED BY (p) (ShowCreateTableExec.scala:162). Removing the if partitioned guards at WriteDistributionAndOrdering.scala:102-104 would still pass here, and since emitted (L55-57) always uses RANGE, no test in this suite emits a HASH or NONE pair.

Could this test use a replay that accepts any clauses, e.g. _ => Some((mode, ordering)), and could the suite add a positive HASH row such as (HASH, Seq(key(id)), Seq(Expressions.identity("p"))) expecting DISTRIBUTED BY PARTITION ORDERED BY (id ASC NULLS FIRST)?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Thanks. The suite now replays the same statement as SHOW CREATE TABLE (replayStatement, shared with ShowCreateTableExec). The no-clause-form test uses a replay that accepts any clauses, so only the guards decide, and it covers a null mode. A new test emits each mode's clause form, including DISTRIBUTED BY PARTITION ORDERED BY (id ASC NULLS FIRST) and LOCALLY ORDERED BY .... Another omits an unspellable key under HASH, RANGE and NONE.

throw QueryCompilationErrors.specifyPartitionNotAllowedWhenTableSchemaNotDefinedError()
}
if (create.writeOrdering.nonEmpty) {
throw QueryCompilationErrors

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.

Nit: SPECIFY_WRITE_ORDERING_IS_NOT_ALLOWED is the only one of the four new conditions without a query context. specifyWriteOrderingNotAllowedWhenTableSchemaNotDefinedError() creates the AnalysisException without an origin, while WRITE_ORDERING_WITH_UNKNOWN_COLUMN and the two parse errors point at the statement, and the converted checkError at CreateTableWriteOrderSuite.scala:329 now pins an empty context. The legacy _LEGACY_ERROR_TEMP_1165 at L331 has the same gap, so this is optional, but changing it later means changing the test again.

Could this error take create.origin?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done. The error now takes create.origin, and the test pins the statement as its context.

* <p>
* A plain column is a {@link org.apache.spark.sql.connector.expressions.NamedReference}; any
* other key is a {@link Transform}, such as {@code bucket(16, id)}. Spark checks that each
* referenced column exists in the table schema, and the parser checks the arguments of

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.

This new sentence promises more validation than the parser does. visitTransform dispatches these names case-sensitively (AstBuilder.scala:5869), so ORDERED BY BUCKET(id, 16) or ORDERED BY DAYS(1) reaches the catalog as a generic transform with no argument check, and DAYS(1) references no column, so the column check does not apply either. For lowercase bucket, the parser checks only the argument kinds: bucket(16) without a column, a zero or negative count, and bucket(4294967300, id), which .toInt truncates to 4 (AstBuilder.scala:5877), are all accepted. WriteDistributionAndOrderingUtilsSuite even lists "bucket without columns" as spellable. A catalog author reading this sentence may skip their own validation.

Could we drop the parser clause, or say that only the lowercase names have their argument kinds checked and that bucket counts and column lists are not validated?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Thanks. It now says that the parser checks the argument kinds of the lowercase bucket, years, months, days and hours, and that Spark does not check orderability, a positive bucket count, a bucket's column list, or the arguments of any other transform.

[ COMMENT table_comment ]
[ TBLPROPERTIES ( key1=val1, key2=val2, ... ) ]
[ DISTRIBUTED BY PARTITION ]
[ [ LOCALLY ] ORDERED BY ( write_order_field [ , ... ] ) | UNORDERED ]

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.

Nit: the parentheses are optional, as L130-131 says and the AstBuilder scaladoc now shows with [(] ... [)], but this syntax block still shows them as required. Could it be [ [ LOCALLY ] ORDERED BY { ( write_order_field [ , ... ] ) | write_order_field [ , ... ] } | UNORDERED ]?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done: [ LOCALLY ] ORDERED BY { ( write_order_field [ , ... ] ) | write_order_field [ , ... ] }.

@anuragmantri anuragmantri left a comment

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Thank you, @dongjoon-hyun. Points 31-39 are addressed in 88882a2, with replies on each thread.

On your other notes:

Point 16: you are right, thanks for the correction. TRUE is in ansiNonReserved, so it always parses back as a column and is never reproduced. Only FALSE and NULL are reproduced under spark.sql.ansi.enforceReservedKeywords. The docs sentence was already right.

The GeneralScalarExpression.toString() recursion: I filed SPARK-60075 for the fix in V2ExpressionSQLBuilder.visitUnexpectedExpr. This PR only avoids it in DESCRIBE.

While checking the class behind 31 and 32, a few more cases came up and are fixed in the same commit.

  • SHOW CREATE TABLE now also omits the clauses when the statement would create a v1 table, because the session catalog with a provider that is not a v2 source rejects them on replay.
  • DESCRIBE no longer recurses through GetArrayItem, VariantGet or a nested sort order. A nested identity(...) is no longer compared as a plain column. And
  • ORDERED BY id.x on an INT column reports WRITE_ORDERING_WITH_UNKNOWN_COLUMN with the statement as its context instead of a context-free INVALID_FIELD_NAME.

* matches the ordering of a CREATE/REPLACE TABLE statement.
*/
def referencesExist(schema: StructType, writeOrdering: Seq[SortOrder]): Boolean = {
writeOrdering.flatMap(_.expression().references()).forall { ref =>

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Good catch. The column check now runs on the replayed ordering, which is what CheckAnalysis sees, so f(identity(missing)) is omitted. A nested identity(...) is also no longer compared as a plain column: only a key's top level counts as identity(col), so f(identity(id)) is omitted too. Both are in the unit suite and in ORDERING_OVERRIDE with omit-and-replay rows.

case g: GeneralScalarExpression =>
g.children().map(toSQL(_, quote)).mkString(s"${g.name}(", ", ", ")")
case other =>
other.children().map(toSQL(_, quote)).mkString(s"${other.getClass.getSimpleName}(", ", ", ")")

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Thanks. DESCRIBE renders connector expressions with ToStringSQLBuilder again, so EXTRACT(YEAR FROM ts), my_udf(id) and id + 1 print as describe printed them. The builder is a subclass that renders literals and references with their types and quoting, and renders GetArrayItem and VariantGet through their children rather than through toString. It renders an expression it does not know as name(children) instead of throwing, so a known parent keeps its name, e.g. my_udf(DATE_TRUNC(id)). Tests cover these cases, plus a sort order and a null child nested in a key.

case (micros: Long, TimestampType) =>
val utc = TimestampFormatter.getFractionFormatter(ZoneOffset.UTC).format(micros)
s"TIMESTAMP_LTZ '${utc}Z'"
case _ => l.sql

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done. literalToSQL pads a TIME with more than microsecond precision to its precision, so TIME '12:00:00.100000000' round-trips as TIME(9). The unit suite has TIME(7) and TIME(9) rows with and without trailing zeros.


private def literalToSQL(l: CatalystLiteral): String = (l.value, l.dataType) match {
case (f: Float, FloatType) if java.lang.Float.isFinite(f) => s"${f}F"
case (micros: Long, TimestampType) =>

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done. A TimestampLTZNanosType literal now renders in UTC with a Z suffix, padded to its precision. A test checks TIMESTAMP_LTZ '2020-11-01 01:30:00.123456789-08:00' in America/Los_Angeles with spark.sql.timestampNanosTypes.enabled, and a nanosecond TIMESTAMP_NTZ row was added too.

test("the pairs with no clause form are not emitted") {
val keys = Seq(key(id))
Seq(
(WriteDistributionMode.HASH, keys, Seq.empty[Transform]),

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Thanks. The suite now replays the same statement as SHOW CREATE TABLE (replayStatement, shared with ShowCreateTableExec). The no-clause-form test uses a replay that accepts any clauses, so only the guards decide, and it covers a null mode. A new test emits each mode's clause form, including DISTRIBUTED BY PARTITION ORDERED BY (id ASC NULLS FIRST) and LOCALLY ORDERED BY .... Another omits an unspellable key under HASH, RANGE and NONE.

throw QueryCompilationErrors.specifyPartitionNotAllowedWhenTableSchemaNotDefinedError()
}
if (create.writeOrdering.nonEmpty) {
throw QueryCompilationErrors

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done. The error now takes create.origin, and the test pins the statement as its context.

* <p>
* A plain column is a {@link org.apache.spark.sql.connector.expressions.NamedReference}; any
* other key is a {@link Transform}, such as {@code bucket(16, id)}. Spark checks that each
* referenced column exists in the table schema, and the parser checks the arguments of

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Thanks. It now says that the parser checks the argument kinds of the lowercase bucket, years, months, days and hours, and that Spark does not check orderability, a positive bucket count, a bucket's column list, or the arguments of any other transform.

[ COMMENT table_comment ]
[ TBLPROPERTIES ( key1=val1, key2=val2, ... ) ]
[ DISTRIBUTED BY PARTITION ]
[ [ LOCALLY ] ORDERED BY ( write_order_field [ , ... ] ) | UNORDERED ]

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done: [ LOCALLY ] ORDERED BY { ( write_order_field [ , ... ] ) | write_order_field [ , ... ] }.

@dongjoon-hyun dongjoon-hyun 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.

Thank you for addressing points 31-38, @anuragmantri. I reviewed the latest commit, 88882a2, and left 3 inline comments, continuing the numbering from points 1-39. None of them is a correctness issue. Point 39 is still open, so I repeated it here as 43. Here is a summary, ordered by file.

Cleanup

  1. WriteDistributionAndOrdering L194: the reason this comment gives for DescribingSQLBuilder no longer holds on master. SPARK-59983 (#59241) already fixed the GeneralScalarExpression recursion in ToStringSQLBuilder.visitUnexpectedExpr, so SPARK-60075 looks like a duplicate.
  2. ShowCreateTableExec L161: createsV1Table copies ResolveSessionCatalog.supportsV1Command and isV2Provider inline, so SHOW CREATE TABLE can drift from what CREATE TABLE actually does.
  3. DescribeTableExec L234: table.writeDistributionMode() and table.writeOrdering() are each called up to three times. Reading them once would avoid repeated connector work and inconsistent snapshots.

Docs

  1. PR description (point 39, not in 88882a2): the description is a separate edit from the commits, and it still says the following.
    • It describes writeDistributionMode as a String with DISTRIBUTION_MODE_HASH/_RANGE/_NONE values, and design decision 1 still says "A String mode", while the API is the WriteDistributionMode enum.
    • It says CreateTableWriteOrderSuite has 30 tests, but it now has 47.
    • It does not mention the two new abstract members of V2CreateTablePlan (writeOrdering, withWriteOrdering), the parse-back check in SHOW CREATE TABLE, or that DelegatingCatalogExtension no longer forwards the new capability.

The most important is 43, because dev/merge_spark_pr.py puts this text into the commit message.

e.children().map(toSQL(_, quote)).mkString(s"$name(", ", ", ")")
}

// `ToStringSQLBuilder` renders a few expressions through their children's `describe`, and for an

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.

[40] The recursion this comment describes is already fixed on master.

SPARK-59983 (#59241, 5926a10) changed ToStringSQLBuilder.visitUnexpectedExpr to render a GeneralScalarExpression with an unknown name as a function call. It went into master after this branch's base, so "recurses without end for a GeneralScalarExpression name it does not know" no longer holds after a rebase. SPARK-60075, which you filed for this in your last reply, looks like a duplicate of SPARK-59983 and can probably be closed.

DescribingSQLBuilder itself is still useful, because it routes literals and references through toSQL and renders other unknown expressions from their children. Could we rebase and reword this comment to give only those reasons?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done. After rebasing on master, the comment gives only the reasons that still hold. DescribingSQLBuilder renders literals and references with toSQL, so they keep their type and quoting, and renders the children of GetArrayItem and VariantGet itself rather than through their toString. It renders an expression ToStringSQLBuilder does not know, and a PartitionPredicate, from its children. The last part is still needed on master: a connector PartitionPredicate that does not override describe still recurses. I rescoped SPARK-60075 to cover this.


// In the session catalog, a CREATE TABLE whose provider is not a v2 source creates a v1 table,
// which cannot record the clauses, so the replay would be rejected.
private def createsV1Table(resolvedTable: ResolvedTable): Boolean = {

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.

[41] createsV1Table copies ResolveSessionCatalog's v1-fallback rule inline.

v1Capable is ResolveSessionCatalog.supportsV1Command (L1030), and the provider check is isV2Provider (L950). Both are private there, so this is a second copy that has to change whenever the rule changes, such as the provider/serde resolution in getStorageFormatAndProvider (L765) or the v1 source list. If the copies drift, SHOW CREATE TABLE either emits clauses that the replayed CREATE TABLE rejects with UNSUPPORTED_FEATURE.TABLE_OPERATION, or drops clauses that would have worked, and the parse-back check cannot catch this because it only parses. Could we move these two predicates into a shared helper, e.g. next to DataSourceV2Utils.getTableProvider, and call it from both places?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Agreed. The copy had in fact already drifted. For a table with no provider it used spark.sql.sources.default, while CREATE TABLE creates a Hive serde table when spark.sql.legacy.createHiveTableByDefault is set. supportsV1Command, isV2Provider, createsHiveTableByDefault and createTableProvider now live in DataSourceV2Utils, and both ResolveSessionCatalog and ShowCreateTableExec call them; ResolveSessionCatalog's behavior is unchanged. A new test covers a provider-less table in a session catalog extension with spark.sql.sources.default set to a v2 source: the clauses are emitted with the legacy conf off and omitted with it on.

* Reports the table's declared write distribution and ordering, whether or not SHOW CREATE TABLE
* can reproduce them.
*/
private def addTableWriteDistributionAndOrdering(rows: ArrayBuffer[InternalRow]): Unit = {

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.

[42] Each accessor is called up to three times.

table.writeDistributionMode() and table.writeOrdering() are each called once for isRequested, once for the null/empty check and once for the row. A connector may build these on every call, for example by converting its own sort order into a new SortOrder[] from refreshed metadata. That repeats the work, and the rows can come from different snapshots than the isRequested check. Reading both into local vals once at the top would avoid both. ShowCreateTableExec.showTableWriteDistributionAndOrdering already does this for writeOrdering.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done. addTableWriteDistributionAndOrdering now reads both into vals once. I also checked the rest of the diff for repeated connector accessor calls; ShowCreateTableExec already reads each one once.

@anuragmantri
anuragmantri force-pushed the SPARK-34586-write-distribution-ordering-create-table branch from 88882a2 to 8338ca9 Compare October 10, 2026 20:20

@anuragmantri anuragmantri left a comment

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Thank you, @dongjoon-hyun. Points 40-42 are addressed in 8338ca9, with replies on each thread. I rebased the branch on master to pick up SPARK-59983.

On SPARK-60075: it is not a duplicate. After SPARK-59983, a connector PartitionPredicate that overrides neither toString nor describe still recurses, because visitPartitionPredicate returns describe(), which defaults to toString(), which builds through the same builder again. I re-scoped SPARK-60075 to that case.

e.children().map(toSQL(_, quote)).mkString(s"$name(", ", ", ")")
}

// `ToStringSQLBuilder` renders a few expressions through their children's `describe`, and for an

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done. After rebasing on master, the comment gives only the reasons that still hold. DescribingSQLBuilder renders literals and references with toSQL, so they keep their type and quoting, and renders the children of GetArrayItem and VariantGet itself rather than through their toString. It renders an expression ToStringSQLBuilder does not know, and a PartitionPredicate, from its children. The last part is still needed on master: a connector PartitionPredicate that does not override describe still recurses. I rescoped SPARK-60075 to cover this.


// In the session catalog, a CREATE TABLE whose provider is not a v2 source creates a v1 table,
// which cannot record the clauses, so the replay would be rejected.
private def createsV1Table(resolvedTable: ResolvedTable): Boolean = {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Agreed. The copy had in fact already drifted. For a table with no provider it used spark.sql.sources.default, while CREATE TABLE creates a Hive serde table when spark.sql.legacy.createHiveTableByDefault is set. supportsV1Command, isV2Provider, createsHiveTableByDefault and createTableProvider now live in DataSourceV2Utils, and both ResolveSessionCatalog and ShowCreateTableExec call them; ResolveSessionCatalog's behavior is unchanged. A new test covers a provider-less table in a session catalog extension with spark.sql.sources.default set to a v2 source: the clauses are emitted with the legacy conf off and omitted with it on.

* Reports the table's declared write distribution and ordering, whether or not SHOW CREATE TABLE
* can reproduce them.
*/
private def addTableWriteDistributionAndOrdering(rows: ArrayBuffer[InternalRow]): Unit = {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done. addTableWriteDistributionAndOrdering now reads both into vals once. I also checked the rest of the diff for repeated connector accessor calls; ShowCreateTableExec already reads each one once.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants