szehon-ho commented on code in PR #17954:
URL: https://github.com/apache/iceberg/pull/17954#discussion_r4150877073
##########
spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/SparkCatalog.java:
##########
@@ -206,11 +213,64 @@ public Table createTable(
Identifier ident, StructType schema, Transform[] transforms, Map<String,
String> properties)
throws TableAlreadyExistsException {
Schema icebergSchema = SparkSchemaUtil.convert(schema);
+ return createTable(
+ ident,
+ icebergSchema,
+ Spark3Util.toPartitionSpec(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
+ // schema, partition spec, properties and sort order, then let
user-specified LIKE options (in
+ // tableInfo) take precedence.
+ Schema icebergSchema;
+ PartitionSpec spec;
+ Map<String, String> properties = Maps.newHashMap();
+ SortOrder sortOrder = SortOrder.unsorted();
+
+ if (sourceTable instanceof SparkTable sparkTable) {
+ org.apache.iceberg.Table sourceIcebergTable = sparkTable.table();
+ icebergSchema = sparkTable.icebergSchema();
+ spec = sourceIcebergTable.spec();
+ properties.putAll(sourceIcebergTable.properties());
+ properties.remove(TableProperties.WRITE_METADATA_LOCATION);
+ properties.remove(TableProperties.WRITE_DATA_LOCATION);
+ properties.remove(TableProperties.OBJECT_STORE_PATH);
+ properties.remove(TableProperties.WRITE_FOLDER_STORAGE_LOCATION);
+ properties.remove(TableProperties.DEFAULT_NAME_MAPPING);
+ properties.put(
+ TableProperties.FORMAT_VERSION,
+ String.valueOf(TableUtil.formatVersion(sourceIcebergTable)));
+ sortOrder =
+ copySortOrder(sourceIcebergTable.schema(), icebergSchema,
sourceIcebergTable.sortOrder());
Review Comment:
Please resolve sort field names from the pinned schema using the source
field IDs. If a sorted column was renamed from `data` to `payload` after a
snapshot, cloning that snapshot uses its old schema but tries to bind
`payload`, so creation fails with “Cannot find field”.
##########
spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/SparkCatalog.java:
##########
@@ -206,11 +213,64 @@ public Table createTable(
Identifier ident, StructType schema, Transform[] transforms, Map<String,
String> properties)
throws TableAlreadyExistsException {
Schema icebergSchema = SparkSchemaUtil.convert(schema);
+ return createTable(
+ ident,
+ icebergSchema,
+ Spark3Util.toPartitionSpec(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
+ // schema, partition spec, properties and sort order, then let
user-specified LIKE options (in
+ // tableInfo) take precedence.
+ Schema icebergSchema;
+ PartitionSpec spec;
+ Map<String, String> properties = Maps.newHashMap();
+ SortOrder sortOrder = SortOrder.unsorted();
+
+ if (sourceTable instanceof SparkTable sparkTable) {
+ org.apache.iceberg.Table sourceIcebergTable = sparkTable.table();
+ icebergSchema = sparkTable.icebergSchema();
+ spec = sourceIcebergTable.spec();
+ properties.putAll(sourceIcebergTable.properties());
+ properties.remove(TableProperties.WRITE_METADATA_LOCATION);
+ properties.remove(TableProperties.WRITE_DATA_LOCATION);
+ properties.remove(TableProperties.OBJECT_STORE_PATH);
+ properties.remove(TableProperties.WRITE_FOLDER_STORAGE_LOCATION);
+ properties.remove(TableProperties.DEFAULT_NAME_MAPPING);
+ properties.put(
+ TableProperties.FORMAT_VERSION,
+ String.valueOf(TableUtil.formatVersion(sourceIcebergTable)));
+ sortOrder =
+ copySortOrder(sourceIcebergTable.schema(), icebergSchema,
sourceIcebergTable.sortOrder());
+ } else {
+ icebergSchema = SparkSchemaUtil.convert(tableInfo.schema());
+ spec = Spark3Util.toPartitionSpec(icebergSchema, tableInfo.partitions());
+ }
+
+ properties.putAll(tableInfo.properties());
Review Comment:
Please apply the `USING` override before merging the copied properties. If
the source has `write.format.default=parquet`, `CREATE TABLE target LIKE source
USING orc` passes both that property and `provider=orc` to
`rebuildCreateProperties`. Its immutable-map builder then receives
`write.format.default` twice and throws instead of using ORC.
##########
spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/SparkCatalog.java:
##########
@@ -206,11 +213,64 @@ public Table createTable(
Identifier ident, StructType schema, Transform[] transforms, Map<String,
String> properties)
throws TableAlreadyExistsException {
Schema icebergSchema = SparkSchemaUtil.convert(schema);
+ return createTable(
+ ident,
+ icebergSchema,
+ Spark3Util.toPartitionSpec(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
+ // schema, partition spec, properties and sort order, then let
user-specified LIKE options (in
+ // tableInfo) take precedence.
+ Schema icebergSchema;
+ PartitionSpec spec;
+ Map<String, String> properties = Maps.newHashMap();
+ SortOrder sortOrder = SortOrder.unsorted();
+
+ if (sourceTable instanceof SparkTable sparkTable) {
+ org.apache.iceberg.Table sourceIcebergTable = sparkTable.table();
+ icebergSchema = sparkTable.icebergSchema();
+ spec = sourceIcebergTable.spec();
+ properties.putAll(sourceIcebergTable.properties());
+ properties.remove(TableProperties.WRITE_METADATA_LOCATION);
+ properties.remove(TableProperties.WRITE_DATA_LOCATION);
+ properties.remove(TableProperties.OBJECT_STORE_PATH);
+ properties.remove(TableProperties.WRITE_FOLDER_STORAGE_LOCATION);
+ properties.remove(TableProperties.DEFAULT_NAME_MAPPING);
+ properties.put(
+ TableProperties.FORMAT_VERSION,
+ String.valueOf(TableUtil.formatVersion(sourceIcebergTable)));
+ sortOrder =
+ copySortOrder(sourceIcebergTable.schema(), icebergSchema,
sourceIcebergTable.sortOrder());
+ } else {
+ icebergSchema = SparkSchemaUtil.convert(tableInfo.schema());
Review Comment:
Please preserve supported default values for non-Iceberg sources too.
`tableInfo.schema()` carries current/write and existence/initial defaults in
`StructField` metadata, but `SparkTypeToType` ignores them, so even a literal
`DEFAULT 42` is silently lost. Please add coverage for this case; defaults that
cannot be represented in Iceberg should be rejected explicitly.
##########
spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/SparkCatalog.java:
##########
@@ -206,11 +213,64 @@ public Table createTable(
Identifier ident, StructType schema, Transform[] transforms, Map<String,
String> properties)
throws TableAlreadyExistsException {
Schema icebergSchema = SparkSchemaUtil.convert(schema);
+ return createTable(
+ ident,
+ icebergSchema,
+ Spark3Util.toPartitionSpec(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
+ // schema, partition spec, properties and sort order, then let
user-specified LIKE options (in
+ // tableInfo) take precedence.
+ Schema icebergSchema;
+ PartitionSpec spec;
+ Map<String, String> properties = Maps.newHashMap();
+ SortOrder sortOrder = SortOrder.unsorted();
+
+ if (sourceTable instanceof SparkTable sparkTable) {
+ org.apache.iceberg.Table sourceIcebergTable = sparkTable.table();
+ icebergSchema = sparkTable.icebergSchema();
+ spec = sourceIcebergTable.spec();
Review Comment:
Please omit void partition fields whose source columns have been deleted
before creating the target. This is a valid state for v1 tables after dropping
a partition field and then its column, but `TableMetadata.newTableMetadata`
assumes each partition source still exists and throws a `NullPointerException`
while assigning fresh IDs.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]