diff --git a/docs/docs/spark-ddl.md b/docs/docs/spark-ddl.md index 9fa6c0e7d3c7..fd2ca04cf6d5 100644 --- a/docs/docs/spark-ddl.md +++ b/docs/docs/spark-ddl.md @@ -44,7 +44,16 @@ Table create commands, including CTAS and RTAS, support the full range of Spark Create commands may also set the default format with the `USING` clause. This is only supported for `SparkCatalog` because Spark handles the `USING` clause differently for the built-in catalog. -`CREATE TABLE ... LIKE ...` syntax is not supported. +Spark 4.2 and later can create an Iceberg table from an existing table: + +```sql +CREATE TABLE prod.db.sample_copy LIKE prod.db.sample; +``` + +`CREATE TABLE ... LIKE ...` copies the source schema, partitioning, sort order, and +table properties. The new table does not copy snapshots, data, metadata history, +or the source table location. Unless `LOCATION` is specified for the new table, +the target catalog assigns its default location. ### `PARTITIONED BY` diff --git a/spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/SparkCatalog.java b/spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/SparkCatalog.java index 67c75c3a63d1..66f853d7c3ba 100644 --- a/spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/SparkCatalog.java +++ b/spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/SparkCatalog.java @@ -22,6 +22,7 @@ import static org.apache.iceberg.TableProperties.GC_ENABLED_DEFAULT; import java.util.Arrays; +import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.Set; @@ -36,6 +37,8 @@ import org.apache.iceberg.HasTableOperations; import org.apache.iceberg.MetadataTableType; import org.apache.iceberg.Schema; +import org.apache.iceberg.SortField; +import org.apache.iceberg.SortOrder; import org.apache.iceberg.Transaction; import org.apache.iceberg.catalog.Catalog; import org.apache.iceberg.catalog.Namespace; @@ -44,6 +47,7 @@ import org.apache.iceberg.catalog.ViewCatalog; import org.apache.iceberg.exceptions.AlreadyExistsException; import org.apache.iceberg.exceptions.ValidationException; +import org.apache.iceberg.expressions.Expressions; import org.apache.iceberg.hadoop.HadoopCatalog; import org.apache.iceberg.hadoop.HadoopTables; import org.apache.iceberg.relocated.com.google.common.base.Preconditions; @@ -82,6 +86,7 @@ import org.apache.spark.sql.connector.catalog.TableChange.ColumnChange; import org.apache.spark.sql.connector.catalog.TableChange.RemoveProperty; import org.apache.spark.sql.connector.catalog.TableChange.SetProperty; +import org.apache.spark.sql.connector.catalog.TableInfo; import org.apache.spark.sql.connector.catalog.TableSummary; import org.apache.spark.sql.connector.catalog.View; import org.apache.spark.sql.connector.expressions.Transform; @@ -206,11 +211,45 @@ public Table createTable( Identifier ident, StructType schema, Transform[] transforms, Map properties) throws TableAlreadyExistsException { Schema icebergSchema = SparkSchemaUtil.convert(schema); + return createTable(ident, icebergSchema, transforms, properties, SortOrder.unsorted()); + } + + @Override + public Table createTableLike(Identifier ident, TableInfo tableInfo, Table sourceTable) + throws TableAlreadyExistsException, NoSuchNamespaceException { + // Spark intentionally excludes the source table's properties from tableInfo and leaves it to + // the connector to decide which to clone via sourceTable. Clone the source Iceberg table's + // properties and sort order, then let user-specified LIKE options (in tableInfo) take + // precedence. + Schema icebergSchema = SparkSchemaUtil.convert(tableInfo.schema()); + Map properties = new HashMap<>(); + SortOrder sortOrder = SortOrder.unsorted(); + + if (sourceTable instanceof SparkTable) { + org.apache.iceberg.Table sourceIcebergTable = ((SparkTable) sourceTable).table(); + properties.putAll(sourceIcebergTable.properties()); + sortOrder = + copySortOrder(sourceIcebergTable.schema(), icebergSchema, sourceIcebergTable.sortOrder()); + } + + properties.putAll(tableInfo.properties()); + + return createTable(ident, icebergSchema, tableInfo.partitions(), properties, sortOrder); + } + + private Table createTable( + Identifier ident, + Schema icebergSchema, + Transform[] transforms, + Map properties, + SortOrder sortOrder) + throws TableAlreadyExistsException { try { Catalog.TableBuilder builder = newBuilder(ident, icebergSchema); org.apache.iceberg.Table icebergTable = builder .withPartitionSpec(Spark3Util.toPartitionSpec(icebergSchema, transforms)) + .withSortOrder(sortOrder) .withLocation(properties.get("location")) .withProperties(Spark3Util.rebuildCreateProperties(properties)) .create(); @@ -220,6 +259,24 @@ public Table createTable( } } + private static SortOrder copySortOrder( + Schema sourceSchema, Schema targetSchema, SortOrder sourceSortOrder) { + if (sourceSortOrder.isUnsorted()) { + return SortOrder.unsorted(); + } + + SortOrder.Builder builder = SortOrder.builderFor(targetSchema); + for (SortField field : sourceSortOrder.fields()) { + String sourceName = sourceSchema.findColumnName(field.sourceId()); + builder.sortBy( + Expressions.transform(sourceName, field.transform()), + field.direction(), + field.nullOrder()); + } + + return builder.build(); + } + @Override public StagedTable stageCreate( Identifier ident, StructType schema, Transform[] transforms, Map properties) diff --git a/spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/SparkSessionCatalog.java b/spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/SparkSessionCatalog.java index d754f84b276d..bc47257e1a21 100644 --- a/spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/SparkSessionCatalog.java +++ b/spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/SparkSessionCatalog.java @@ -52,6 +52,7 @@ import org.apache.spark.sql.connector.catalog.Table; import org.apache.spark.sql.connector.catalog.TableCatalog; import org.apache.spark.sql.connector.catalog.TableChange; +import org.apache.spark.sql.connector.catalog.TableInfo; import org.apache.spark.sql.connector.catalog.TableSummary; import org.apache.spark.sql.connector.catalog.View; import org.apache.spark.sql.connector.catalog.ViewCatalog; @@ -252,6 +253,19 @@ public Table createTable( } } + @Override + public Table createTableLike(Identifier ident, TableInfo tableInfo, Table sourceTable) + throws TableAlreadyExistsException, NoSuchNamespaceException { + checkViewNotExists(ident); + + String provider = tableInfo.properties().get("provider"); + if (useIceberg(provider)) { + return icebergCatalog.createTableLike(ident, tableInfo, sourceTable); + } else { + return getSessionCatalog().createTableLike(ident, tableInfo, sourceTable); + } + } + @Override public StagedTable stageCreate( Identifier ident, StructType schema, Transform[] partitions, Map properties) diff --git a/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/sql/TestCreateTable.java b/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/sql/TestCreateTable.java index afb4e9964842..6bdc2334e91b 100644 --- a/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/sql/TestCreateTable.java +++ b/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/sql/TestCreateTable.java @@ -31,9 +31,12 @@ import org.apache.iceberg.ParameterizedTestExtension; import org.apache.iceberg.PartitionSpec; import org.apache.iceberg.Schema; +import org.apache.iceberg.SortOrder; import org.apache.iceberg.Table; import org.apache.iceberg.TableOperations; import org.apache.iceberg.TableProperties; +import org.apache.iceberg.catalog.Namespace; +import org.apache.iceberg.catalog.TableIdentifier; import org.apache.iceberg.exceptions.ValidationException; import org.apache.iceberg.hadoop.HadoopCatalog; import org.apache.iceberg.spark.CatalogTestBase; @@ -69,6 +72,102 @@ public void testTransformIgnoreCase() { assertThat(validationCatalog.tableExists(tableIdent)).as("Table should already exist").isTrue(); } + @TestTemplate + public void testCreateTableLike() { + String sourceName = tableName("source"); + TableIdentifier sourceIdent = TableIdentifier.of(Namespace.of("default"), "source"); + Schema schema = + new Schema( + NestedField.required(1, "id", Types.LongType.get()), + NestedField.optional(2, "category", Types.StringType.get()), + NestedField.optional(3, "data", Types.StringType.get())); + PartitionSpec spec = PartitionSpec.builderFor(schema).identity("category").build(); + SortOrder order = SortOrder.builderFor(schema).desc("id").asc("data").build(); + + try { + validationCatalog + .buildTable(sourceIdent, schema) + .withPartitionSpec(spec) + .withSortOrder(order) + .withProperty("custom-property", "custom-value") + .create(); + + Table source = validationCatalog.loadTable(sourceIdent); + sql("CREATE TABLE %s LIKE %s", tableName, sourceName); + + Table target = validationCatalog.loadTable(tableIdent); + assertThat(target.schema().asStruct()).isEqualTo(source.schema().asStruct()); + assertThat(target.spec()).isEqualTo(source.spec()); + assertThat(target.sortOrder().sameOrder(source.sortOrder())).isTrue(); + assertThat(target.properties()).containsEntry("custom-property", "custom-value"); + assertThat(target.location()).isNotEqualTo(source.location()); + } finally { + sql("DROP TABLE IF EXISTS %s", sourceName); + } + } + + @TestTemplate + public void testCreateTableLikeClonesAndOverridesProperties() { + String sourceName = tableName("source"); + TableIdentifier sourceIdent = TableIdentifier.of(Namespace.of("default"), "source"); + Schema schema = new Schema(NestedField.required(1, "id", Types.LongType.get())); + + try { + validationCatalog + .buildTable(sourceIdent, schema) + .withProperty("clone-me", "from-source") + .withProperty("override-me", "from-source") + .create(); + + sql( + "CREATE TABLE %s LIKE %s TBLPROPERTIES ('override-me'='from-target')", + tableName, sourceName); + + Table target = validationCatalog.loadTable(tableIdent); + assertThat(target.properties()) + .containsEntry("clone-me", "from-source") + .containsEntry("override-me", "from-target"); + } finally { + sql("DROP TABLE IF EXISTS %s", sourceName); + } + } + + @TestTemplate + public void testCreateTableLikeIfNotExists() { + String sourceName = tableName("source"); + + try { + sql( + "CREATE TABLE %s (id BIGINT, data STRING) " + + "USING iceberg TBLPROPERTIES ('source-property'='source')", + sourceName); + sql( + "CREATE TABLE %s (id BIGINT) " + + "USING iceberg TBLPROPERTIES ('target-property'='target')", + tableName); + + sql("CREATE TABLE IF NOT EXISTS %s LIKE %s", tableName, sourceName); + + Table target = validationCatalog.loadTable(tableIdent); + assertThat(target.schema().columns()).hasSize(1); + assertThat(target.properties()) + .containsEntry("target-property", "target") + .doesNotContainKey("source-property"); + } finally { + sql("DROP TABLE IF EXISTS %s", sourceName); + } + } + + @TestTemplate + public void testCreateTableLikeMissingSource() { + String missingSource = tableName("missing_source"); + + assertThatThrownBy(() -> sql("CREATE TABLE %s LIKE %s", tableName, missingSource)) + .isInstanceOf(org.apache.spark.sql.AnalysisException.class) + .hasMessageContaining("missing_source"); + assertThat(validationCatalog.tableExists(tableIdent)).isFalse(); + } + @TestTemplate public void testTransformSingularForm() { assertThat(validationCatalog.tableExists(tableIdent))