From a2e44ca30170c53d275c49911156aa54d0c884bf Mon Sep 17 00:00:00 2001 From: Dzeri96 Date: Wed, 1 Jul 2026 10:36:54 +0200 Subject: [PATCH 1/9] Issue 14424: Deprecated default-namespace config. Now defaultDatabase is used. --- docs/docs/spark-configuration.md | 22 ++--- .../apache/iceberg/spark/SparkCatalog.java | 15 +++- .../iceberg/spark/SparkSessionCatalog.java | 3 +- .../iceberg/spark/CustomSparkTestBase.java | 83 +++++++++++++++++++ .../spark/sql/TestStaticCatalogConfigs.java | 78 +++++++++++++++++ 5 files changed, 187 insertions(+), 14 deletions(-) create mode 100644 spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java create mode 100644 spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java diff --git a/docs/docs/spark-configuration.md b/docs/docs/spark-configuration.md index 8d9dd6dd2756..9dc49b4f522a 100644 --- a/docs/docs/spark-configuration.md +++ b/docs/docs/spark-configuration.md @@ -63,22 +63,22 @@ Iceberg supplies two implementations: Both catalogs are configured using properties nested under the catalog name. Common configuration properties for Hive and Hadoop are: -| Property | Values | Description | -| -------------------------------------------------- |----------------------------------------------------|----------------------------------------------------------------------------------------------------------------------------------------------------------------------| -| spark.sql.catalog._catalog-name_.type | `hive`, `hadoop`, `rest`, `glue`, `jdbc` or `nessie` | The underlying Iceberg catalog implementation, `HiveCatalog`, `HadoopCatalog`, `RESTCatalog`, `GlueCatalog`, `JdbcCatalog`, `NessieCatalog` or left unset if using a custom catalog | -| spark.sql.catalog._catalog-name_.catalog-impl | | The custom Iceberg catalog implementation. If `type` is null, `catalog-impl` must not be null. | +| Property | Values | Description | +|---------------------------------------------------------------|----------------------------------------------------|----------------------------------------------------------------------------------------------------------------------------------------------------------------------| +| spark.sql.catalog._catalog-name_.type | `hive`, `hadoop`, `rest`, `glue`, `jdbc` or `nessie` | The underlying Iceberg catalog implementation, `HiveCatalog`, `HadoopCatalog`, `RESTCatalog`, `GlueCatalog`, `JdbcCatalog`, `NessieCatalog` or left unset if using a custom catalog | +| spark.sql.catalog._catalog-name_.catalog-impl | | The custom Iceberg catalog implementation. If `type` is null, `catalog-impl` must not be null. | | spark.sql.catalog._catalog-name_.io-impl | | The custom FileIO implementation. | | spark.sql.catalog._catalog-name_.metrics-reporter-impl | | The custom MetricsReporter implementation. | -| spark.sql.catalog._catalog-name_.default-namespace | default | The default current namespace for the catalog | -| spark.sql.catalog._catalog-name_.uri | thrift://host:port | Hive metastore URL for hive typed catalog, REST URL for REST typed catalog | -| spark.sql.catalog._catalog-name_.warehouse | hdfs://nn:8020/warehouse/path | Base path for the warehouse directory | -| spark.sql.catalog._catalog-name_.cache-enabled | `true` or `false` | Whether to enable catalog cache, default value is `true` | +| spark.sql.catalog._catalog-name_.defaultDatabase | default | The default current namespace for the catalog | +| spark.sql.catalog._catalog-name_.uri | thrift://host:port | Hive metastore URL for hive typed catalog, REST URL for REST typed catalog | +| spark.sql.catalog._catalog-name_.warehouse | hdfs://nn:8020/warehouse/path | Base path for the warehouse directory | +| spark.sql.catalog._catalog-name_.cache-enabled | `true` or `false` | Whether to enable catalog cache, default value is `true` | | spark.sql.catalog._catalog-name_.cache.expiration-interval-ms | `30000` (30 seconds) | Duration after which cached catalog entries are expired; Only effective if `cache-enabled` is `true`. `-1` disables cache expiration and `0` disables caching entirely, irrespective of `cache-enabled`. Default is `30000` (30 seconds) | | spark.sql.catalog._catalog-name_.table-default._propertyKey_ | | Default Iceberg table property value for property key _propertyKey_, which will be set on tables created by this catalog if not overridden | | spark.sql.catalog._catalog-name_.table-override._propertyKey_ | | Enforced Iceberg table property value for property key _propertyKey_, which cannot be overridden on table creation by user | -| spark.sql.catalog._catalog-name_.view-default._propertyKey_ | | Default Iceberg view property value for property key _propertyKey_, which will be set on views created by this catalog if not overridden | -| spark.sql.catalog._catalog-name_.view-override._propertyKey_ | | Enforced Iceberg view property value for property key _propertyKey_, which cannot be overridden on view creation by user | -| spark.sql.catalog._catalog-name_.use-nullable-query-schema | `true` or `false` | Whether to preserve fields' nullability when creating the table using CTAS and RTAS. If set to `true`, all fields will be marked as nullable. If set to `false`, fields' nullability will be preserved. The default value is `true`. Available in Spark 3.5 and above. | +| spark.sql.catalog._catalog-name_.view-default._propertyKey_ | | Default Iceberg view property value for property key _propertyKey_, which will be set on views created by this catalog if not overridden | +| spark.sql.catalog._catalog-name_.view-override._propertyKey_ | | Enforced Iceberg view property value for property key _propertyKey_, which cannot be overridden on view creation by user | +| spark.sql.catalog._catalog-name_.use-nullable-query-schema | `true` or `false` | Whether to preserve fields' nullability when creating the table using CTAS and RTAS. If set to `true`, all fields will be marked as nullable. If set to `false`, fields' nullability will be preserved. The default value is `true`. Available in Spark 3.5 and above. | Additional properties can be found in common [catalog configuration](catalog-properties.md). diff --git a/spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/SparkCatalog.java b/spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/SparkCatalog.java index da22607d05b0..1cfc3e3af8cc 100644 --- a/spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/SparkCatalog.java +++ b/spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/SparkCatalog.java @@ -91,6 +91,8 @@ import org.apache.spark.sql.connector.expressions.Transform; import org.apache.spark.sql.types.StructType; import org.apache.spark.sql.util.CaseInsensitiveStringMap; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; /** * A Spark TableCatalog implementation that wraps an Iceberg {@link Catalog}. @@ -106,7 +108,9 @@ *
  • io-impl - a custom {@link org.apache.iceberg.io.FileIO} implementation to use *
  • metrics-reporter-impl - a custom {@link * org.apache.iceberg.metrics.MetricsReporter} implementation to use - *
  • default-namespace - a namespace to use as the default + *
  • default-namespace - DEPRECATED: use defaultDatabase + * instead + *
  • defaultDatabase - a namespace/database to use as the default *
  • cache-enabled - whether to enable catalog cache *
  • cache.case-sensitive - whether the catalog cache should compare table * identifiers in a case sensitive way @@ -122,6 +126,8 @@ *

    */ public class SparkCatalog extends BaseCatalog { + private static final Logger LOG = LoggerFactory.getLogger(SparkCatalog.class); + private static final Set DEFAULT_NS_KEYS = ImmutableSet.of(TableCatalog.PROP_OWNER); private static final Splitter COMMA = Splitter.on(","); private static final Joiner COMMA_JOINER = Joiner.on(","); @@ -783,7 +789,14 @@ public final void initialize(String name, CaseInsensitiveStringMap options) { : catalog; if (catalog instanceof SupportsNamespaces) { this.asNamespaceCatalog = (SupportsNamespaces) catalog; + if (options.containsKey("defaultDatabase")) { + this.defaultNamespace = + Splitter.on('.').splitToList(options.get("defaultDatabase")).toArray(new String[0]); + } if (options.containsKey("default-namespace")) { + LOG.warn( + "The `default-namespace` property is deprecated and will be removed in the next major version" + + " Please use Spark's `defaultDatabase` instead."); this.defaultNamespace = Splitter.on('.').splitToList(options.get("default-namespace")).toArray(new String[0]); } diff --git a/spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/SparkSessionCatalog.java b/spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/SparkSessionCatalog.java index 1e66f10cc20c..5d09a3c7a8f4 100644 --- a/spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/SparkSessionCatalog.java +++ b/spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/SparkSessionCatalog.java @@ -64,7 +64,6 @@ public class SparkSessionCatalog< T extends TableCatalog & FunctionCatalog & SupportsNamespaces & ViewCatalog> extends BaseCatalog implements CatalogExtension { - private static final String[] DEFAULT_NAMESPACE = new String[] {"default"}; private String catalogName = null; private TableCatalog icebergCatalog = null; @@ -93,7 +92,7 @@ protected TableCatalog buildSparkCatalog(String name, CaseInsensitiveStringMap o @Override public String[] defaultNamespace() { - return DEFAULT_NAMESPACE; + return getSessionCatalog().defaultNamespace(); } @Override diff --git a/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java b/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java new file mode 100644 index 000000000000..078106685f0c --- /dev/null +++ b/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java @@ -0,0 +1,83 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.iceberg.spark; + +import static org.apache.hadoop.hive.conf.HiveConf.ConfVars.METASTOREURIS; +import static org.assertj.core.api.Assertions.assertThat; + +import java.util.Map; +import java.util.function.Consumer; +import org.apache.hadoop.hive.conf.HiveConf; +import org.apache.iceberg.hive.TestHiveMetastore; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; +import org.apache.spark.sql.SparkSession; +import org.apache.spark.sql.internal.SQLConf; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.BeforeAll; + +public class CustomSparkTestBase { + + protected static TestHiveMetastore metastore = null; + protected static HiveConf hiveConf = null; + + @BeforeAll + public static void startMetastore() { + metastore = new TestHiveMetastore(); + metastore.start(); + hiveConf = metastore.hiveConf(); + } + + @AfterAll + public static void stopMetastore() throws Exception { + if (metastore != null) { + metastore.stop(); + metastore = null; + } + } + + protected SparkSession.Builder baseBuilder() { + Map disableUIConfig = + ImmutableMap.of( + "spark.ui.enabled", + "false", + "spark.metrics.conf.*.sink.servlet.class", + "org.apache.iceberg.spark.DummyMetricsServlet"); + return SparkSession.builder() + .master("local[1]") + .config(SQLConf.PARTITION_OVERWRITE_MODE().key(), "dynamic") + .config("spark.hadoop." + METASTOREURIS.varname, hiveConf.get(METASTOREURIS.varname)) + .config("spark.sql.legacy.respectNullabilityInTextDatasetConversion", "true") + .config("spark.sql.catalog.spark_catalog.type", "hive") + .config(disableUIConfig) + .appName("icebergCustomSparkTest") + .enableHiveSupport(); + } + + protected void withCustomSpark(Map overrides, Consumer test) + throws Exception { + assertThat(SparkSession.getActiveSession().isEmpty()) + .withFailMessage("A Spark session is already active!") + .isTrue(); + SparkSession.Builder sparkBuilder = baseBuilder(); + overrides.forEach((sparkBuilder::config)); + try (var sparkSession = sparkBuilder.create()) { + test.accept(sparkSession); + } + } +} diff --git a/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java b/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java new file mode 100644 index 000000000000..a5e0b0136a28 --- /dev/null +++ b/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java @@ -0,0 +1,78 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.iceberg.spark.sql; + +import static org.assertj.core.api.Assertions.assertThat; + +import java.util.Map; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; +import org.apache.iceberg.spark.CustomSparkTestBase; +import org.junit.jupiter.api.Test; + +public class TestStaticCatalogConfigs extends CustomSparkTestBase { + + @Test + public void sessionCatalogPicksUpDefaultDatabaseConfig() throws Exception { + Map overrides = + ImmutableMap.of( + "spark.sql.catalog.spark_catalog", "org.apache.iceberg.spark.SparkSessionCatalog", + "spark.sql.catalog.spark_catalog.defaultDatabase", "testDefaultDB"); + withCustomSpark( + overrides, + sparkSession -> { + String[] foundDefaultNamespace = + sparkSession.sessionState().catalogManager().v2SessionCatalog().defaultNamespace(); + + assertThat(foundDefaultNamespace).containsExactly("testDefaultDB"); + }); + } + + @Test + public void sparkCatalogPicksUpDefaultDatabaseConfig() throws Exception { + Map overrides = + ImmutableMap.of( + "spark.sql.catalog.spark_catalog", "org.apache.iceberg.spark.SparkCatalog", + "spark.sql.catalog.spark_catalog.defaultDatabase", "testDefaultDB"); + withCustomSpark( + overrides, + sparkSession -> { + String[] foundDefaultNamespace = + sparkSession.sessionState().catalogManager().v2SessionCatalog().defaultNamespace(); + + assertThat(foundDefaultNamespace).containsExactly("testDefaultDB"); + }); + } + + @Test + void sparkCatalogStillPrefersDefaultNamespaceConfig() throws Exception { + Map overrides = + ImmutableMap.of( + "spark.sql.catalog.spark_catalog", "org.apache.iceberg.spark.SparkCatalog", + "spark.sql.catalog.spark_catalog.defaultDatabase", "dbToBeOverridden", + "spark.sql.catalog.spark_catalog.default-namespace", "testDefaultDB"); + withCustomSpark( + overrides, + sparkSession -> { + String[] foundDefaultNamespace = + sparkSession.sessionState().catalogManager().v2SessionCatalog().defaultNamespace(); + + assertThat(foundDefaultNamespace).containsExactly("testDefaultDB"); + }); + } +} From 004a9224b296ef04cff76db9a937f25225ec448d Mon Sep 17 00:00:00 2001 From: Dzeri96 Date: Tue, 8 Sep 2026 16:57:00 +0200 Subject: [PATCH 2/9] Issue 14424: Removed changes to SparkCatalog. Will be a different PR --- docs/docs/spark-configuration.md | 3 +- .../apache/iceberg/spark/SparkCatalog.java | 11 +------ .../spark/sql/TestStaticCatalogConfigs.java | 33 ------------------- 3 files changed, 3 insertions(+), 44 deletions(-) diff --git a/docs/docs/spark-configuration.md b/docs/docs/spark-configuration.md index e1dedddb4b30..996aa82415f2 100644 --- a/docs/docs/spark-configuration.md +++ b/docs/docs/spark-configuration.md @@ -69,7 +69,8 @@ Both catalogs are configured using properties nested under the catalog name. Com | spark.sql.catalog._catalog-name_.catalog-impl | | The custom Iceberg catalog implementation. If `type` is null, `catalog-impl` must not be null. | | spark.sql.catalog._catalog-name_.io-impl | | The custom FileIO implementation. | | spark.sql.catalog._catalog-name_.metrics-reporter-impl | | The custom MetricsReporter implementation. | -| spark.sql.catalog._catalog-name_.defaultDatabase | default | The default current namespace for the catalog | +| spark.sql.catalog._catalog-name_.default-namespace | default | The default current namespace for the `SparkCatalog`. | +| spark.sql.catalog._catalog-name_.defaultDatabase | default | The default current namespace for the `SparkSessionCatalog`. | | spark.sql.catalog._catalog-name_.uri | thrift://host:port | Hive metastore URL for hive typed catalog, REST URL for REST typed catalog | | spark.sql.catalog._catalog-name_.warehouse | hdfs://nn:8020/warehouse/path | Base path for the warehouse directory | | spark.sql.catalog._catalog-name_.cache-enabled | `true` or `false` | Whether to enable catalog cache, default value is `true` | diff --git a/spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/SparkCatalog.java b/spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/SparkCatalog.java index 6d9fe36bdfb5..92825d763503 100644 --- a/spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/SparkCatalog.java +++ b/spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/SparkCatalog.java @@ -109,9 +109,7 @@ *

  • io-impl - a custom {@link org.apache.iceberg.io.FileIO} implementation to use *
  • metrics-reporter-impl - a custom {@link * org.apache.iceberg.metrics.MetricsReporter} implementation to use - *
  • default-namespace - DEPRECATED: use defaultDatabase - * instead - *
  • defaultDatabase - a namespace/database to use as the default + *
  • default-namespace - a namespace to use as the default *
  • cache-enabled - whether to enable catalog cache *
  • cache.case-sensitive - whether the catalog cache should compare table * identifiers in a case sensitive way @@ -815,14 +813,7 @@ public final void initialize(String name, CaseInsensitiveStringMap options) { : catalog; if (catalog instanceof SupportsNamespaces) { this.asNamespaceCatalog = (SupportsNamespaces) catalog; - if (options.containsKey("defaultDatabase")) { - this.defaultNamespace = - Splitter.on('.').splitToList(options.get("defaultDatabase")).toArray(new String[0]); - } if (options.containsKey("default-namespace")) { - LOG.warn( - "The `default-namespace` property is deprecated and will be removed in the next major version" - + " Please use Spark's `defaultDatabase` instead."); this.defaultNamespace = Splitter.on('.').splitToList(options.get("default-namespace")).toArray(new String[0]); } diff --git a/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java b/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java index a5e0b0136a28..508c9a7e755c 100644 --- a/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java +++ b/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java @@ -42,37 +42,4 @@ public void sessionCatalogPicksUpDefaultDatabaseConfig() throws Exception { assertThat(foundDefaultNamespace).containsExactly("testDefaultDB"); }); } - - @Test - public void sparkCatalogPicksUpDefaultDatabaseConfig() throws Exception { - Map overrides = - ImmutableMap.of( - "spark.sql.catalog.spark_catalog", "org.apache.iceberg.spark.SparkCatalog", - "spark.sql.catalog.spark_catalog.defaultDatabase", "testDefaultDB"); - withCustomSpark( - overrides, - sparkSession -> { - String[] foundDefaultNamespace = - sparkSession.sessionState().catalogManager().v2SessionCatalog().defaultNamespace(); - - assertThat(foundDefaultNamespace).containsExactly("testDefaultDB"); - }); - } - - @Test - void sparkCatalogStillPrefersDefaultNamespaceConfig() throws Exception { - Map overrides = - ImmutableMap.of( - "spark.sql.catalog.spark_catalog", "org.apache.iceberg.spark.SparkCatalog", - "spark.sql.catalog.spark_catalog.defaultDatabase", "dbToBeOverridden", - "spark.sql.catalog.spark_catalog.default-namespace", "testDefaultDB"); - withCustomSpark( - overrides, - sparkSession -> { - String[] foundDefaultNamespace = - sparkSession.sessionState().catalogManager().v2SessionCatalog().defaultNamespace(); - - assertThat(foundDefaultNamespace).containsExactly("testDefaultDB"); - }); - } } From 848c5204e0e6895e80fcf390356038c03ba864cf Mon Sep 17 00:00:00 2001 From: Dzeri96 Date: Tue, 8 Sep 2026 17:09:10 +0200 Subject: [PATCH 3/9] Issue 14424: Added changes to all Spark versions --- .../iceberg/spark/SparkSessionCatalog.java | 3 +- .../iceberg/spark/CustomSparkTestBase.java | 83 +++++++++++++++++++ .../spark/sql/TestStaticCatalogConfigs.java | 45 ++++++++++ .../iceberg/spark/SparkSessionCatalog.java | 3 +- .../iceberg/spark/CustomSparkTestBase.java | 83 +++++++++++++++++++ .../spark/sql/TestStaticCatalogConfigs.java | 45 ++++++++++ .../iceberg/spark/SparkSessionCatalog.java | 3 +- .../iceberg/spark/CustomSparkTestBase.java | 83 +++++++++++++++++++ .../spark/sql/TestStaticCatalogConfigs.java | 45 ++++++++++ 9 files changed, 387 insertions(+), 6 deletions(-) create mode 100644 spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java create mode 100644 spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java create mode 100644 spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java create mode 100644 spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java create mode 100644 spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java create mode 100644 spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java diff --git a/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/SparkSessionCatalog.java b/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/SparkSessionCatalog.java index eebe44946a9e..02945144b4f4 100644 --- a/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/SparkSessionCatalog.java +++ b/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/SparkSessionCatalog.java @@ -63,7 +63,6 @@ public class SparkSessionCatalog< T extends TableCatalog & FunctionCatalog & SupportsNamespaces & ViewCatalog> extends BaseCatalog implements CatalogExtension { - private static final String[] DEFAULT_NAMESPACE = new String[] {"default"}; private String catalogName = null; private TableCatalog icebergCatalog = null; @@ -92,7 +91,7 @@ protected TableCatalog buildSparkCatalog(String name, CaseInsensitiveStringMap o @Override public String[] defaultNamespace() { - return DEFAULT_NAMESPACE; + return getSessionCatalog().defaultNamespace(); } @Override diff --git a/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java new file mode 100644 index 000000000000..593739615b11 --- /dev/null +++ b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java @@ -0,0 +1,83 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.iceberg.spark; + +import static org.apache.hadoop.hive.conf.HiveConf.ConfVars.METASTOREURIS; +import static org.assertj.core.api.Assertions.assertThat; + +import java.util.Map; +import java.util.function.Consumer; +import org.apache.hadoop.hive.conf.HiveConf; +import org.apache.iceberg.hive.TestHiveMetastore; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; +import org.apache.spark.sql.SparkSession; +import org.apache.spark.sql.internal.SQLConf; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.BeforeAll; + +public class CustomSparkTestBase { + + protected static TestHiveMetastore metastore = null; + protected static HiveConf hiveConf = null; + + @BeforeAll + public static void startMetastore() { + metastore = new TestHiveMetastore(); + metastore.start(); + hiveConf = metastore.hiveConf(); + } + + @AfterAll + public static void stopMetastore() throws Exception { + if (metastore != null) { + metastore.stop(); + metastore = null; + } + } + + protected SparkSession.Builder baseBuilder() { + Map disableUIConfig = + ImmutableMap.of( + "spark.ui.enabled", + "false", + "spark.metrics.conf.*.sink.servlet.class", + "org.apache.iceberg.spark.DummyMetricsServlet"); + return SparkSession.builder() + .master("local[1]") + .config(SQLConf.PARTITION_OVERWRITE_MODE().key(), "dynamic") + .config("spark.hadoop." + METASTOREURIS.varname, hiveConf.get(METASTOREURIS.varname)) + .config("spark.sql.legacy.respectNullabilityInTextDatasetConversion", "true") + .config("spark.sql.catalog.spark_catalog.type", "hive") + .config(disableUIConfig) + .appName("icebergCustomSparkTest") + .enableHiveSupport(); + } + + protected void withCustomSpark(Map overrides, Consumer test) + throws Exception { + assertThat(SparkSession.getActiveSession().isEmpty()) + .withFailMessage("A Spark session is already active!") + .isTrue(); + SparkSession.Builder sparkBuilder = baseBuilder(); + overrides.forEach((sparkBuilder::config)); + try (var sparkSession = sparkBuilder.getOrCreate()) { + test.accept(sparkSession); + } + } +} diff --git a/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java new file mode 100644 index 000000000000..508c9a7e755c --- /dev/null +++ b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java @@ -0,0 +1,45 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.iceberg.spark.sql; + +import static org.assertj.core.api.Assertions.assertThat; + +import java.util.Map; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; +import org.apache.iceberg.spark.CustomSparkTestBase; +import org.junit.jupiter.api.Test; + +public class TestStaticCatalogConfigs extends CustomSparkTestBase { + + @Test + public void sessionCatalogPicksUpDefaultDatabaseConfig() throws Exception { + Map overrides = + ImmutableMap.of( + "spark.sql.catalog.spark_catalog", "org.apache.iceberg.spark.SparkSessionCatalog", + "spark.sql.catalog.spark_catalog.defaultDatabase", "testDefaultDB"); + withCustomSpark( + overrides, + sparkSession -> { + String[] foundDefaultNamespace = + sparkSession.sessionState().catalogManager().v2SessionCatalog().defaultNamespace(); + + assertThat(foundDefaultNamespace).containsExactly("testDefaultDB"); + }); + } +} diff --git a/spark/v4.1/spark/src/main/java/org/apache/iceberg/spark/SparkSessionCatalog.java b/spark/v4.1/spark/src/main/java/org/apache/iceberg/spark/SparkSessionCatalog.java index ec37435ba119..6ec99714129c 100644 --- a/spark/v4.1/spark/src/main/java/org/apache/iceberg/spark/SparkSessionCatalog.java +++ b/spark/v4.1/spark/src/main/java/org/apache/iceberg/spark/SparkSessionCatalog.java @@ -65,7 +65,6 @@ public class SparkSessionCatalog< T extends TableCatalog & FunctionCatalog & SupportsNamespaces & ViewCatalog> extends BaseCatalog implements CatalogExtension { - private static final String[] DEFAULT_NAMESPACE = new String[] {"default"}; private String catalogName = null; private TableCatalog icebergCatalog = null; @@ -94,7 +93,7 @@ protected TableCatalog buildSparkCatalog(String name, CaseInsensitiveStringMap o @Override public String[] defaultNamespace() { - return DEFAULT_NAMESPACE; + return getSessionCatalog().defaultNamespace(); } @Override diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java new file mode 100644 index 000000000000..078106685f0c --- /dev/null +++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java @@ -0,0 +1,83 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.iceberg.spark; + +import static org.apache.hadoop.hive.conf.HiveConf.ConfVars.METASTOREURIS; +import static org.assertj.core.api.Assertions.assertThat; + +import java.util.Map; +import java.util.function.Consumer; +import org.apache.hadoop.hive.conf.HiveConf; +import org.apache.iceberg.hive.TestHiveMetastore; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; +import org.apache.spark.sql.SparkSession; +import org.apache.spark.sql.internal.SQLConf; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.BeforeAll; + +public class CustomSparkTestBase { + + protected static TestHiveMetastore metastore = null; + protected static HiveConf hiveConf = null; + + @BeforeAll + public static void startMetastore() { + metastore = new TestHiveMetastore(); + metastore.start(); + hiveConf = metastore.hiveConf(); + } + + @AfterAll + public static void stopMetastore() throws Exception { + if (metastore != null) { + metastore.stop(); + metastore = null; + } + } + + protected SparkSession.Builder baseBuilder() { + Map disableUIConfig = + ImmutableMap.of( + "spark.ui.enabled", + "false", + "spark.metrics.conf.*.sink.servlet.class", + "org.apache.iceberg.spark.DummyMetricsServlet"); + return SparkSession.builder() + .master("local[1]") + .config(SQLConf.PARTITION_OVERWRITE_MODE().key(), "dynamic") + .config("spark.hadoop." + METASTOREURIS.varname, hiveConf.get(METASTOREURIS.varname)) + .config("spark.sql.legacy.respectNullabilityInTextDatasetConversion", "true") + .config("spark.sql.catalog.spark_catalog.type", "hive") + .config(disableUIConfig) + .appName("icebergCustomSparkTest") + .enableHiveSupport(); + } + + protected void withCustomSpark(Map overrides, Consumer test) + throws Exception { + assertThat(SparkSession.getActiveSession().isEmpty()) + .withFailMessage("A Spark session is already active!") + .isTrue(); + SparkSession.Builder sparkBuilder = baseBuilder(); + overrides.forEach((sparkBuilder::config)); + try (var sparkSession = sparkBuilder.create()) { + test.accept(sparkSession); + } + } +} diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java new file mode 100644 index 000000000000..508c9a7e755c --- /dev/null +++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java @@ -0,0 +1,45 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.iceberg.spark.sql; + +import static org.assertj.core.api.Assertions.assertThat; + +import java.util.Map; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; +import org.apache.iceberg.spark.CustomSparkTestBase; +import org.junit.jupiter.api.Test; + +public class TestStaticCatalogConfigs extends CustomSparkTestBase { + + @Test + public void sessionCatalogPicksUpDefaultDatabaseConfig() throws Exception { + Map overrides = + ImmutableMap.of( + "spark.sql.catalog.spark_catalog", "org.apache.iceberg.spark.SparkSessionCatalog", + "spark.sql.catalog.spark_catalog.defaultDatabase", "testDefaultDB"); + withCustomSpark( + overrides, + sparkSession -> { + String[] foundDefaultNamespace = + sparkSession.sessionState().catalogManager().v2SessionCatalog().defaultNamespace(); + + assertThat(foundDefaultNamespace).containsExactly("testDefaultDB"); + }); + } +} 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..a2bee3dff369 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 @@ -69,7 +69,6 @@ public class SparkSessionCatalog< T extends TableCatalog & FunctionCatalog & SupportsNamespaces & ViewCatalog> extends BaseCatalog implements CatalogExtension { - private static final String[] DEFAULT_NAMESPACE = new String[] {"default"}; private String catalogName = null; private TableCatalog icebergCatalog = null; @@ -98,7 +97,7 @@ protected TableCatalog buildSparkCatalog(String name, CaseInsensitiveStringMap o @Override public String[] defaultNamespace() { - return DEFAULT_NAMESPACE; + return getSessionCatalog().defaultNamespace(); } @Override diff --git a/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java b/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java new file mode 100644 index 000000000000..078106685f0c --- /dev/null +++ b/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java @@ -0,0 +1,83 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.iceberg.spark; + +import static org.apache.hadoop.hive.conf.HiveConf.ConfVars.METASTOREURIS; +import static org.assertj.core.api.Assertions.assertThat; + +import java.util.Map; +import java.util.function.Consumer; +import org.apache.hadoop.hive.conf.HiveConf; +import org.apache.iceberg.hive.TestHiveMetastore; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; +import org.apache.spark.sql.SparkSession; +import org.apache.spark.sql.internal.SQLConf; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.BeforeAll; + +public class CustomSparkTestBase { + + protected static TestHiveMetastore metastore = null; + protected static HiveConf hiveConf = null; + + @BeforeAll + public static void startMetastore() { + metastore = new TestHiveMetastore(); + metastore.start(); + hiveConf = metastore.hiveConf(); + } + + @AfterAll + public static void stopMetastore() throws Exception { + if (metastore != null) { + metastore.stop(); + metastore = null; + } + } + + protected SparkSession.Builder baseBuilder() { + Map disableUIConfig = + ImmutableMap.of( + "spark.ui.enabled", + "false", + "spark.metrics.conf.*.sink.servlet.class", + "org.apache.iceberg.spark.DummyMetricsServlet"); + return SparkSession.builder() + .master("local[1]") + .config(SQLConf.PARTITION_OVERWRITE_MODE().key(), "dynamic") + .config("spark.hadoop." + METASTOREURIS.varname, hiveConf.get(METASTOREURIS.varname)) + .config("spark.sql.legacy.respectNullabilityInTextDatasetConversion", "true") + .config("spark.sql.catalog.spark_catalog.type", "hive") + .config(disableUIConfig) + .appName("icebergCustomSparkTest") + .enableHiveSupport(); + } + + protected void withCustomSpark(Map overrides, Consumer test) + throws Exception { + assertThat(SparkSession.getActiveSession().isEmpty()) + .withFailMessage("A Spark session is already active!") + .isTrue(); + SparkSession.Builder sparkBuilder = baseBuilder(); + overrides.forEach((sparkBuilder::config)); + try (var sparkSession = sparkBuilder.create()) { + test.accept(sparkSession); + } + } +} diff --git a/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java b/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java new file mode 100644 index 000000000000..508c9a7e755c --- /dev/null +++ b/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java @@ -0,0 +1,45 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.iceberg.spark.sql; + +import static org.assertj.core.api.Assertions.assertThat; + +import java.util.Map; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; +import org.apache.iceberg.spark.CustomSparkTestBase; +import org.junit.jupiter.api.Test; + +public class TestStaticCatalogConfigs extends CustomSparkTestBase { + + @Test + public void sessionCatalogPicksUpDefaultDatabaseConfig() throws Exception { + Map overrides = + ImmutableMap.of( + "spark.sql.catalog.spark_catalog", "org.apache.iceberg.spark.SparkSessionCatalog", + "spark.sql.catalog.spark_catalog.defaultDatabase", "testDefaultDB"); + withCustomSpark( + overrides, + sparkSession -> { + String[] foundDefaultNamespace = + sparkSession.sessionState().catalogManager().v2SessionCatalog().defaultNamespace(); + + assertThat(foundDefaultNamespace).containsExactly("testDefaultDB"); + }); + } +} From eae6f5b9253a2e4a81eac54ffd7e0f6220184c50 Mon Sep 17 00:00:00 2001 From: Dzeri96 Date: Wed, 9 Sep 2026 12:29:17 +0200 Subject: [PATCH 4/9] Issue 14424: Fixed active spark session detection in CustomSparkTestBase for Spark 3.5 --- .../apache/iceberg/spark/CustomSparkTestBase.java | 12 ++++++++---- 1 file changed, 8 insertions(+), 4 deletions(-) diff --git a/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java index 593739615b11..85c69ac9ebe5 100644 --- a/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java +++ b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java @@ -69,11 +69,15 @@ protected SparkSession.Builder baseBuilder() { .enableHiveSupport(); } - protected void withCustomSpark(Map overrides, Consumer test) - throws Exception { - assertThat(SparkSession.getActiveSession().isEmpty()) + protected void withCustomSpark(Map overrides, Consumer test) { + // Needed in v3.5 because getActiveSession is less strict than in Spark v4+ + boolean hasValidActiveSession = + SparkSession.getActiveSession() + .fold(() -> false, session -> !session.sparkContext().isStopped()); + assertThat(hasValidActiveSession) .withFailMessage("A Spark session is already active!") - .isTrue(); + .isFalse(); + SparkSession.Builder sparkBuilder = baseBuilder(); overrides.forEach((sparkBuilder::config)); try (var sparkSession = sparkBuilder.getOrCreate()) { From 6fabe44192211079834ce1ab34a8624388298fd0 Mon Sep 17 00:00:00 2001 From: Dzeri96 Date: Wed, 9 Sep 2026 22:48:45 +0200 Subject: [PATCH 5/9] Issue 14424: Removed unnecessary exception from test signature in TestStaticCatalogConfigs.java for Spark 3.5 --- .../org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java index 508c9a7e755c..c9a6f596de11 100644 --- a/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java +++ b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java @@ -28,7 +28,7 @@ public class TestStaticCatalogConfigs extends CustomSparkTestBase { @Test - public void sessionCatalogPicksUpDefaultDatabaseConfig() throws Exception { + public void sessionCatalogPicksUpDefaultDatabaseConfig() { Map overrides = ImmutableMap.of( "spark.sql.catalog.spark_catalog", "org.apache.iceberg.spark.SparkSessionCatalog", From cd5ce09fdd2f87cf17eb44102aaa8a7e254c7257 Mon Sep 17 00:00:00 2001 From: Dzeri96 Date: Fri, 25 Sep 2026 10:14:05 +0200 Subject: [PATCH 6/9] Issue 14424: Reworked defaultDatabase config test --- docs/docs/spark-configuration.md | 24 ++--- gradle.properties | 2 +- .../iceberg/spark/CustomSparkTestBase.java | 87 ------------------- .../spark/TestSparkSessionCatalog.java | 17 ++++ .../spark/sql/TestStaticCatalogConfigs.java | 45 ---------- .../iceberg/spark/CustomSparkTestBase.java | 83 ------------------ .../spark/TestSparkSessionCatalog.java | 17 ++++ .../spark/sql/TestStaticCatalogConfigs.java | 45 ---------- .../iceberg/spark/CustomSparkTestBase.java | 83 ------------------ .../spark/TestSparkSessionCatalog.java | 17 ++++ .../spark/sql/TestStaticCatalogConfigs.java | 45 ---------- .../iceberg/spark/CustomSparkTestBase.java | 83 ------------------ .../spark/TestSparkSessionCatalog.java | 17 ++++ .../spark/sql/TestStaticCatalogConfigs.java | 45 ---------- 14 files changed, 81 insertions(+), 529 deletions(-) delete mode 100644 spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java delete mode 100644 spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java delete mode 100644 spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java delete mode 100644 spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java delete mode 100644 spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java delete mode 100644 spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java delete mode 100644 spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java delete mode 100644 spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java diff --git a/docs/docs/spark-configuration.md b/docs/docs/spark-configuration.md index 996aa82415f2..ba578badb7a6 100644 --- a/docs/docs/spark-configuration.md +++ b/docs/docs/spark-configuration.md @@ -63,23 +63,23 @@ Iceberg supplies two implementations: Both catalogs are configured using properties nested under the catalog name. Common configuration properties for Hive and Hadoop are: -| Property | Values | Description | -|---------------------------------------------------------------|----------------------------------------------------|----------------------------------------------------------------------------------------------------------------------------------------------------------------------| -| spark.sql.catalog._catalog-name_.type | `hive`, `hadoop`, `rest`, `glue`, `jdbc` or `nessie` | The underlying Iceberg catalog implementation, `HiveCatalog`, `HadoopCatalog`, `RESTCatalog`, `GlueCatalog`, `JdbcCatalog`, `NessieCatalog` or left unset if using a custom catalog | -| spark.sql.catalog._catalog-name_.catalog-impl | | The custom Iceberg catalog implementation. If `type` is null, `catalog-impl` must not be null. | +| Property | Values | Description | +| -------------------------------------------------- |----------------------------------------------------|----------------------------------------------------------------------------------------------------------------------------------------------------------------------| +| spark.sql.catalog._catalog-name_.type | `hive`, `hadoop`, `rest`, `glue`, `jdbc` or `nessie` | The underlying Iceberg catalog implementation, `HiveCatalog`, `HadoopCatalog`, `RESTCatalog`, `GlueCatalog`, `JdbcCatalog`, `NessieCatalog` or left unset if using a custom catalog | +| spark.sql.catalog._catalog-name_.catalog-impl | | The custom Iceberg catalog implementation. If `type` is null, `catalog-impl` must not be null. | | spark.sql.catalog._catalog-name_.io-impl | | The custom FileIO implementation. | | spark.sql.catalog._catalog-name_.metrics-reporter-impl | | The custom MetricsReporter implementation. | -| spark.sql.catalog._catalog-name_.default-namespace | default | The default current namespace for the `SparkCatalog`. | -| spark.sql.catalog._catalog-name_.defaultDatabase | default | The default current namespace for the `SparkSessionCatalog`. | -| spark.sql.catalog._catalog-name_.uri | thrift://host:port | Hive metastore URL for hive typed catalog, REST URL for REST typed catalog | -| spark.sql.catalog._catalog-name_.warehouse | hdfs://nn:8020/warehouse/path | Base path for the warehouse directory | -| spark.sql.catalog._catalog-name_.cache-enabled | `true` or `false` | Whether to enable catalog cache, default value is `true` | +| spark.sql.catalog._catalog-name_.default-namespace | default | The default current namespace for the `SparkCatalog`. | +| spark.sql.catalog._catalog-name_.defaultDatabase | default | The default current namespace for the `SparkSessionCatalog`. | +| spark.sql.catalog._catalog-name_.uri | thrift://host:port | Hive metastore URL for hive typed catalog, REST URL for REST typed catalog | +| spark.sql.catalog._catalog-name_.warehouse | hdfs://nn:8020/warehouse/path | Base path for the warehouse directory | +| spark.sql.catalog._catalog-name_.cache-enabled | `true` or `false` | Whether to enable catalog cache, default value is `true` | | spark.sql.catalog._catalog-name_.cache.expiration-interval-ms | `30000` (30 seconds) | Duration after which cached catalog entries are expired; Only effective if `cache-enabled` is `true`. `-1` disables cache expiration and `0` disables caching entirely, irrespective of `cache-enabled`. Default is `30000` (30 seconds) | | spark.sql.catalog._catalog-name_.table-default._propertyKey_ | | Default Iceberg table property value for property key _propertyKey_, which will be set on tables created by this catalog if not overridden | | spark.sql.catalog._catalog-name_.table-override._propertyKey_ | | Enforced Iceberg table property value for property key _propertyKey_, which cannot be overridden on table creation by user | -| spark.sql.catalog._catalog-name_.view-default._propertyKey_ | | Default Iceberg view property value for property key _propertyKey_, which will be set on views created by this catalog if not overridden | -| spark.sql.catalog._catalog-name_.view-override._propertyKey_ | | Enforced Iceberg view property value for property key _propertyKey_, which cannot be overridden on view creation by user | -| spark.sql.catalog._catalog-name_.use-nullable-query-schema | `true` or `false` | Whether to preserve fields' nullability when creating the table using CTAS and RTAS. If set to `true`, all fields will be marked as nullable. If set to `false`, fields' nullability will be preserved. The default value is `true`. Available in Spark 3.5 and above. | +| spark.sql.catalog._catalog-name_.view-default._propertyKey_ | | Default Iceberg view property value for property key _propertyKey_, which will be set on views created by this catalog if not overridden | +| spark.sql.catalog._catalog-name_.view-override._propertyKey_ | | Enforced Iceberg view property value for property key _propertyKey_, which cannot be overridden on view creation by user | +| spark.sql.catalog._catalog-name_.use-nullable-query-schema | `true` or `false` | Whether to preserve fields' nullability when creating the table using CTAS and RTAS. If set to `true`, all fields will be marked as nullable. If set to `false`, fields' nullability will be preserved. The default value is `true`. Available in Spark 3.5 and above. | Additional properties can be found in common [catalog configuration](catalog-properties.md). diff --git a/gradle.properties b/gradle.properties index 3546db6f063a..c87fa88aeb70 100644 --- a/gradle.properties +++ b/gradle.properties @@ -18,7 +18,7 @@ jmhJsonOutputPath=build/reports/jmh/results.json jmhIncludeRegex=.* systemProp.defaultFlinkVersions=2.3 systemProp.knownFlinkVersions=1.20,2.1,2.2,2.3 -systemProp.defaultSparkVersions=4.2 +systemProp.defaultSparkVersions=3.5,4.0,4.1,4.2 systemProp.knownSparkVersions=3.5,4.0,4.1,4.2 systemProp.defaultKafkaVersions=3 systemProp.knownKafkaVersions=3 diff --git a/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java deleted file mode 100644 index 85c69ac9ebe5..000000000000 --- a/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java +++ /dev/null @@ -1,87 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you under the Apache License, Version 2.0 (the - * "License"); you may not use this file except in compliance - * with the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, - * software distributed under the License is distributed on an - * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY - * KIND, either express or implied. See the License for the - * specific language governing permissions and limitations - * under the License. - */ -package org.apache.iceberg.spark; - -import static org.apache.hadoop.hive.conf.HiveConf.ConfVars.METASTOREURIS; -import static org.assertj.core.api.Assertions.assertThat; - -import java.util.Map; -import java.util.function.Consumer; -import org.apache.hadoop.hive.conf.HiveConf; -import org.apache.iceberg.hive.TestHiveMetastore; -import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; -import org.apache.spark.sql.SparkSession; -import org.apache.spark.sql.internal.SQLConf; -import org.junit.jupiter.api.AfterAll; -import org.junit.jupiter.api.BeforeAll; - -public class CustomSparkTestBase { - - protected static TestHiveMetastore metastore = null; - protected static HiveConf hiveConf = null; - - @BeforeAll - public static void startMetastore() { - metastore = new TestHiveMetastore(); - metastore.start(); - hiveConf = metastore.hiveConf(); - } - - @AfterAll - public static void stopMetastore() throws Exception { - if (metastore != null) { - metastore.stop(); - metastore = null; - } - } - - protected SparkSession.Builder baseBuilder() { - Map disableUIConfig = - ImmutableMap.of( - "spark.ui.enabled", - "false", - "spark.metrics.conf.*.sink.servlet.class", - "org.apache.iceberg.spark.DummyMetricsServlet"); - return SparkSession.builder() - .master("local[1]") - .config(SQLConf.PARTITION_OVERWRITE_MODE().key(), "dynamic") - .config("spark.hadoop." + METASTOREURIS.varname, hiveConf.get(METASTOREURIS.varname)) - .config("spark.sql.legacy.respectNullabilityInTextDatasetConversion", "true") - .config("spark.sql.catalog.spark_catalog.type", "hive") - .config(disableUIConfig) - .appName("icebergCustomSparkTest") - .enableHiveSupport(); - } - - protected void withCustomSpark(Map overrides, Consumer test) { - // Needed in v3.5 because getActiveSession is less strict than in Spark v4+ - boolean hasValidActiveSession = - SparkSession.getActiveSession() - .fold(() -> false, session -> !session.sparkContext().isStopped()); - assertThat(hasValidActiveSession) - .withFailMessage("A Spark session is already active!") - .isFalse(); - - SparkSession.Builder sparkBuilder = baseBuilder(); - overrides.forEach((sparkBuilder::config)); - try (var sparkSession = sparkBuilder.getOrCreate()) { - test.accept(sparkSession); - } - } -} diff --git a/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java index d6d5a237d2b0..7615dac7369c 100644 --- a/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java +++ b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java @@ -32,6 +32,9 @@ import org.apache.spark.sql.connector.catalog.SupportsNamespaces; import org.apache.spark.sql.connector.catalog.TableCatalog; import org.apache.spark.sql.connector.catalog.ViewCatalog; +import org.apache.spark.sql.execution.datasources.v2.V2SessionCatalog; +import org.apache.spark.sql.internal.SQLConf; +import org.apache.spark.sql.internal.StaticSQLConf; import org.apache.spark.sql.util.CaseInsensitiveStringMap; import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.BeforeEach; @@ -129,6 +132,20 @@ public void listViewsReturnsSessionCatalogViews() throws NoSuchNamespaceExceptio assertThat(catalog.listViews("default")).containsExactly(viewIdent); } + @Test + public void sessionCatalogPicksUpDefaultDatabaseConfig() { + SQLConf sqlConf = new SQLConf(); + sqlConf.setConf(StaticSQLConf.CATALOG_DEFAULT_DATABASE(), "testDefaultDB"); + SQLConf.setSQLConfGetter(() -> sqlConf); + + var v1SessionCatalogMock = mock(org.apache.spark.sql.catalyst.catalog.SessionCatalog.class); + V2SessionCatalog v2SessionCatalog = new V2SessionCatalog(v1SessionCatalogMock); + + SparkSessionCatalog icebergSessionCatalog = new NoViewCatalog<>(); + icebergSessionCatalog.setDelegateCatalog(v2SessionCatalog); + assertThat(icebergSessionCatalog.defaultNamespace()).containsExactly("testDefaultDB"); + } + private static class NoViewCatalog< T extends TableCatalog & FunctionCatalog & SupportsNamespaces & ViewCatalog> extends SparkSessionCatalog { diff --git a/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java deleted file mode 100644 index c9a6f596de11..000000000000 --- a/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java +++ /dev/null @@ -1,45 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you under the Apache License, Version 2.0 (the - * "License"); you may not use this file except in compliance - * with the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, - * software distributed under the License is distributed on an - * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY - * KIND, either express or implied. See the License for the - * specific language governing permissions and limitations - * under the License. - */ -package org.apache.iceberg.spark.sql; - -import static org.assertj.core.api.Assertions.assertThat; - -import java.util.Map; -import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; -import org.apache.iceberg.spark.CustomSparkTestBase; -import org.junit.jupiter.api.Test; - -public class TestStaticCatalogConfigs extends CustomSparkTestBase { - - @Test - public void sessionCatalogPicksUpDefaultDatabaseConfig() { - Map overrides = - ImmutableMap.of( - "spark.sql.catalog.spark_catalog", "org.apache.iceberg.spark.SparkSessionCatalog", - "spark.sql.catalog.spark_catalog.defaultDatabase", "testDefaultDB"); - withCustomSpark( - overrides, - sparkSession -> { - String[] foundDefaultNamespace = - sparkSession.sessionState().catalogManager().v2SessionCatalog().defaultNamespace(); - - assertThat(foundDefaultNamespace).containsExactly("testDefaultDB"); - }); - } -} diff --git a/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java b/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java deleted file mode 100644 index 078106685f0c..000000000000 --- a/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java +++ /dev/null @@ -1,83 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you under the Apache License, Version 2.0 (the - * "License"); you may not use this file except in compliance - * with the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, - * software distributed under the License is distributed on an - * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY - * KIND, either express or implied. See the License for the - * specific language governing permissions and limitations - * under the License. - */ -package org.apache.iceberg.spark; - -import static org.apache.hadoop.hive.conf.HiveConf.ConfVars.METASTOREURIS; -import static org.assertj.core.api.Assertions.assertThat; - -import java.util.Map; -import java.util.function.Consumer; -import org.apache.hadoop.hive.conf.HiveConf; -import org.apache.iceberg.hive.TestHiveMetastore; -import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; -import org.apache.spark.sql.SparkSession; -import org.apache.spark.sql.internal.SQLConf; -import org.junit.jupiter.api.AfterAll; -import org.junit.jupiter.api.BeforeAll; - -public class CustomSparkTestBase { - - protected static TestHiveMetastore metastore = null; - protected static HiveConf hiveConf = null; - - @BeforeAll - public static void startMetastore() { - metastore = new TestHiveMetastore(); - metastore.start(); - hiveConf = metastore.hiveConf(); - } - - @AfterAll - public static void stopMetastore() throws Exception { - if (metastore != null) { - metastore.stop(); - metastore = null; - } - } - - protected SparkSession.Builder baseBuilder() { - Map disableUIConfig = - ImmutableMap.of( - "spark.ui.enabled", - "false", - "spark.metrics.conf.*.sink.servlet.class", - "org.apache.iceberg.spark.DummyMetricsServlet"); - return SparkSession.builder() - .master("local[1]") - .config(SQLConf.PARTITION_OVERWRITE_MODE().key(), "dynamic") - .config("spark.hadoop." + METASTOREURIS.varname, hiveConf.get(METASTOREURIS.varname)) - .config("spark.sql.legacy.respectNullabilityInTextDatasetConversion", "true") - .config("spark.sql.catalog.spark_catalog.type", "hive") - .config(disableUIConfig) - .appName("icebergCustomSparkTest") - .enableHiveSupport(); - } - - protected void withCustomSpark(Map overrides, Consumer test) - throws Exception { - assertThat(SparkSession.getActiveSession().isEmpty()) - .withFailMessage("A Spark session is already active!") - .isTrue(); - SparkSession.Builder sparkBuilder = baseBuilder(); - overrides.forEach((sparkBuilder::config)); - try (var sparkSession = sparkBuilder.create()) { - test.accept(sparkSession); - } - } -} diff --git a/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java b/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java index d6d5a237d2b0..7615dac7369c 100644 --- a/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java +++ b/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java @@ -32,6 +32,9 @@ import org.apache.spark.sql.connector.catalog.SupportsNamespaces; import org.apache.spark.sql.connector.catalog.TableCatalog; import org.apache.spark.sql.connector.catalog.ViewCatalog; +import org.apache.spark.sql.execution.datasources.v2.V2SessionCatalog; +import org.apache.spark.sql.internal.SQLConf; +import org.apache.spark.sql.internal.StaticSQLConf; import org.apache.spark.sql.util.CaseInsensitiveStringMap; import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.BeforeEach; @@ -129,6 +132,20 @@ public void listViewsReturnsSessionCatalogViews() throws NoSuchNamespaceExceptio assertThat(catalog.listViews("default")).containsExactly(viewIdent); } + @Test + public void sessionCatalogPicksUpDefaultDatabaseConfig() { + SQLConf sqlConf = new SQLConf(); + sqlConf.setConf(StaticSQLConf.CATALOG_DEFAULT_DATABASE(), "testDefaultDB"); + SQLConf.setSQLConfGetter(() -> sqlConf); + + var v1SessionCatalogMock = mock(org.apache.spark.sql.catalyst.catalog.SessionCatalog.class); + V2SessionCatalog v2SessionCatalog = new V2SessionCatalog(v1SessionCatalogMock); + + SparkSessionCatalog icebergSessionCatalog = new NoViewCatalog<>(); + icebergSessionCatalog.setDelegateCatalog(v2SessionCatalog); + assertThat(icebergSessionCatalog.defaultNamespace()).containsExactly("testDefaultDB"); + } + private static class NoViewCatalog< T extends TableCatalog & FunctionCatalog & SupportsNamespaces & ViewCatalog> extends SparkSessionCatalog { diff --git a/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java b/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java deleted file mode 100644 index 508c9a7e755c..000000000000 --- a/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java +++ /dev/null @@ -1,45 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you under the Apache License, Version 2.0 (the - * "License"); you may not use this file except in compliance - * with the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, - * software distributed under the License is distributed on an - * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY - * KIND, either express or implied. See the License for the - * specific language governing permissions and limitations - * under the License. - */ -package org.apache.iceberg.spark.sql; - -import static org.assertj.core.api.Assertions.assertThat; - -import java.util.Map; -import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; -import org.apache.iceberg.spark.CustomSparkTestBase; -import org.junit.jupiter.api.Test; - -public class TestStaticCatalogConfigs extends CustomSparkTestBase { - - @Test - public void sessionCatalogPicksUpDefaultDatabaseConfig() throws Exception { - Map overrides = - ImmutableMap.of( - "spark.sql.catalog.spark_catalog", "org.apache.iceberg.spark.SparkSessionCatalog", - "spark.sql.catalog.spark_catalog.defaultDatabase", "testDefaultDB"); - withCustomSpark( - overrides, - sparkSession -> { - String[] foundDefaultNamespace = - sparkSession.sessionState().catalogManager().v2SessionCatalog().defaultNamespace(); - - assertThat(foundDefaultNamespace).containsExactly("testDefaultDB"); - }); - } -} diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java deleted file mode 100644 index 078106685f0c..000000000000 --- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java +++ /dev/null @@ -1,83 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you under the Apache License, Version 2.0 (the - * "License"); you may not use this file except in compliance - * with the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, - * software distributed under the License is distributed on an - * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY - * KIND, either express or implied. See the License for the - * specific language governing permissions and limitations - * under the License. - */ -package org.apache.iceberg.spark; - -import static org.apache.hadoop.hive.conf.HiveConf.ConfVars.METASTOREURIS; -import static org.assertj.core.api.Assertions.assertThat; - -import java.util.Map; -import java.util.function.Consumer; -import org.apache.hadoop.hive.conf.HiveConf; -import org.apache.iceberg.hive.TestHiveMetastore; -import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; -import org.apache.spark.sql.SparkSession; -import org.apache.spark.sql.internal.SQLConf; -import org.junit.jupiter.api.AfterAll; -import org.junit.jupiter.api.BeforeAll; - -public class CustomSparkTestBase { - - protected static TestHiveMetastore metastore = null; - protected static HiveConf hiveConf = null; - - @BeforeAll - public static void startMetastore() { - metastore = new TestHiveMetastore(); - metastore.start(); - hiveConf = metastore.hiveConf(); - } - - @AfterAll - public static void stopMetastore() throws Exception { - if (metastore != null) { - metastore.stop(); - metastore = null; - } - } - - protected SparkSession.Builder baseBuilder() { - Map disableUIConfig = - ImmutableMap.of( - "spark.ui.enabled", - "false", - "spark.metrics.conf.*.sink.servlet.class", - "org.apache.iceberg.spark.DummyMetricsServlet"); - return SparkSession.builder() - .master("local[1]") - .config(SQLConf.PARTITION_OVERWRITE_MODE().key(), "dynamic") - .config("spark.hadoop." + METASTOREURIS.varname, hiveConf.get(METASTOREURIS.varname)) - .config("spark.sql.legacy.respectNullabilityInTextDatasetConversion", "true") - .config("spark.sql.catalog.spark_catalog.type", "hive") - .config(disableUIConfig) - .appName("icebergCustomSparkTest") - .enableHiveSupport(); - } - - protected void withCustomSpark(Map overrides, Consumer test) - throws Exception { - assertThat(SparkSession.getActiveSession().isEmpty()) - .withFailMessage("A Spark session is already active!") - .isTrue(); - SparkSession.Builder sparkBuilder = baseBuilder(); - overrides.forEach((sparkBuilder::config)); - try (var sparkSession = sparkBuilder.create()) { - test.accept(sparkSession); - } - } -} diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java index d6d5a237d2b0..7615dac7369c 100644 --- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java +++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java @@ -32,6 +32,9 @@ import org.apache.spark.sql.connector.catalog.SupportsNamespaces; import org.apache.spark.sql.connector.catalog.TableCatalog; import org.apache.spark.sql.connector.catalog.ViewCatalog; +import org.apache.spark.sql.execution.datasources.v2.V2SessionCatalog; +import org.apache.spark.sql.internal.SQLConf; +import org.apache.spark.sql.internal.StaticSQLConf; import org.apache.spark.sql.util.CaseInsensitiveStringMap; import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.BeforeEach; @@ -129,6 +132,20 @@ public void listViewsReturnsSessionCatalogViews() throws NoSuchNamespaceExceptio assertThat(catalog.listViews("default")).containsExactly(viewIdent); } + @Test + public void sessionCatalogPicksUpDefaultDatabaseConfig() { + SQLConf sqlConf = new SQLConf(); + sqlConf.setConf(StaticSQLConf.CATALOG_DEFAULT_DATABASE(), "testDefaultDB"); + SQLConf.setSQLConfGetter(() -> sqlConf); + + var v1SessionCatalogMock = mock(org.apache.spark.sql.catalyst.catalog.SessionCatalog.class); + V2SessionCatalog v2SessionCatalog = new V2SessionCatalog(v1SessionCatalogMock); + + SparkSessionCatalog icebergSessionCatalog = new NoViewCatalog<>(); + icebergSessionCatalog.setDelegateCatalog(v2SessionCatalog); + assertThat(icebergSessionCatalog.defaultNamespace()).containsExactly("testDefaultDB"); + } + private static class NoViewCatalog< T extends TableCatalog & FunctionCatalog & SupportsNamespaces & ViewCatalog> extends SparkSessionCatalog { diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java deleted file mode 100644 index 508c9a7e755c..000000000000 --- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java +++ /dev/null @@ -1,45 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you under the Apache License, Version 2.0 (the - * "License"); you may not use this file except in compliance - * with the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, - * software distributed under the License is distributed on an - * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY - * KIND, either express or implied. See the License for the - * specific language governing permissions and limitations - * under the License. - */ -package org.apache.iceberg.spark.sql; - -import static org.assertj.core.api.Assertions.assertThat; - -import java.util.Map; -import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; -import org.apache.iceberg.spark.CustomSparkTestBase; -import org.junit.jupiter.api.Test; - -public class TestStaticCatalogConfigs extends CustomSparkTestBase { - - @Test - public void sessionCatalogPicksUpDefaultDatabaseConfig() throws Exception { - Map overrides = - ImmutableMap.of( - "spark.sql.catalog.spark_catalog", "org.apache.iceberg.spark.SparkSessionCatalog", - "spark.sql.catalog.spark_catalog.defaultDatabase", "testDefaultDB"); - withCustomSpark( - overrides, - sparkSession -> { - String[] foundDefaultNamespace = - sparkSession.sessionState().catalogManager().v2SessionCatalog().defaultNamespace(); - - assertThat(foundDefaultNamespace).containsExactly("testDefaultDB"); - }); - } -} diff --git a/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java b/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java deleted file mode 100644 index 078106685f0c..000000000000 --- a/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/CustomSparkTestBase.java +++ /dev/null @@ -1,83 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you under the Apache License, Version 2.0 (the - * "License"); you may not use this file except in compliance - * with the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, - * software distributed under the License is distributed on an - * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY - * KIND, either express or implied. See the License for the - * specific language governing permissions and limitations - * under the License. - */ -package org.apache.iceberg.spark; - -import static org.apache.hadoop.hive.conf.HiveConf.ConfVars.METASTOREURIS; -import static org.assertj.core.api.Assertions.assertThat; - -import java.util.Map; -import java.util.function.Consumer; -import org.apache.hadoop.hive.conf.HiveConf; -import org.apache.iceberg.hive.TestHiveMetastore; -import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; -import org.apache.spark.sql.SparkSession; -import org.apache.spark.sql.internal.SQLConf; -import org.junit.jupiter.api.AfterAll; -import org.junit.jupiter.api.BeforeAll; - -public class CustomSparkTestBase { - - protected static TestHiveMetastore metastore = null; - protected static HiveConf hiveConf = null; - - @BeforeAll - public static void startMetastore() { - metastore = new TestHiveMetastore(); - metastore.start(); - hiveConf = metastore.hiveConf(); - } - - @AfterAll - public static void stopMetastore() throws Exception { - if (metastore != null) { - metastore.stop(); - metastore = null; - } - } - - protected SparkSession.Builder baseBuilder() { - Map disableUIConfig = - ImmutableMap.of( - "spark.ui.enabled", - "false", - "spark.metrics.conf.*.sink.servlet.class", - "org.apache.iceberg.spark.DummyMetricsServlet"); - return SparkSession.builder() - .master("local[1]") - .config(SQLConf.PARTITION_OVERWRITE_MODE().key(), "dynamic") - .config("spark.hadoop." + METASTOREURIS.varname, hiveConf.get(METASTOREURIS.varname)) - .config("spark.sql.legacy.respectNullabilityInTextDatasetConversion", "true") - .config("spark.sql.catalog.spark_catalog.type", "hive") - .config(disableUIConfig) - .appName("icebergCustomSparkTest") - .enableHiveSupport(); - } - - protected void withCustomSpark(Map overrides, Consumer test) - throws Exception { - assertThat(SparkSession.getActiveSession().isEmpty()) - .withFailMessage("A Spark session is already active!") - .isTrue(); - SparkSession.Builder sparkBuilder = baseBuilder(); - overrides.forEach((sparkBuilder::config)); - try (var sparkSession = sparkBuilder.create()) { - test.accept(sparkSession); - } - } -} diff --git a/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java b/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java index 46897c2f9a44..7d52e60e4d23 100644 --- a/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java +++ b/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java @@ -48,6 +48,9 @@ import org.apache.spark.sql.connector.catalog.ViewCatalog; import org.apache.spark.sql.connector.expressions.Transform; import org.apache.spark.sql.connector.metric.CustomTaskMetric; +import org.apache.spark.sql.execution.datasources.v2.V2SessionCatalog; +import org.apache.spark.sql.internal.SQLConf; +import org.apache.spark.sql.internal.StaticSQLConf; import org.apache.spark.sql.types.StructType; import org.apache.spark.sql.util.CaseInsensitiveStringMap; import org.junit.jupiter.api.BeforeAll; @@ -492,6 +495,20 @@ public void createOrReplaceViewUsesExistingSessionCatalogView() verify(icebergViewCatalog, never()).createOrReplaceView(viewIdent, replacement); } + @Test + public void sessionCatalogPicksUpDefaultDatabaseConfig() { + SQLConf sqlConf = new SQLConf(); + sqlConf.setConf(StaticSQLConf.CATALOG_DEFAULT_DATABASE(), "testDefaultDB"); + SQLConf.setSQLConfGetter(() -> sqlConf); + + var v1SessionCatalogMock = mock(org.apache.spark.sql.catalyst.catalog.SessionCatalog.class); + V2SessionCatalog v2SessionCatalog = new V2SessionCatalog(v1SessionCatalogMock); + + SparkSessionCatalog icebergSessionCatalog = new NoViewCatalog<>(); + icebergSessionCatalog.setDelegateCatalog(v2SessionCatalog); + assertThat(icebergSessionCatalog.defaultNamespace()).containsExactly("testDefaultDB"); + } + private static TableCatalog sessionCatalogWithViews() { return mock( TableCatalog.class, diff --git a/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java b/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java deleted file mode 100644 index 508c9a7e755c..000000000000 --- a/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/sql/TestStaticCatalogConfigs.java +++ /dev/null @@ -1,45 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you under the Apache License, Version 2.0 (the - * "License"); you may not use this file except in compliance - * with the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, - * software distributed under the License is distributed on an - * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY - * KIND, either express or implied. See the License for the - * specific language governing permissions and limitations - * under the License. - */ -package org.apache.iceberg.spark.sql; - -import static org.assertj.core.api.Assertions.assertThat; - -import java.util.Map; -import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; -import org.apache.iceberg.spark.CustomSparkTestBase; -import org.junit.jupiter.api.Test; - -public class TestStaticCatalogConfigs extends CustomSparkTestBase { - - @Test - public void sessionCatalogPicksUpDefaultDatabaseConfig() throws Exception { - Map overrides = - ImmutableMap.of( - "spark.sql.catalog.spark_catalog", "org.apache.iceberg.spark.SparkSessionCatalog", - "spark.sql.catalog.spark_catalog.defaultDatabase", "testDefaultDB"); - withCustomSpark( - overrides, - sparkSession -> { - String[] foundDefaultNamespace = - sparkSession.sessionState().catalogManager().v2SessionCatalog().defaultNamespace(); - - assertThat(foundDefaultNamespace).containsExactly("testDefaultDB"); - }); - } -} From e8a11683c79551a263e9d83fa3a2a3be408ba8d1 Mon Sep 17 00:00:00 2001 From: Dzeri96 Date: Fri, 25 Sep 2026 10:29:38 +0200 Subject: [PATCH 7/9] Issue 14424: Reverted accidentally-commited gradle.properties --- gradle.properties | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/gradle.properties b/gradle.properties index c87fa88aeb70..3546db6f063a 100644 --- a/gradle.properties +++ b/gradle.properties @@ -18,7 +18,7 @@ jmhJsonOutputPath=build/reports/jmh/results.json jmhIncludeRegex=.* systemProp.defaultFlinkVersions=2.3 systemProp.knownFlinkVersions=1.20,2.1,2.2,2.3 -systemProp.defaultSparkVersions=3.5,4.0,4.1,4.2 +systemProp.defaultSparkVersions=4.2 systemProp.knownSparkVersions=3.5,4.0,4.1,4.2 systemProp.defaultKafkaVersions=3 systemProp.knownKafkaVersions=3 From 686f3aad6f6c45e19c5c77a02a00885dcd774772 Mon Sep 17 00:00:00 2001 From: Dzeri96 Date: Fri, 25 Sep 2026 16:48:55 +0200 Subject: [PATCH 8/9] Issue 14424: Re-worked defaultDatabase tests with SQLConf.withExistingConf --- .../spark/TestSparkSessionCatalog.java | 23 ++++++++++++------- .../spark/TestSparkSessionCatalog.java | 23 ++++++++++++------- .../spark/TestSparkSessionCatalog.java | 23 ++++++++++++------- .../spark/TestSparkSessionCatalog.java | 23 ++++++++++++------- 4 files changed, 60 insertions(+), 32 deletions(-) diff --git a/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java index 7615dac7369c..5019100e6041 100644 --- a/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java +++ b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java @@ -39,6 +39,7 @@ import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import scala.Function0; public class TestSparkSessionCatalog extends TestBase { private final String envHmsUriKey = "spark.hadoop." + METASTOREURIS.varname; @@ -136,14 +137,20 @@ public void listViewsReturnsSessionCatalogViews() throws NoSuchNamespaceExceptio public void sessionCatalogPicksUpDefaultDatabaseConfig() { SQLConf sqlConf = new SQLConf(); sqlConf.setConf(StaticSQLConf.CATALOG_DEFAULT_DATABASE(), "testDefaultDB"); - SQLConf.setSQLConfGetter(() -> sqlConf); - - var v1SessionCatalogMock = mock(org.apache.spark.sql.catalyst.catalog.SessionCatalog.class); - V2SessionCatalog v2SessionCatalog = new V2SessionCatalog(v1SessionCatalogMock); - - SparkSessionCatalog icebergSessionCatalog = new NoViewCatalog<>(); - icebergSessionCatalog.setDelegateCatalog(v2SessionCatalog); - assertThat(icebergSessionCatalog.defaultNamespace()).containsExactly("testDefaultDB"); + String[] result = + SQLConf.withExistingConf( + sqlConf, + (Function0) + () -> { + var v1SessionCatalogMock = + mock(org.apache.spark.sql.catalyst.catalog.SessionCatalog.class); + V2SessionCatalog v2SessionCatalog = new V2SessionCatalog(v1SessionCatalogMock); + + SparkSessionCatalog icebergSessionCatalog = new NoViewCatalog<>(); + icebergSessionCatalog.setDelegateCatalog(v2SessionCatalog); + return icebergSessionCatalog.defaultNamespace(); + }); + assertThat(result).containsExactly("testDefaultDB"); } private static class NoViewCatalog< diff --git a/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java b/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java index 7615dac7369c..5019100e6041 100644 --- a/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java +++ b/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java @@ -39,6 +39,7 @@ import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import scala.Function0; public class TestSparkSessionCatalog extends TestBase { private final String envHmsUriKey = "spark.hadoop." + METASTOREURIS.varname; @@ -136,14 +137,20 @@ public void listViewsReturnsSessionCatalogViews() throws NoSuchNamespaceExceptio public void sessionCatalogPicksUpDefaultDatabaseConfig() { SQLConf sqlConf = new SQLConf(); sqlConf.setConf(StaticSQLConf.CATALOG_DEFAULT_DATABASE(), "testDefaultDB"); - SQLConf.setSQLConfGetter(() -> sqlConf); - - var v1SessionCatalogMock = mock(org.apache.spark.sql.catalyst.catalog.SessionCatalog.class); - V2SessionCatalog v2SessionCatalog = new V2SessionCatalog(v1SessionCatalogMock); - - SparkSessionCatalog icebergSessionCatalog = new NoViewCatalog<>(); - icebergSessionCatalog.setDelegateCatalog(v2SessionCatalog); - assertThat(icebergSessionCatalog.defaultNamespace()).containsExactly("testDefaultDB"); + String[] result = + SQLConf.withExistingConf( + sqlConf, + (Function0) + () -> { + var v1SessionCatalogMock = + mock(org.apache.spark.sql.catalyst.catalog.SessionCatalog.class); + V2SessionCatalog v2SessionCatalog = new V2SessionCatalog(v1SessionCatalogMock); + + SparkSessionCatalog icebergSessionCatalog = new NoViewCatalog<>(); + icebergSessionCatalog.setDelegateCatalog(v2SessionCatalog); + return icebergSessionCatalog.defaultNamespace(); + }); + assertThat(result).containsExactly("testDefaultDB"); } private static class NoViewCatalog< diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java index 7615dac7369c..5019100e6041 100644 --- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java +++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java @@ -39,6 +39,7 @@ import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import scala.Function0; public class TestSparkSessionCatalog extends TestBase { private final String envHmsUriKey = "spark.hadoop." + METASTOREURIS.varname; @@ -136,14 +137,20 @@ public void listViewsReturnsSessionCatalogViews() throws NoSuchNamespaceExceptio public void sessionCatalogPicksUpDefaultDatabaseConfig() { SQLConf sqlConf = new SQLConf(); sqlConf.setConf(StaticSQLConf.CATALOG_DEFAULT_DATABASE(), "testDefaultDB"); - SQLConf.setSQLConfGetter(() -> sqlConf); - - var v1SessionCatalogMock = mock(org.apache.spark.sql.catalyst.catalog.SessionCatalog.class); - V2SessionCatalog v2SessionCatalog = new V2SessionCatalog(v1SessionCatalogMock); - - SparkSessionCatalog icebergSessionCatalog = new NoViewCatalog<>(); - icebergSessionCatalog.setDelegateCatalog(v2SessionCatalog); - assertThat(icebergSessionCatalog.defaultNamespace()).containsExactly("testDefaultDB"); + String[] result = + SQLConf.withExistingConf( + sqlConf, + (Function0) + () -> { + var v1SessionCatalogMock = + mock(org.apache.spark.sql.catalyst.catalog.SessionCatalog.class); + V2SessionCatalog v2SessionCatalog = new V2SessionCatalog(v1SessionCatalogMock); + + SparkSessionCatalog icebergSessionCatalog = new NoViewCatalog<>(); + icebergSessionCatalog.setDelegateCatalog(v2SessionCatalog); + return icebergSessionCatalog.defaultNamespace(); + }); + assertThat(result).containsExactly("testDefaultDB"); } private static class NoViewCatalog< diff --git a/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java b/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java index 7d52e60e4d23..5e643915ccd6 100644 --- a/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java +++ b/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java @@ -56,6 +56,7 @@ import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import scala.Function0; public class TestSparkSessionCatalog extends TestBase { private final String envHmsUriKey = "spark.hadoop." + METASTOREURIS.varname; @@ -499,14 +500,20 @@ public void createOrReplaceViewUsesExistingSessionCatalogView() public void sessionCatalogPicksUpDefaultDatabaseConfig() { SQLConf sqlConf = new SQLConf(); sqlConf.setConf(StaticSQLConf.CATALOG_DEFAULT_DATABASE(), "testDefaultDB"); - SQLConf.setSQLConfGetter(() -> sqlConf); - - var v1SessionCatalogMock = mock(org.apache.spark.sql.catalyst.catalog.SessionCatalog.class); - V2SessionCatalog v2SessionCatalog = new V2SessionCatalog(v1SessionCatalogMock); - - SparkSessionCatalog icebergSessionCatalog = new NoViewCatalog<>(); - icebergSessionCatalog.setDelegateCatalog(v2SessionCatalog); - assertThat(icebergSessionCatalog.defaultNamespace()).containsExactly("testDefaultDB"); + String[] result = + SQLConf.withExistingConf( + sqlConf, + (Function0) + () -> { + var v1SessionCatalogMock = + mock(org.apache.spark.sql.catalyst.catalog.SessionCatalog.class); + V2SessionCatalog v2SessionCatalog = new V2SessionCatalog(v1SessionCatalogMock); + + SparkSessionCatalog icebergSessionCatalog = new NoViewCatalog<>(); + icebergSessionCatalog.setDelegateCatalog(v2SessionCatalog); + return icebergSessionCatalog.defaultNamespace(); + }); + assertThat(result).containsExactly("testDefaultDB"); } private static TableCatalog sessionCatalogWithViews() { From e5412960aa3657e6f0bbc03b688fa28496cc620d Mon Sep 17 00:00:00 2001 From: Dzeri96 Date: Wed, 7 Oct 2026 00:04:02 +0200 Subject: [PATCH 9/9] Issue 14424: Completely removed tests for defaultDatabase calls in SparkSessionCatalog --- .../spark/TestSparkSessionCatalog.java | 24 ------------------- .../spark/TestSparkSessionCatalog.java | 24 ------------------- .../spark/TestSparkSessionCatalog.java | 24 ------------------- .../spark/TestSparkSessionCatalog.java | 24 ------------------- 4 files changed, 96 deletions(-) diff --git a/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java index 5019100e6041..d6d5a237d2b0 100644 --- a/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java +++ b/spark/v3.5/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java @@ -32,14 +32,10 @@ import org.apache.spark.sql.connector.catalog.SupportsNamespaces; import org.apache.spark.sql.connector.catalog.TableCatalog; import org.apache.spark.sql.connector.catalog.ViewCatalog; -import org.apache.spark.sql.execution.datasources.v2.V2SessionCatalog; -import org.apache.spark.sql.internal.SQLConf; -import org.apache.spark.sql.internal.StaticSQLConf; import org.apache.spark.sql.util.CaseInsensitiveStringMap; import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; -import scala.Function0; public class TestSparkSessionCatalog extends TestBase { private final String envHmsUriKey = "spark.hadoop." + METASTOREURIS.varname; @@ -133,26 +129,6 @@ public void listViewsReturnsSessionCatalogViews() throws NoSuchNamespaceExceptio assertThat(catalog.listViews("default")).containsExactly(viewIdent); } - @Test - public void sessionCatalogPicksUpDefaultDatabaseConfig() { - SQLConf sqlConf = new SQLConf(); - sqlConf.setConf(StaticSQLConf.CATALOG_DEFAULT_DATABASE(), "testDefaultDB"); - String[] result = - SQLConf.withExistingConf( - sqlConf, - (Function0) - () -> { - var v1SessionCatalogMock = - mock(org.apache.spark.sql.catalyst.catalog.SessionCatalog.class); - V2SessionCatalog v2SessionCatalog = new V2SessionCatalog(v1SessionCatalogMock); - - SparkSessionCatalog icebergSessionCatalog = new NoViewCatalog<>(); - icebergSessionCatalog.setDelegateCatalog(v2SessionCatalog); - return icebergSessionCatalog.defaultNamespace(); - }); - assertThat(result).containsExactly("testDefaultDB"); - } - private static class NoViewCatalog< T extends TableCatalog & FunctionCatalog & SupportsNamespaces & ViewCatalog> extends SparkSessionCatalog { diff --git a/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java b/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java index 5019100e6041..d6d5a237d2b0 100644 --- a/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java +++ b/spark/v4.0/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java @@ -32,14 +32,10 @@ import org.apache.spark.sql.connector.catalog.SupportsNamespaces; import org.apache.spark.sql.connector.catalog.TableCatalog; import org.apache.spark.sql.connector.catalog.ViewCatalog; -import org.apache.spark.sql.execution.datasources.v2.V2SessionCatalog; -import org.apache.spark.sql.internal.SQLConf; -import org.apache.spark.sql.internal.StaticSQLConf; import org.apache.spark.sql.util.CaseInsensitiveStringMap; import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; -import scala.Function0; public class TestSparkSessionCatalog extends TestBase { private final String envHmsUriKey = "spark.hadoop." + METASTOREURIS.varname; @@ -133,26 +129,6 @@ public void listViewsReturnsSessionCatalogViews() throws NoSuchNamespaceExceptio assertThat(catalog.listViews("default")).containsExactly(viewIdent); } - @Test - public void sessionCatalogPicksUpDefaultDatabaseConfig() { - SQLConf sqlConf = new SQLConf(); - sqlConf.setConf(StaticSQLConf.CATALOG_DEFAULT_DATABASE(), "testDefaultDB"); - String[] result = - SQLConf.withExistingConf( - sqlConf, - (Function0) - () -> { - var v1SessionCatalogMock = - mock(org.apache.spark.sql.catalyst.catalog.SessionCatalog.class); - V2SessionCatalog v2SessionCatalog = new V2SessionCatalog(v1SessionCatalogMock); - - SparkSessionCatalog icebergSessionCatalog = new NoViewCatalog<>(); - icebergSessionCatalog.setDelegateCatalog(v2SessionCatalog); - return icebergSessionCatalog.defaultNamespace(); - }); - assertThat(result).containsExactly("testDefaultDB"); - } - private static class NoViewCatalog< T extends TableCatalog & FunctionCatalog & SupportsNamespaces & ViewCatalog> extends SparkSessionCatalog { diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java index 5019100e6041..d6d5a237d2b0 100644 --- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java +++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java @@ -32,14 +32,10 @@ import org.apache.spark.sql.connector.catalog.SupportsNamespaces; import org.apache.spark.sql.connector.catalog.TableCatalog; import org.apache.spark.sql.connector.catalog.ViewCatalog; -import org.apache.spark.sql.execution.datasources.v2.V2SessionCatalog; -import org.apache.spark.sql.internal.SQLConf; -import org.apache.spark.sql.internal.StaticSQLConf; import org.apache.spark.sql.util.CaseInsensitiveStringMap; import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; -import scala.Function0; public class TestSparkSessionCatalog extends TestBase { private final String envHmsUriKey = "spark.hadoop." + METASTOREURIS.varname; @@ -133,26 +129,6 @@ public void listViewsReturnsSessionCatalogViews() throws NoSuchNamespaceExceptio assertThat(catalog.listViews("default")).containsExactly(viewIdent); } - @Test - public void sessionCatalogPicksUpDefaultDatabaseConfig() { - SQLConf sqlConf = new SQLConf(); - sqlConf.setConf(StaticSQLConf.CATALOG_DEFAULT_DATABASE(), "testDefaultDB"); - String[] result = - SQLConf.withExistingConf( - sqlConf, - (Function0) - () -> { - var v1SessionCatalogMock = - mock(org.apache.spark.sql.catalyst.catalog.SessionCatalog.class); - V2SessionCatalog v2SessionCatalog = new V2SessionCatalog(v1SessionCatalogMock); - - SparkSessionCatalog icebergSessionCatalog = new NoViewCatalog<>(); - icebergSessionCatalog.setDelegateCatalog(v2SessionCatalog); - return icebergSessionCatalog.defaultNamespace(); - }); - assertThat(result).containsExactly("testDefaultDB"); - } - private static class NoViewCatalog< T extends TableCatalog & FunctionCatalog & SupportsNamespaces & ViewCatalog> extends SparkSessionCatalog { diff --git a/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java b/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java index 5e643915ccd6..46897c2f9a44 100644 --- a/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java +++ b/spark/v4.2/spark/src/test/java/org/apache/iceberg/spark/TestSparkSessionCatalog.java @@ -48,15 +48,11 @@ import org.apache.spark.sql.connector.catalog.ViewCatalog; import org.apache.spark.sql.connector.expressions.Transform; import org.apache.spark.sql.connector.metric.CustomTaskMetric; -import org.apache.spark.sql.execution.datasources.v2.V2SessionCatalog; -import org.apache.spark.sql.internal.SQLConf; -import org.apache.spark.sql.internal.StaticSQLConf; import org.apache.spark.sql.types.StructType; import org.apache.spark.sql.util.CaseInsensitiveStringMap; import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; -import scala.Function0; public class TestSparkSessionCatalog extends TestBase { private final String envHmsUriKey = "spark.hadoop." + METASTOREURIS.varname; @@ -496,26 +492,6 @@ public void createOrReplaceViewUsesExistingSessionCatalogView() verify(icebergViewCatalog, never()).createOrReplaceView(viewIdent, replacement); } - @Test - public void sessionCatalogPicksUpDefaultDatabaseConfig() { - SQLConf sqlConf = new SQLConf(); - sqlConf.setConf(StaticSQLConf.CATALOG_DEFAULT_DATABASE(), "testDefaultDB"); - String[] result = - SQLConf.withExistingConf( - sqlConf, - (Function0) - () -> { - var v1SessionCatalogMock = - mock(org.apache.spark.sql.catalyst.catalog.SessionCatalog.class); - V2SessionCatalog v2SessionCatalog = new V2SessionCatalog(v1SessionCatalogMock); - - SparkSessionCatalog icebergSessionCatalog = new NoViewCatalog<>(); - icebergSessionCatalog.setDelegateCatalog(v2SessionCatalog); - return icebergSessionCatalog.defaultNamespace(); - }); - assertThat(result).containsExactly("testDefaultDB"); - } - private static TableCatalog sessionCatalogWithViews() { return mock( TableCatalog.class,