Skip to content
Open
22 changes: 11 additions & 11 deletions docs/docs/spark-configuration.md
Original file line number Diff line number Diff line change
Expand Up @@ -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).

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -108,7 +108,9 @@
* <li><code>io-impl</code> - a custom {@link org.apache.iceberg.io.FileIO} implementation to use
* <li><code>metrics-reporter-impl</code> - a custom {@link
* org.apache.iceberg.metrics.MetricsReporter} implementation to use
* <li><code>default-namespace</code> - a namespace to use as the default
* <li><code>default-namespace</code> - <b>DEPRECATED:</b> use <code>defaultDatabase</code>
* instead
* <li><code>defaultDatabase</code> - a namespace/database to use as the default
* <li><code>cache-enabled</code> - whether to enable catalog cache
* <li><code>cache.case-sensitive</code> - whether the catalog cache should compare table
* identifiers in a case sensitive way
Expand Down Expand Up @@ -813,7 +815,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]);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -92,7 +91,7 @@ protected TableCatalog buildSparkCatalog(String name, CaseInsensitiveStringMap o

@Override
public String[] defaultNamespace() {
return DEFAULT_NAMESPACE;
return getSessionCatalog().defaultNamespace();
}

@Override
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,87 @@
/*
* 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<String, Object> 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<String, String> overrides, Consumer<SparkSession> 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);
}
}
}
Original file line number Diff line number Diff line change
@@ -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() {
Map<String, String> 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() {
Map<String, String> 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() {
Map<String, String> 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");
});
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -109,7 +109,9 @@
* <li><code>io-impl</code> - a custom {@link org.apache.iceberg.io.FileIO} implementation to use
* <li><code>metrics-reporter-impl</code> - a custom {@link
* org.apache.iceberg.metrics.MetricsReporter} implementation to use
* <li><code>default-namespace</code> - a namespace to use as the default
* <li><code>default-namespace</code> - <b>DEPRECATED:</b> use <code>defaultDatabase</code>
* instead
* <li><code>defaultDatabase</code> - a namespace/database to use as the default
* <li><code>cache-enabled</code> - whether to enable catalog cache
* <li><code>cache.case-sensitive</code> - whether the catalog cache should compare table
* identifiers in a case sensitive way
Expand Down Expand Up @@ -813,7 +815,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]);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -93,7 +92,7 @@ protected TableCatalog buildSparkCatalog(String name, CaseInsensitiveStringMap o

@Override
public String[] defaultNamespace() {
return DEFAULT_NAMESPACE;
return getSessionCatalog().defaultNamespace();
}

@Override
Expand Down
Loading
Loading