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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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 =
Expand Down Expand Up @@ -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()));
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Comment thread
nielspardon marked this conversation as resolved.
return new CreateTable(namedWrite.getNames(), toRowType(namedWrite.getTableSchema()), relNode);
return new CreateTable(
namedWrite.getNames(),
toRowType(namedWrite.getTableSchema()),
relNode,
namedWrite.getCreateMode());
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
Original file line number Diff line number Diff line change
@@ -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;
Expand All @@ -13,38 +15,72 @@ public class CreateTable extends SingleRel {

private final List<String> tableName;
private final RelDataType tableSchema;
private final CreateMode createMode;

private CreateTable(
RelOptCluster cluster,
RelTraitSet traitSet,
List<String> 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");
Comment thread
nielspardon marked this conversation as resolved.
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<String> 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<String> 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<String> tableName, RelDataType tableSchema, RelNode input, CreateMode createMode) {
this(input.getCluster(), input.getTraitSet(), tableName, tableSchema, input, createMode);
}

/**
Expand All @@ -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());
}

/**
Expand All @@ -85,7 +122,8 @@ public RelNode copy(RelTraitSet traitSet, List<RelNode> 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);
}

/**
Expand All @@ -107,4 +145,13 @@ public List<String> 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;
}
}
Original file line number Diff line number Diff line change
@@ -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;
Expand Down Expand Up @@ -90,20 +91,34 @@ protected RelRoot handleNonDdl(final SqlNode sqlNode) {
/**
* Handles {@code CREATE TABLE AS SELECT} statements.
*
* <p>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) {
Comment thread
nielspardon marked this conversation as resolved.
throw new IllegalArgumentException(
Comment thread
nielspardon marked this conversation as resolved.
"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());
}

Expand Down
73 changes: 73 additions & 0 deletions isthmus/src/test/java/io/substrait/isthmus/DdlRoundtripTest.java
Original file line number Diff line number Diff line change
Expand Up @@ -2,15 +2,18 @@

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;

import io.substrait.expression.ExpressionCreator;
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;
Expand All @@ -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 =
Expand All @@ -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";
Expand Down
Loading
Loading