diff --git a/isthmus-cli/src/main/java/io/substrait/isthmus/cli/IsthmusExecutionExceptionHandler.java b/isthmus-cli/src/main/java/io/substrait/isthmus/cli/IsthmusExecutionExceptionHandler.java index 0f717a1ef..67696fa59 100644 --- a/isthmus-cli/src/main/java/io/substrait/isthmus/cli/IsthmusExecutionExceptionHandler.java +++ b/isthmus-cli/src/main/java/io/substrait/isthmus/cli/IsthmusExecutionExceptionHandler.java @@ -47,6 +47,10 @@ class IsthmusExecutionExceptionHandler implements CommandLine.IExecutionExceptio /** The message the DDL converter reports for a CREATE TABLE without a query. */ private static final String CTAS_ONLY = "Only create table as select statements are supported"; + /** The message for SQL flags that request conflicting table creation policies. */ + private static final String CONFLICTING_CTAS_POLICIES = + "CREATE TABLE cannot combine OR REPLACE and IF NOT EXISTS"; + // The hints are hard-wrapped for a terminal rather than joined into single long lines. private static final String CREATE_HINT = @@ -138,7 +142,10 @@ private static boolean isInputError(final Exception ex) { return ((SqlParseException) ex).getPos() != null; } // CalciteContextException is, by construction, a complaint about the SQL at a line and column. - return ex instanceof CalciteContextException || isPlainCreateTableQuery(ex); + return ex instanceof CalciteContextException + || isPlainCreateTableQuery(ex) + || (ex instanceof IllegalArgumentException + && CONFLICTING_CTAS_POLICIES.equals(ex.getMessage())); } /** diff --git a/isthmus-cli/src/test/java/io/substrait/isthmus/cli/IsthmusEntryPointTest.java b/isthmus-cli/src/test/java/io/substrait/isthmus/cli/IsthmusEntryPointTest.java index d27849b14..55d018024 100644 --- a/isthmus-cli/src/test/java/io/substrait/isthmus/cli/IsthmusEntryPointTest.java +++ b/isthmus-cli/src/test/java/io/substrait/isthmus/cli/IsthmusEntryPointTest.java @@ -163,6 +163,24 @@ void createTableAsQuerySuggestsCreateOption() { run.assertNoStackTrace(); } + @Test + void conflictingCreationPoliciesAreReportedWithoutAStackTrace() { + Run run = run("CREATE OR REPLACE TABLE IF NOT EXISTS dst AS SELECT 99 AS v"); + + assertEquals(CommandLine.ExitCode.SOFTWARE, run.statusCode); + run.assertErrContains("CREATE TABLE cannot combine OR REPLACE and IF NOT EXISTS"); + run.assertNoStackTrace(); + } + + @Test + void conflictingCreationPoliciesRetainTheRequestedStackTrace() { + Run run = run("CREATE OR REPLACE TABLE IF NOT EXISTS dst AS SELECT 99 AS v", "--stacktrace"); + + assertEquals(CommandLine.ExitCode.SOFTWARE, run.statusCode); + run.assertErrContains("CREATE TABLE cannot combine OR REPLACE and IF NOT EXISTS"); + run.assertStackTrace(); + } + @Test void queryPassedToCreateOptionSuggestsQueryArgument() { Run run = run("SELECT * FROM foo", "-c", "SELECT 1"); diff --git a/isthmus/src/main/java/io/substrait/isthmus/SubstraitRelNodeConverter.java b/isthmus/src/main/java/io/substrait/isthmus/SubstraitRelNodeConverter.java index fd2332e3b..62bfcf22c 100644 --- a/isthmus/src/main/java/io/substrait/isthmus/SubstraitRelNodeConverter.java +++ b/isthmus/src/main/java/io/substrait/isthmus/SubstraitRelNodeConverter.java @@ -1196,20 +1196,29 @@ private RelDataType toRowType(NamedStruct schema) { } private RelNode handleCreateTableAs(NamedWrite namedWrite, Context context) { - if (namedWrite.getCreateMode() != AbstractWriteRel.CreateMode.REPLACE_IF_EXISTS - || namedWrite.getOutputMode() != AbstractWriteRel.OutputMode.NO_OUTPUT) { + if (namedWrite.getOutputMode() != AbstractWriteRel.OutputMode.NO_OUTPUT) { throw new UnsupportedOperationException( String.format( - "Can only handle CTAS NamedWrite with (%s, %s), given (%s, %s)", - AbstractWriteRel.CreateMode.REPLACE_IF_EXISTS, - AbstractWriteRel.OutputMode.NO_OUTPUT, - namedWrite.getCreateMode(), - namedWrite.getOutputMode())); + "Can only handle CTAS NamedWrite with output mode %s, given %s", + AbstractWriteRel.OutputMode.NO_OUTPUT, namedWrite.getOutputMode())); + } + switch (namedWrite.getCreateMode()) { + case ERROR_IF_EXISTS: + case IGNORE_IF_EXISTS: + case REPLACE_IF_EXISTS: + break; + default: + throw new UnsupportedOperationException( + "Cannot convert CTAS creation mode to Calcite: " + namedWrite.getCreateMode()); } Rel input = namedWrite.getInput(); RelNode relNode = input.accept(this, context); - return new CreateTable(namedWrite.getNames(), toRowType(namedWrite.getTableSchema()), relNode); + return new CreateTable( + namedWrite.getNames(), + toRowType(namedWrite.getTableSchema()), + relNode, + namedWrite.getCreateMode()); } @Override diff --git a/isthmus/src/main/java/io/substrait/isthmus/SubstraitRelVisitor.java b/isthmus/src/main/java/io/substrait/isthmus/SubstraitRelVisitor.java index 757788ea6..d1abe24dd 100644 --- a/isthmus/src/main/java/io/substrait/isthmus/SubstraitRelVisitor.java +++ b/isthmus/src/main/java/io/substrait/isthmus/SubstraitRelVisitor.java @@ -927,7 +927,7 @@ public Rel handleCreateTable(CreateTable createTable) { .input(inputRel) .tableSchema(schema) .operation(AbstractWriteRel.WriteOp.CTAS) - .createMode(AbstractWriteRel.CreateMode.REPLACE_IF_EXISTS) + .createMode(createTable.getCreateMode()) .outputMode(AbstractWriteRel.OutputMode.NO_OUTPUT) .names(createTable.getTableName()) .build(); diff --git a/isthmus/src/main/java/io/substrait/isthmus/calcite/rel/CreateTable.java b/isthmus/src/main/java/io/substrait/isthmus/calcite/rel/CreateTable.java index 6e536417f..68a8e1bb9 100644 --- a/isthmus/src/main/java/io/substrait/isthmus/calcite/rel/CreateTable.java +++ b/isthmus/src/main/java/io/substrait/isthmus/calcite/rel/CreateTable.java @@ -1,6 +1,8 @@ package io.substrait.isthmus.calcite.rel; +import io.substrait.relation.AbstractWriteRel.CreateMode; import java.util.List; +import java.util.Objects; import org.apache.calcite.plan.RelOptCluster; import org.apache.calcite.plan.RelTraitSet; import org.apache.calcite.rel.RelNode; @@ -13,38 +15,72 @@ public class CreateTable extends SingleRel { private final List tableName; private final RelDataType tableSchema; + private final CreateMode createMode; private CreateTable( RelOptCluster cluster, RelTraitSet traitSet, List tableName, RelDataType tableSchema, - RelNode input) { + RelNode input, + CreateMode createMode) { super(cluster, traitSet, input); this.tableName = tableName; this.tableSchema = DdlSchemas.requireFilledBy(tableSchema, input, "table"); + this.createMode = Objects.requireNonNull(createMode, "createMode"); + switch (createMode) { + case ERROR_IF_EXISTS: + case IGNORE_IF_EXISTS: + case REPLACE_IF_EXISTS: + break; + default: + throw new IllegalArgumentException("Unsupported CTAS creation mode: " + createMode); + } } /** * CreateTable Constructor, taking the row type of the input as the schema of the table to create. + * Retains the historical replace-if-exists behavior; use the overload with an explicit mode to + * choose another policy. * * @param tableName tablename components * @param input RelNode input + * @deprecated Use {@link #CreateTable(List, RelDataType, RelNode, CreateMode)} to choose the + * creation policy explicitly. */ + @Deprecated public CreateTable(List tableName, RelNode input) { - this(input.getCluster(), input.getTraitSet(), tableName, input.getRowType(), input); + this(tableName, input.getRowType(), input, CreateMode.REPLACE_IF_EXISTS); } /** - * CreateTable Constructor. + * CreateTable Constructor. Retains the historical replace-if-exists behavior; use the overload + * with an explicit mode to choose another policy. * * @param tableName tablename components * @param tableSchema the schema of the table to create, which the input fills but need not name * the same way * @param input RelNode input + * @deprecated Use {@link #CreateTable(List, RelDataType, RelNode, CreateMode)} to choose the + * creation policy explicitly. */ + @Deprecated public CreateTable(List tableName, RelDataType tableSchema, RelNode input) { - this(input.getCluster(), input.getTraitSet(), tableName, tableSchema, input); + this(tableName, tableSchema, input, CreateMode.REPLACE_IF_EXISTS); + } + + /** + * Creates a table with a declared schema and an explicit policy for an existing table. + * + * @param tableName table name components + * @param tableSchema the declared schema of the table + * @param input the query filling the table + * @param createMode ERROR_IF_EXISTS, IGNORE_IF_EXISTS, or REPLACE_IF_EXISTS + * @throws IllegalArgumentException if the creation mode is unsupported by Isthmus + */ + public CreateTable( + List tableName, RelDataType tableSchema, RelNode input, CreateMode createMode) { + this(input.getCluster(), input.getTraitSet(), tableName, tableSchema, input, createMode); } /** @@ -68,7 +104,8 @@ protected RelDataType deriveRowType() { public RelWriter explainTerms(RelWriter pw) { return super.explainTerms(pw) .item("tableName", getTableName()) - .item("tableSchema", getTableSchema().getFullTypeString()); + .item("tableSchema", getTableSchema().getFullTypeString()) + .item("createMode", getCreateMode().name()); } /** @@ -85,7 +122,8 @@ public RelNode copy(RelTraitSet traitSet, List inputs) { throw new IllegalArgumentException( "CreateTable requires exactly one input, but got " + inputs.size()); } - return new CreateTable(getCluster(), traitSet, tableName, tableSchema, inputs.get(0)); + return new CreateTable( + getCluster(), traitSet, tableName, tableSchema, inputs.get(0), createMode); } /** @@ -107,4 +145,13 @@ public List getTableName() { public RelDataType getTableSchema() { return tableSchema; } + + /** + * Returns the policy to apply when the target table already exists. + * + * @return the creation mode + */ + public CreateMode getCreateMode() { + return createMode; + } } diff --git a/isthmus/src/main/java/io/substrait/isthmus/calcite/rel/DdlSqlToRelConverter.java b/isthmus/src/main/java/io/substrait/isthmus/calcite/rel/DdlSqlToRelConverter.java index 60e6de30d..e7ca6f86e 100644 --- a/isthmus/src/main/java/io/substrait/isthmus/calcite/rel/DdlSqlToRelConverter.java +++ b/isthmus/src/main/java/io/substrait/isthmus/calcite/rel/DdlSqlToRelConverter.java @@ -1,5 +1,6 @@ package io.substrait.isthmus.calcite.rel; +import io.substrait.relation.AbstractWriteRel.CreateMode; import java.util.ArrayList; import java.util.List; import java.util.Map; @@ -90,20 +91,34 @@ protected RelRoot handleNonDdl(final SqlNode sqlNode) { /** * Handles {@code CREATE TABLE AS SELECT} statements. * + *

Isthmus rejects the combination of OR REPLACE and IF NOT EXISTS rather than choosing which + * policy takes precedence. Substrait does not define how these SQL flags interact. + * * @param sqlCreateTable the CREATE TABLE node * @return a {@link RelRoot} wrapping a synthetic {@code CreateTable} relational node - * @throws IllegalArgumentException if the statement is not a CTAS + * @throws IllegalArgumentException if the statement is not a CTAS or combines OR REPLACE and IF + * NOT EXISTS */ protected RelRoot handleCreateTable(final SqlCreateTable sqlCreateTable) { if (sqlCreateTable.query == null) { throw new IllegalArgumentException("Only create table as select statements are supported"); } + if (sqlCreateTable.getReplace() && sqlCreateTable.ifNotExists) { + throw new IllegalArgumentException( + "CREATE TABLE cannot combine OR REPLACE and IF NOT EXISTS"); + } + final CreateMode createMode = + sqlCreateTable.getReplace() + ? CreateMode.REPLACE_IF_EXISTS + : sqlCreateTable.ifNotExists ? CreateMode.IGNORE_IF_EXISTS : CreateMode.ERROR_IF_EXISTS; final RelNode input = converter.convertQuery(sqlCreateTable.query, true, true).rel; final RelDataType schema = declaredSchema(sqlCreateTable.columnList, input); return RelRoot.of( - schema == null - ? new CreateTable(sqlCreateTable.name.names, input) - : new CreateTable(sqlCreateTable.name.names, schema, input), + new CreateTable( + sqlCreateTable.name.names, + schema == null ? input.getRowType() : schema, + input, + createMode), sqlCreateTable.getKind()); } diff --git a/isthmus/src/test/java/io/substrait/isthmus/DdlRoundtripTest.java b/isthmus/src/test/java/io/substrait/isthmus/DdlRoundtripTest.java index 5b8dc3b87..97ae80be1 100644 --- a/isthmus/src/test/java/io/substrait/isthmus/DdlRoundtripTest.java +++ b/isthmus/src/test/java/io/substrait/isthmus/DdlRoundtripTest.java @@ -2,6 +2,7 @@ import static org.junit.jupiter.api.Assertions.assertAll; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertInstanceOf; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; @@ -9,8 +10,10 @@ import io.substrait.isthmus.sql.SubstraitCreateStatementParser; import io.substrait.isthmus.sql.SubstraitSqlToCalcite; import io.substrait.plan.Plan; +import io.substrait.proto.WriteRel; import io.substrait.relation.AbstractDdlRel; import io.substrait.relation.AbstractWriteRel; +import io.substrait.relation.AbstractWriteRel.CreateMode; import io.substrait.relation.NamedDdl; import io.substrait.relation.NamedWrite; import io.substrait.relation.Project; @@ -22,6 +25,10 @@ import org.apache.calcite.rel.RelRoot; import org.apache.calcite.sql.parser.SqlParseException; import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.CsvSource; +import org.junit.jupiter.params.provider.EnumSource; +import org.junit.jupiter.params.provider.ValueSource; class DdlRoundtripTest extends PlanTestBase { final Prepare.CatalogReader catalogReader = @@ -39,6 +46,72 @@ void testCreateTable() throws Exception { assertFullRoundTrip(sql, catalogReader); } + @ParameterizedTest + @CsvSource({ + "CREATE TABLE dst AS SELECT 99 AS v, CREATE_MODE_ERROR_IF_EXISTS", + "CREATE TABLE IF NOT EXISTS dst AS SELECT 99 AS v, CREATE_MODE_IGNORE_IF_EXISTS", + "CREATE OR REPLACE TABLE dst AS SELECT 99 AS v, CREATE_MODE_REPLACE_IF_EXISTS", + "CREATE TABLE dst(v INTEGER) AS SELECT 99 AS v, CREATE_MODE_ERROR_IF_EXISTS", + "CREATE TABLE IF NOT EXISTS dst(v INTEGER) AS SELECT 99 AS v, CREATE_MODE_IGNORE_IF_EXISTS", + "CREATE OR REPLACE TABLE dst(v INTEGER) AS SELECT 99 AS v, CREATE_MODE_REPLACE_IF_EXISTS" + }) + void preservesCreationPolicyWithOrWithoutAnExistingTarget( + String sql, WriteRel.CreateMode expected) throws SqlParseException { + for (Prepare.CatalogReader catalog : + List.of( + SubstraitCreateStatementParser.processCreateStatementsToCatalog(), + SubstraitCreateStatementParser.processCreateStatementsToCatalog( + "CREATE TABLE dst(v INTEGER)"))) { + assertFullRoundTrip(sql, catalog); + Plan plan = new SqlToSubstrait(converterProvider).convert(sql, catalog); + assertEquals( + expected, toProto(plan).getRelations(0).getRoot().getInput().getWrite().getCreateMode()); + } + } + + @Test + void rejectsConflictingCreationPolicies() { + assertEquals( + "CREATE TABLE cannot combine OR REPLACE and IF NOT EXISTS", + assertThrows( + IllegalArgumentException.class, + () -> + new SqlToSubstrait(converterProvider) + .convert( + "CREATE OR REPLACE TABLE IF NOT EXISTS dst AS SELECT 99 AS v", + catalogReader)) + .getMessage()); + } + + @ParameterizedTest + @ValueSource(strings = {"CREATE TABLE", "CREATE TABLE IF NOT EXISTS", "CREATE OR REPLACE TABLE"}) + void createWithoutAQueryRemainsUnsupported(String prefix) { + assertThrows( + IllegalArgumentException.class, + () -> + new SqlToSubstrait(converterProvider) + .convert(prefix + " dst(v INTEGER)", catalogReader)); + } + + @ParameterizedTest + @EnumSource( + value = CreateMode.class, + names = {"ERROR_IF_EXISTS", "IGNORE_IF_EXISTS", "REPLACE_IF_EXISTS"}, + mode = EnumSource.Mode.EXCLUDE) + void unsupportedCreationModesAreRefused(CreateMode mode) throws SqlParseException { + Plan plan = + new SqlToSubstrait(converterProvider) + .convert("CREATE TABLE dst AS SELECT 99 AS v", catalogReader); + NamedWrite write = assertInstanceOf(NamedWrite.class, plan.getRoots().get(0).getInput()); + NamedWrite unsupported = NamedWrite.builder().from(write).createMode(mode).build(); + assertEquals( + "Cannot convert CTAS creation mode to Calcite: " + mode, + assertThrows( + UnsupportedOperationException.class, + () -> new SubstraitToCalcite(converterProvider, catalogReader).convert(unsupported)) + .getMessage()); + } + @Test void testCreateView() throws Exception { String sql = "create view dst1 as select * from src1"; diff --git a/isthmus/src/test/java/io/substrait/isthmus/calcite/rel/DdlRelCopyTest.java b/isthmus/src/test/java/io/substrait/isthmus/calcite/rel/DdlRelCopyTest.java index d3a20d43b..dc41bb49b 100644 --- a/isthmus/src/test/java/io/substrait/isthmus/calcite/rel/DdlRelCopyTest.java +++ b/isthmus/src/test/java/io/substrait/isthmus/calcite/rel/DdlRelCopyTest.java @@ -1,14 +1,21 @@ package io.substrait.isthmus.calcite.rel; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertInstanceOf; +import static org.junit.jupiter.api.Assertions.assertNotEquals; import static org.junit.jupiter.api.Assertions.assertThrows; +import com.fasterxml.jackson.databind.ObjectMapper; import io.substrait.isthmus.PlanTestBase; +import io.substrait.relation.AbstractWriteRel.CreateMode; import java.util.List; import org.apache.calcite.rel.RelNode; +import org.apache.calcite.rel.externalize.RelJsonWriter; import org.apache.calcite.rel.type.RelDataType; import org.apache.calcite.sql.type.SqlTypeName; import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.EnumSource; /** * The DDL relations are single-input, their {@code copy()} rejects any other input count, and the @@ -24,6 +31,65 @@ private RelDataType declaredSchema() { List.of(typeFactory.createSqlType(SqlTypeName.BIGINT)), List.of("declared")); } + @ParameterizedTest + @EnumSource( + value = CreateMode.class, + names = {"ERROR_IF_EXISTS", "IGNORE_IF_EXISTS", "REPLACE_IF_EXISTS"}) + void creationPolicySurvivesPlannerCopies(CreateMode mode) { + CreateTable original = new CreateTable(List.of("DST"), declaredSchema(), input, mode); + CreateTable copied = + assertInstanceOf( + CreateTable.class, original.copy(original.getTraitSet(), List.of(otherInput))); + assertEquals(mode, copied.getCreateMode()); + assertEquals(declaredSchema(), copied.getTableSchema()); + assertEquals(otherInput, copied.getInput()); + } + + @ParameterizedTest + @EnumSource( + value = CreateMode.class, + names = {"ERROR_IF_EXISTS", "IGNORE_IF_EXISTS", "REPLACE_IF_EXISTS"}) + void creationPolicyCanBeExplainedAsJson(CreateMode mode) throws Exception { + CreateTable table = new CreateTable(List.of("DST"), declaredSchema(), input, mode); + RelJsonWriter writer = new RelJsonWriter(); + table.explain(writer); + assertEquals( + List.of(mode.name()), + new ObjectMapper().readTree(writer.asString()).findValuesAsText("createMode")); + } + + @ParameterizedTest + @EnumSource( + value = CreateMode.class, + names = {"ERROR_IF_EXISTS", "IGNORE_IF_EXISTS", "REPLACE_IF_EXISTS"}, + mode = EnumSource.Mode.EXCLUDE) + void constructorsRejectUnsupportedCreationPolicies(CreateMode mode) { + assertEquals( + "Unsupported CTAS creation mode: " + mode, + assertThrows( + IllegalArgumentException.class, + () -> new CreateTable(List.of("DST"), declaredSchema(), input, mode)) + .getMessage()); + } + + @Test + void existingConstructorsRetainTheirCreationPolicy() { + assertEquals( + CreateMode.REPLACE_IF_EXISTS, new CreateTable(List.of("DST"), input).getCreateMode()); + assertEquals( + CreateMode.REPLACE_IF_EXISTS, + new CreateTable(List.of("DST"), declaredSchema(), input).getCreateMode()); + } + + @Test + void differentCreationPoliciesHaveDifferentPlannerDigests() { + CreateTable plain = + new CreateTable(List.of("DST"), declaredSchema(), input, CreateMode.ERROR_IF_EXISTS); + CreateTable replace = + new CreateTable(List.of("DST"), declaredSchema(), input, CreateMode.REPLACE_IF_EXISTS); + assertNotEquals(plain.getDigest(), replace.getDigest()); + } + @Test void createTableCopiesItsSingleInput() { CreateTable createTable = new CreateTable(List.of("FOO"), input);