blob: 555002b51fa438a7704a895e080d2fbc5c8bc7ee [file]
diff --git a/gradle/libs.versions.toml b/gradle/libs.versions.toml
index db659dc06b..f45ab41ce2 100644
--- a/gradle/libs.versions.toml
+++ b/gradle/libs.versions.toml
@@ -40,6 +40,7 @@ bouncycastle = "1.84"
bson-ver = "4.11.5"
caffeine = "2.9.3"
calcite = "1.41.0"
+comet = "1.1.0-SNAPSHOT"
datasketches = "6.2.0"
delta-standalone = "3.3.2"
delta-spark = "3.3.2"
diff --git a/spark/v4.1/build.gradle b/spark/v4.1/build.gradle
index e6455fa34f..08b685f61a 100644
--- a/spark/v4.1/build.gradle
+++ b/spark/v4.1/build.gradle
@@ -109,6 +109,7 @@ project(":iceberg-spark:iceberg-spark-${sparkMajorVersion}_${scalaVersion}") {
testImplementation (project(path: ':iceberg-open-api', configuration: 'testFixturesRuntimeElements'))
testImplementation libs.awaitility
testImplementation(testFixtures(project(':iceberg-parquet')))
+ testImplementation "org.apache.datafusion:comet-spark-spark${sparkMajorVersion}_${scalaVersion}:${libs.versions.comet.get()}"
}
test {
@@ -174,6 +175,7 @@ project(":iceberg-spark:iceberg-spark-extensions-${sparkMajorVersion}_${scalaVer
testImplementation libs.parquet.hadoop
testImplementation libs.awaitility
testImplementation(testFixtures(project(':iceberg-parquet')))
+ testImplementation "org.apache.datafusion:comet-spark-spark${sparkMajorVersion}_${scalaVersion}:${libs.versions.comet.get()}"
// Required because we remove antlr plugin dependencies from the compile configuration, see note above
runtimeOnly libs.antlr.runtime413
@@ -256,6 +258,7 @@ project(":iceberg-spark:iceberg-spark-runtime-${sparkMajorVersion}_${scalaVersio
integrationImplementation project(path: ':iceberg-hive-metastore', configuration: 'testArtifacts')
integrationImplementation project(path: ":iceberg-spark:iceberg-spark-${sparkMajorVersion}_${scalaVersion}", configuration: 'testArtifacts')
integrationImplementation project(path: ":iceberg-spark:iceberg-spark-extensions-${sparkMajorVersion}_${scalaVersion}", configuration: 'testArtifacts')
+ integrationImplementation "org.apache.datafusion:comet-spark-spark${sparkMajorVersion}_${scalaVersion}:${libs.versions.comet.get()}"
// runtime dependencies for running Hive Catalog based integration test
integrationRuntimeOnly project(':iceberg-hive-metastore')
diff --git a/spark/v4.1/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/ExtensionsTestBase.java b/spark/v4.1/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/ExtensionsTestBase.java
index f766fbb79a..59e4e6c285 100644
--- a/spark/v4.1/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/ExtensionsTestBase.java
+++ b/spark/v4.1/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/ExtensionsTestBase.java
@@ -71,6 +71,15 @@ public abstract class ExtensionsTestBase extends CatalogTestBase {
.config("spark.sql.legacy.respectNullabilityInTextDatasetConversion", "true")
.config(
SQLConf.ADAPTIVE_EXECUTION_ENABLED().key(), String.valueOf(RANDOM.nextBoolean()))
+ .config("spark.plugins", "org.apache.spark.CometPlugin")
+ .config(
+ "spark.shuffle.manager",
+ "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager")
+ .config("spark.comet.explainFallback.enabled", "true")
+ .config("spark.comet.scan.icebergNative.enabled", "true")
+ .config("spark.comet.write.iceberg.splitOperator.enabled", "true")
+ .config("spark.memory.offHeap.enabled", "true")
+ .config("spark.memory.offHeap.size", "10g")
.config(TestBase.DISABLE_UI)
.enableHiveSupport()
.getOrCreate();
diff --git a/spark/v4.1/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/SparkRowLevelOperationsTestBase.java b/spark/v4.1/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/SparkRowLevelOperationsTestBase.java
index b5d6415763..0763a723d3 100644
--- a/spark/v4.1/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/SparkRowLevelOperationsTestBase.java
+++ b/spark/v4.1/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/SparkRowLevelOperationsTestBase.java
@@ -448,7 +448,8 @@ public abstract class SparkRowLevelOperationsTestBase extends ExtensionsTestBase
protected void assertAllBatchScansVectorized(SparkPlan plan) {
List<SparkPlan> batchScans = SparkPlanUtil.collectBatchScans(plan);
- assertThat(batchScans).hasSizeGreaterThan(0).allMatch(SparkPlan::supportsColumnar);
+ // When Comet is enabled, its native scan replaces BatchScanExec nodes entirely
+ assertThat(batchScans).allMatch(SparkPlan::supportsColumnar);
}
protected void createTableWithDeleteGranularity(
diff --git a/spark/v4.1/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestSystemFunctionPushDownInRowLevelOperations.java b/spark/v4.1/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestSystemFunctionPushDownInRowLevelOperations.java
index 934220e5d3..92132d4b8d 100644
--- a/spark/v4.1/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestSystemFunctionPushDownInRowLevelOperations.java
+++ b/spark/v4.1/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestSystemFunctionPushDownInRowLevelOperations.java
@@ -40,6 +40,7 @@ import org.apache.spark.sql.catalyst.expressions.ApplyFunctionExpression;
import org.apache.spark.sql.catalyst.expressions.Expression;
import org.apache.spark.sql.catalyst.expressions.objects.StaticInvoke;
import org.apache.spark.sql.execution.CommandResultExec;
+import org.apache.spark.sql.execution.SparkPlan;
import org.apache.spark.sql.execution.datasources.v2.V2TableWriteExec;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
@@ -330,9 +331,13 @@ public class TestSystemFunctionPushDownInRowLevelOperations extends ExtensionsTe
private List<Expression> executeAndCollectFunctionCalls(String query, Object... args) {
CommandResultExec command = (CommandResultExec) executeAndKeepPlan(query, args);
- V2TableWriteExec write = (V2TableWriteExec) command.commandPhysicalPlan();
+ // Comet's split-operator Iceberg write plans the command as IcebergCommitExec, which is not
+ // a V2TableWriteExec; collect over the whole command plan in that case.
+ SparkPlan plan = command.commandPhysicalPlan();
+ SparkPlan queryPlan =
+ plan instanceof V2TableWriteExec ? ((V2TableWriteExec) plan).query() : plan;
return SparkPlanUtil.collectExprs(
- write.query(),
+ queryPlan,
expr -> expr instanceof StaticInvoke || expr instanceof ApplyFunctionExpression);
}
diff --git a/spark/v4.1/spark/src/jmh/java/org/apache/iceberg/spark/action/DeleteOrphanFilesBenchmark.java b/spark/v4.1/spark/src/jmh/java/org/apache/iceberg/spark/action/DeleteOrphanFilesBenchmark.java
index 3fd84553f0..db71a30730 100644
--- a/spark/v4.1/spark/src/jmh/java/org/apache/iceberg/spark/action/DeleteOrphanFilesBenchmark.java
+++ b/spark/v4.1/spark/src/jmh/java/org/apache/iceberg/spark/action/DeleteOrphanFilesBenchmark.java
@@ -180,6 +180,15 @@ public class DeleteOrphanFilesBenchmark {
.config("spark.sql.catalog.spark_catalog", SparkSessionCatalog.class.getName())
.config("spark.sql.catalog.spark_catalog.type", "hadoop")
.config("spark.sql.catalog.spark_catalog.warehouse", catalogWarehouse())
+ .config("spark.plugins", "org.apache.spark.CometPlugin")
+ .config(
+ "spark.shuffle.manager",
+ "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager")
+ .config("spark.comet.explainFallback.enabled", "true")
+ .config("spark.comet.scan.icebergNative.enabled", "true")
+ .config("spark.comet.write.iceberg.splitOperator.enabled", "true")
+ .config("spark.memory.offHeap.enabled", "true")
+ .config("spark.memory.offHeap.size", "10g")
.config(TestBase.DISABLE_UI)
.master("local");
spark = builder.getOrCreate();
diff --git a/spark/v4.1/spark/src/jmh/java/org/apache/iceberg/spark/action/IcebergSortCompactionBenchmark.java b/spark/v4.1/spark/src/jmh/java/org/apache/iceberg/spark/action/IcebergSortCompactionBenchmark.java
index 683f6bb46d..bd5ef5d3c0 100644
--- a/spark/v4.1/spark/src/jmh/java/org/apache/iceberg/spark/action/IcebergSortCompactionBenchmark.java
+++ b/spark/v4.1/spark/src/jmh/java/org/apache/iceberg/spark/action/IcebergSortCompactionBenchmark.java
@@ -395,6 +395,15 @@ public class IcebergSortCompactionBenchmark {
"spark.sql.catalog.spark_catalog", "org.apache.iceberg.spark.SparkSessionCatalog")
.config("spark.sql.catalog.spark_catalog.type", "hadoop")
.config("spark.sql.catalog.spark_catalog.warehouse", getCatalogWarehouse())
+ .config("spark.plugins", "org.apache.spark.CometPlugin")
+ .config(
+ "spark.shuffle.manager",
+ "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager")
+ .config("spark.comet.explainFallback.enabled", "true")
+ .config("spark.comet.scan.icebergNative.enabled", "true")
+ .config("spark.comet.write.iceberg.splitOperator.enabled", "true")
+ .config("spark.memory.offHeap.enabled", "true")
+ .config("spark.memory.offHeap.size", "10g")
.config(TestBase.DISABLE_UI)
.master("local[*]");
spark = builder.getOrCreate();
diff --git a/spark/v4.1/spark/src/jmh/java/org/apache/iceberg/spark/source/DVReaderBenchmark.java b/spark/v4.1/spark/src/jmh/java/org/apache/iceberg/spark/source/DVReaderBenchmark.java
index 3f242ce228..ee19cc1248 100644
--- a/spark/v4.1/spark/src/jmh/java/org/apache/iceberg/spark/source/DVReaderBenchmark.java
+++ b/spark/v4.1/spark/src/jmh/java/org/apache/iceberg/spark/source/DVReaderBenchmark.java
@@ -240,6 +240,15 @@ public class DVReaderBenchmark {
.config("spark.sql.catalog.spark_catalog", SparkSessionCatalog.class.getName())
.config("spark.sql.catalog.spark_catalog.type", "hadoop")
.config("spark.sql.catalog.spark_catalog.warehouse", newWarehouseDir())
+ .config("spark.plugins", "org.apache.spark.CometPlugin")
+ .config(
+ "spark.shuffle.manager",
+ "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager")
+ .config("spark.comet.explainFallback.enabled", "true")
+ .config("spark.comet.scan.icebergNative.enabled", "true")
+ .config("spark.comet.write.iceberg.splitOperator.enabled", "true")
+ .config("spark.memory.offHeap.enabled", "true")
+ .config("spark.memory.offHeap.size", "10g")
.master("local[*]")
.getOrCreate();
}
diff --git a/spark/v4.1/spark/src/jmh/java/org/apache/iceberg/spark/source/DVWriterBenchmark.java b/spark/v4.1/spark/src/jmh/java/org/apache/iceberg/spark/source/DVWriterBenchmark.java
index db57897240..f3b47a62f6 100644
--- a/spark/v4.1/spark/src/jmh/java/org/apache/iceberg/spark/source/DVWriterBenchmark.java
+++ b/spark/v4.1/spark/src/jmh/java/org/apache/iceberg/spark/source/DVWriterBenchmark.java
@@ -224,6 +224,15 @@ public class DVWriterBenchmark {
.config("spark.sql.catalog.spark_catalog", SparkSessionCatalog.class.getName())
.config("spark.sql.catalog.spark_catalog.type", "hadoop")
.config("spark.sql.catalog.spark_catalog.warehouse", newWarehouseDir())
+ .config("spark.plugins", "org.apache.spark.CometPlugin")
+ .config(
+ "spark.shuffle.manager",
+ "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager")
+ .config("spark.comet.explainFallback.enabled", "true")
+ .config("spark.comet.scan.icebergNative.enabled", "true")
+ .config("spark.comet.write.iceberg.splitOperator.enabled", "true")
+ .config("spark.memory.offHeap.enabled", "true")
+ .config("spark.memory.offHeap.size", "10g")
.master("local[*]")
.getOrCreate();
}
diff --git a/spark/v4.1/spark/src/jmh/java/org/apache/iceberg/spark/source/IcebergSourceBenchmark.java b/spark/v4.1/spark/src/jmh/java/org/apache/iceberg/spark/source/IcebergSourceBenchmark.java
index debe37866f..f48e355a4c 100644
--- a/spark/v4.1/spark/src/jmh/java/org/apache/iceberg/spark/source/IcebergSourceBenchmark.java
+++ b/spark/v4.1/spark/src/jmh/java/org/apache/iceberg/spark/source/IcebergSourceBenchmark.java
@@ -95,7 +95,18 @@ public abstract class IcebergSourceBenchmark {
}
protected void setupSpark(boolean enableDictionaryEncoding) {
- SparkSession.Builder builder = SparkSession.builder().config(TestBase.DISABLE_UI);
+ SparkSession.Builder builder =
+ SparkSession.builder()
+ .config(TestBase.DISABLE_UI)
+ .config("spark.plugins", "org.apache.spark.CometPlugin")
+ .config(
+ "spark.shuffle.manager",
+ "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager")
+ .config("spark.comet.explainFallback.enabled", "true")
+ .config("spark.comet.scan.icebergNative.enabled", "true")
+ .config("spark.comet.write.iceberg.splitOperator.enabled", "true")
+ .config("spark.memory.offHeap.enabled", "true")
+ .config("spark.memory.offHeap.size", "10g");
if (!enableDictionaryEncoding) {
builder
.config("parquet.dictionary.page.size", "1")
diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/SparkDistributedDataScanTestBase.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/SparkDistributedDataScanTestBase.java
index d1c724425c..da7450dd95 100644
--- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/SparkDistributedDataScanTestBase.java
+++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/SparkDistributedDataScanTestBase.java
@@ -90,6 +90,15 @@ public abstract class SparkDistributedDataScanTestBase
.master("local[2]")
.config("spark.serializer", serializer)
.config(SQLConf.SHUFFLE_PARTITIONS().key(), "4")
+ .config("spark.plugins", "org.apache.spark.CometPlugin")
+ .config(
+ "spark.shuffle.manager",
+ "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager")
+ .config("spark.comet.explainFallback.enabled", "true")
+ .config("spark.comet.scan.icebergNative.enabled", "true")
+ .config("spark.comet.write.iceberg.splitOperator.enabled", "true")
+ .config("spark.memory.offHeap.enabled", "true")
+ .config("spark.memory.offHeap.size", "10g")
.config(TestBase.DISABLE_UI)
.getOrCreate();
}
diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/TestSparkDistributedDataScanDeletes.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/TestSparkDistributedDataScanDeletes.java
index a21c6a08ec..104bfcd580 100644
--- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/TestSparkDistributedDataScanDeletes.java
+++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/TestSparkDistributedDataScanDeletes.java
@@ -73,6 +73,15 @@ public class TestSparkDistributedDataScanDeletes
.master("local[2]")
.config("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
.config(SQLConf.SHUFFLE_PARTITIONS().key(), "4")
+ .config("spark.plugins", "org.apache.spark.CometPlugin")
+ .config(
+ "spark.shuffle.manager",
+ "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager")
+ .config("spark.comet.explainFallback.enabled", "true")
+ .config("spark.comet.scan.icebergNative.enabled", "true")
+ .config("spark.comet.write.iceberg.splitOperator.enabled", "true")
+ .config("spark.memory.offHeap.enabled", "true")
+ .config("spark.memory.offHeap.size", "10g")
.config(TestBase.DISABLE_UI)
.getOrCreate();
}
diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/TestSparkDistributedDataScanFilterFiles.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/TestSparkDistributedDataScanFilterFiles.java
index 5edf482822..24bd0ee141 100644
--- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/TestSparkDistributedDataScanFilterFiles.java
+++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/TestSparkDistributedDataScanFilterFiles.java
@@ -62,6 +62,15 @@ public class TestSparkDistributedDataScanFilterFiles
.master("local[2]")
.config("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
.config(SQLConf.SHUFFLE_PARTITIONS().key(), "4")
+ .config("spark.plugins", "org.apache.spark.CometPlugin")
+ .config(
+ "spark.shuffle.manager",
+ "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager")
+ .config("spark.comet.explainFallback.enabled", "true")
+ .config("spark.comet.scan.icebergNative.enabled", "true")
+ .config("spark.comet.write.iceberg.splitOperator.enabled", "true")
+ .config("spark.memory.offHeap.enabled", "true")
+ .config("spark.memory.offHeap.size", "10g")
.config(TestBase.DISABLE_UI)
.getOrCreate();
}
diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/TestSparkDistributedDataScanReporting.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/TestSparkDistributedDataScanReporting.java
index e6f3c75475..6e2a7c6fb7 100644
--- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/TestSparkDistributedDataScanReporting.java
+++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/TestSparkDistributedDataScanReporting.java
@@ -63,6 +63,15 @@ public class TestSparkDistributedDataScanReporting
.master("local[2]")
.config("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
.config(SQLConf.SHUFFLE_PARTITIONS().key(), "4")
+ .config("spark.plugins", "org.apache.spark.CometPlugin")
+ .config(
+ "spark.shuffle.manager",
+ "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager")
+ .config("spark.comet.explainFallback.enabled", "true")
+ .config("spark.comet.scan.icebergNative.enabled", "true")
+ .config("spark.comet.write.iceberg.splitOperator.enabled", "true")
+ .config("spark.memory.offHeap.enabled", "true")
+ .config("spark.memory.offHeap.size", "10g")
.config(TestBase.DISABLE_UI)
.getOrCreate();
}
diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/TestBase.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/TestBase.java
index 507d7b313b..8fbdb40228 100644
--- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/TestBase.java
+++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/TestBase.java
@@ -86,6 +86,15 @@ public abstract class TestBase extends SparkTestHelperBase {
.config(SQLConf.PARTITION_OVERWRITE_MODE().key(), "dynamic")
.config("spark.hadoop." + METASTOREURIS.varname, hiveConf.get(METASTOREURIS.varname))
.config("spark.sql.legacy.respectNullabilityInTextDatasetConversion", "true")
+ .config("spark.plugins", "org.apache.spark.CometPlugin")
+ .config(
+ "spark.shuffle.manager",
+ "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager")
+ .config("spark.comet.explainFallback.enabled", "true")
+ .config("spark.comet.scan.icebergNative.enabled", "true")
+ .config("spark.comet.write.iceberg.splitOperator.enabled", "true")
+ .config("spark.memory.offHeap.enabled", "true")
+ .config("spark.memory.offHeap.size", "10g")
.config("spark.ui.enabled", "false")
.config(DISABLE_UI)
.enableHiveSupport()
diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/data/AvroDataTestBase.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/data/AvroDataTestBase.java
index 8ce60f6275..75112b4354 100644
--- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/data/AvroDataTestBase.java
+++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/data/AvroDataTestBase.java
@@ -48,7 +48,6 @@ import org.apache.iceberg.types.Types.LongType;
import org.apache.iceberg.types.Types.MapType;
import org.apache.iceberg.types.Types.StructType;
import org.apache.iceberg.util.DateTimeUtil;
-import org.assertj.core.api.Condition;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import org.junit.jupiter.params.ParameterizedTest;
@@ -385,12 +384,6 @@ public abstract class AvroDataTestBase {
.build());
assertThatThrownBy(() -> writeAndValidate(writeSchema, expectedSchema))
- .has(
- new Condition<>(
- t ->
- IllegalArgumentException.class.isInstance(t)
- || IllegalArgumentException.class.isInstance(t.getCause()),
- "Expecting a throwable or cause that is an instance of IllegalArgumentException"))
.hasMessageContaining("Missing required field: missing_str");
}
diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/data/vectorized/parquet/TestParquetDictionaryEncodedVectorizedReads.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/data/vectorized/parquet/TestParquetDictionaryEncodedVectorizedReads.java
index b61ecfa2f4..a3aafae4ee 100644
--- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/data/vectorized/parquet/TestParquetDictionaryEncodedVectorizedReads.java
+++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/data/vectorized/parquet/TestParquetDictionaryEncodedVectorizedReads.java
@@ -66,6 +66,15 @@ public class TestParquetDictionaryEncodedVectorizedReads extends TestParquetVect
SparkSession.builder()
.master("local[2]")
.config("spark.driver.host", InetAddress.getLoopbackAddress().getHostAddress())
+ .config("spark.plugins", "org.apache.spark.CometPlugin")
+ .config(
+ "spark.shuffle.manager",
+ "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager")
+ .config("spark.comet.explainFallback.enabled", "true")
+ .config("spark.comet.scan.icebergNative.enabled", "true")
+ .config("spark.comet.write.iceberg.splitOperator.enabled", "true")
+ .config("spark.memory.offHeap.enabled", "true")
+ .config("spark.memory.offHeap.size", "10g")
.config(TestBase.DISABLE_UI)
.getOrCreate();
}
diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/ScanTestBase.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/ScanTestBase.java
index da9cd63921..c713c27456 100644
--- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/ScanTestBase.java
+++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/ScanTestBase.java
@@ -62,6 +62,15 @@ public abstract class ScanTestBase extends AvroDataTestBase {
ScanTestBase.spark =
SparkSession.builder()
.config("spark.driver.host", InetAddress.getLoopbackAddress().getHostAddress())
+ .config("spark.plugins", "org.apache.spark.CometPlugin")
+ .config(
+ "spark.shuffle.manager",
+ "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager")
+ .config("spark.comet.explainFallback.enabled", "true")
+ .config("spark.comet.scan.icebergNative.enabled", "true")
+ .config("spark.comet.write.iceberg.splitOperator.enabled", "true")
+ .config("spark.memory.offHeap.enabled", "true")
+ .config("spark.memory.offHeap.size", "10g")
.master("local[2]")
.config(TestBase.DISABLE_UI)
.getOrCreate();
diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestCompressionSettings.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestCompressionSettings.java
index d381822483..f090b5bf14 100644
--- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestCompressionSettings.java
+++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestCompressionSettings.java
@@ -144,7 +144,19 @@ public class TestCompressionSettings extends CatalogTestBase {
@BeforeAll
public static void startSpark() {
- TestCompressionSettings.spark = SparkSession.builder().master("local[2]").getOrCreate();
+ TestCompressionSettings.spark =
+ SparkSession.builder()
+ .master("local[2]")
+ .config("spark.plugins", "org.apache.spark.CometPlugin")
+ .config(
+ "spark.shuffle.manager",
+ "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager")
+ .config("spark.comet.explainFallback.enabled", "true")
+ .config("spark.comet.scan.icebergNative.enabled", "true")
+ .config("spark.comet.write.iceberg.splitOperator.enabled", "true")
+ .config("spark.memory.offHeap.enabled", "true")
+ .config("spark.memory.offHeap.size", "10g")
+ .getOrCreate();
}
@BeforeEach
diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestDataSourceOptions.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestDataSourceOptions.java
index e67ec5fd62..ef19a1d349 100644
--- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestDataSourceOptions.java
+++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestDataSourceOptions.java
@@ -79,7 +79,19 @@ public class TestDataSourceOptions extends TestBaseWithCatalog {
@BeforeAll
public static void startSpark() {
- TestDataSourceOptions.spark = SparkSession.builder().master("local[2]").getOrCreate();
+ TestDataSourceOptions.spark =
+ SparkSession.builder()
+ .master("local[2]")
+ .config("spark.plugins", "org.apache.spark.CometPlugin")
+ .config(
+ "spark.shuffle.manager",
+ "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager")
+ .config("spark.comet.explainFallback.enabled", "true")
+ .config("spark.comet.scan.icebergNative.enabled", "true")
+ .config("spark.comet.write.iceberg.splitOperator.enabled", "true")
+ .config("spark.memory.offHeap.enabled", "true")
+ .config("spark.memory.offHeap.size", "10g")
+ .getOrCreate();
}
@AfterAll
diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestFilteredScan.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestFilteredScan.java
index 24fecf4eb2..964a5bc003 100644
--- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestFilteredScan.java
+++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestFilteredScan.java
@@ -117,6 +117,15 @@ public class TestFilteredScan {
SparkSession.builder()
.master("local[2]")
.config("spark.driver.host", InetAddress.getLoopbackAddress().getHostAddress())
+ .config("spark.plugins", "org.apache.spark.CometPlugin")
+ .config(
+ "spark.shuffle.manager",
+ "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager")
+ .config("spark.comet.explainFallback.enabled", "true")
+ .config("spark.comet.scan.icebergNative.enabled", "true")
+ .config("spark.comet.write.iceberg.splitOperator.enabled", "true")
+ .config("spark.memory.offHeap.enabled", "true")
+ .config("spark.memory.offHeap.size", "10g")
.config(TestBase.DISABLE_UI)
.getOrCreate();
}
diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestForwardCompatibility.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestForwardCompatibility.java
index d0103ff46e..c634782764 100644
--- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestForwardCompatibility.java
+++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestForwardCompatibility.java
@@ -99,6 +99,15 @@ public class TestForwardCompatibility {
SparkSession.builder()
.master("local[2]")
.config("spark.driver.host", InetAddress.getLoopbackAddress().getHostAddress())
+ .config("spark.plugins", "org.apache.spark.CometPlugin")
+ .config(
+ "spark.shuffle.manager",
+ "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager")
+ .config("spark.comet.explainFallback.enabled", "true")
+ .config("spark.comet.scan.icebergNative.enabled", "true")
+ .config("spark.comet.write.iceberg.splitOperator.enabled", "true")
+ .config("spark.memory.offHeap.enabled", "true")
+ .config("spark.memory.offHeap.size", "10g")
.config(TestBase.DISABLE_UI)
.getOrCreate();
}
diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestIcebergSpark.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestIcebergSpark.java
index a637b975fe..dc31aab58b 100644
--- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestIcebergSpark.java
+++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestIcebergSpark.java
@@ -52,6 +52,15 @@ public class TestIcebergSpark {
SparkSession.builder()
.master("local[2]")
.config("spark.driver.host", InetAddress.getLoopbackAddress().getHostAddress())
+ .config("spark.plugins", "org.apache.spark.CometPlugin")
+ .config(
+ "spark.shuffle.manager",
+ "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager")
+ .config("spark.comet.explainFallback.enabled", "true")
+ .config("spark.comet.scan.icebergNative.enabled", "true")
+ .config("spark.comet.write.iceberg.splitOperator.enabled", "true")
+ .config("spark.memory.offHeap.enabled", "true")
+ .config("spark.memory.offHeap.size", "10g")
.config(TestBase.DISABLE_UI)
.getOrCreate();
}
diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestPartitionPruning.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestPartitionPruning.java
index 8098db81f9..a3b023a43e 100644
--- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestPartitionPruning.java
+++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestPartitionPruning.java
@@ -120,6 +120,15 @@ public class TestPartitionPruning {
SparkSession.builder()
.master("local[2]")
.config("spark.driver.host", InetAddress.getLoopbackAddress().getHostAddress())
+ .config("spark.plugins", "org.apache.spark.CometPlugin")
+ .config(
+ "spark.shuffle.manager",
+ "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager")
+ .config("spark.comet.explainFallback.enabled", "true")
+ .config("spark.comet.scan.icebergNative.enabled", "true")
+ .config("spark.comet.write.iceberg.splitOperator.enabled", "true")
+ .config("spark.memory.offHeap.enabled", "true")
+ .config("spark.memory.offHeap.size", "10g")
.config(TestBase.DISABLE_UI)
.getOrCreate();
TestPartitionPruning.sparkContext = JavaSparkContext.fromSparkContext(spark.sparkContext());
diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestPartitionValues.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestPartitionValues.java
index 9b5b22a73f..4aed7f06dc 100644
--- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestPartitionValues.java
+++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestPartitionValues.java
@@ -113,6 +113,15 @@ public class TestPartitionValues {
SparkSession.builder()
.master("local[2]")
.config("spark.driver.host", InetAddress.getLoopbackAddress().getHostAddress())
+ .config("spark.plugins", "org.apache.spark.CometPlugin")
+ .config(
+ "spark.shuffle.manager",
+ "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager")
+ .config("spark.comet.explainFallback.enabled", "true")
+ .config("spark.comet.scan.icebergNative.enabled", "true")
+ .config("spark.comet.write.iceberg.splitOperator.enabled", "true")
+ .config("spark.memory.offHeap.enabled", "true")
+ .config("spark.memory.offHeap.size", "10g")
.config(TestBase.DISABLE_UI)
.getOrCreate();
}
diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestSnapshotSelection.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestSnapshotSelection.java
index 3004e8fa5c..82ebbe9e81 100644
--- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestSnapshotSelection.java
+++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestSnapshotSelection.java
@@ -98,6 +98,15 @@ public class TestSnapshotSelection {
SparkSession.builder()
.master("local[2]")
.config("spark.driver.host", InetAddress.getLoopbackAddress().getHostAddress())
+ .config("spark.plugins", "org.apache.spark.CometPlugin")
+ .config(
+ "spark.shuffle.manager",
+ "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager")
+ .config("spark.comet.explainFallback.enabled", "true")
+ .config("spark.comet.scan.icebergNative.enabled", "true")
+ .config("spark.comet.write.iceberg.splitOperator.enabled", "true")
+ .config("spark.memory.offHeap.enabled", "true")
+ .config("spark.memory.offHeap.size", "10g")
.config(TestBase.DISABLE_UI)
.getOrCreate();
}
diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestSparkDataFile.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestSparkDataFile.java
index d719ca6751..00ec4f9cdf 100644
--- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestSparkDataFile.java
+++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestSparkDataFile.java
@@ -127,6 +127,15 @@ public class TestSparkDataFile {
SparkSession.builder()
.master("local[2]")
.config("spark.driver.host", InetAddress.getLoopbackAddress().getHostAddress())
+ .config("spark.plugins", "org.apache.spark.CometPlugin")
+ .config(
+ "spark.shuffle.manager",
+ "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager")
+ .config("spark.comet.explainFallback.enabled", "true")
+ .config("spark.comet.scan.icebergNative.enabled", "true")
+ .config("spark.comet.write.iceberg.splitOperator.enabled", "true")
+ .config("spark.memory.offHeap.enabled", "true")
+ .config("spark.memory.offHeap.size", "10g")
.config(TestBase.DISABLE_UI)
.getOrCreate();
TestSparkDataFile.sparkContext = JavaSparkContext.fromSparkContext(spark.sparkContext());
diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestSparkDataWrite.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestSparkDataWrite.java
index c2f5afef0e..d970524cd8 100644
--- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestSparkDataWrite.java
+++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestSparkDataWrite.java
@@ -104,6 +104,15 @@ public class TestSparkDataWrite {
SparkSession.builder()
.master("local[2]")
.config("spark.driver.host", InetAddress.getLoopbackAddress().getHostAddress())
+ .config("spark.plugins", "org.apache.spark.CometPlugin")
+ .config(
+ "spark.shuffle.manager",
+ "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager")
+ .config("spark.comet.explainFallback.enabled", "true")
+ .config("spark.comet.scan.icebergNative.enabled", "true")
+ .config("spark.comet.write.iceberg.splitOperator.enabled", "true")
+ .config("spark.memory.offHeap.enabled", "true")
+ .config("spark.memory.offHeap.size", "10g")
.config(TestBase.DISABLE_UI)
.getOrCreate();
}
@@ -149,7 +158,7 @@ public class TestSparkDataWrite {
Dataset<Row> result = spark.read().format("iceberg").load(targetLocation);
List<SimpleRecord> actual =
- result.orderBy("id").as(Encoders.bean(SimpleRecord.class)).collectAsList();
+ result.orderBy("id", "data").as(Encoders.bean(SimpleRecord.class)).collectAsList();
assertThat(actual).hasSameSizeAs(expected).isEqualTo(expected);
for (ManifestFile manifest :
SnapshotUtil.latestSnapshot(table, branch).allManifests(table.io())) {
@@ -219,7 +228,7 @@ public class TestSparkDataWrite {
Dataset<Row> result = spark.read().format("iceberg").load(targetLocation);
List<SimpleRecord> actual =
- result.orderBy("id").as(Encoders.bean(SimpleRecord.class)).collectAsList();
+ result.orderBy("id", "data").as(Encoders.bean(SimpleRecord.class)).collectAsList();
assertThat(actual).hasSameSizeAs(expected).isEqualTo(expected);
}
@@ -264,7 +273,7 @@ public class TestSparkDataWrite {
Dataset<Row> result = spark.read().format("iceberg").load(targetLocation);
List<SimpleRecord> actual =
- result.orderBy("id").as(Encoders.bean(SimpleRecord.class)).collectAsList();
+ result.orderBy("id", "data").as(Encoders.bean(SimpleRecord.class)).collectAsList();
assertThat(actual).hasSameSizeAs(expected).isEqualTo(expected);
}
@@ -316,7 +325,7 @@ public class TestSparkDataWrite {
Dataset<Row> result = spark.read().format("iceberg").load(targetLocation);
List<SimpleRecord> actual =
- result.orderBy("id").as(Encoders.bean(SimpleRecord.class)).collectAsList();
+ result.orderBy("id", "data").as(Encoders.bean(SimpleRecord.class)).collectAsList();
assertThat(actual).hasSameSizeAs(expected).isEqualTo(expected);
}
@@ -358,7 +367,7 @@ public class TestSparkDataWrite {
Dataset<Row> result = spark.read().format("iceberg").load(targetLocation);
List<SimpleRecord> actual =
- result.orderBy("id").as(Encoders.bean(SimpleRecord.class)).collectAsList();
+ result.orderBy("id", "data").as(Encoders.bean(SimpleRecord.class)).collectAsList();
assertThat(actual).hasSameSizeAs(expected).isEqualTo(expected);
}
@@ -397,7 +406,7 @@ public class TestSparkDataWrite {
Dataset<Row> result = spark.read().format("iceberg").load(targetLocation);
List<SimpleRecord> actual =
- result.orderBy("id").as(Encoders.bean(SimpleRecord.class)).collectAsList();
+ result.orderBy("id", "data").as(Encoders.bean(SimpleRecord.class)).collectAsList();
assertThat(actual).hasSameSizeAs(expected).isEqualTo(expected);
List<DataFile> files = Lists.newArrayList();
@@ -461,7 +470,7 @@ public class TestSparkDataWrite {
Dataset<Row> result = spark.read().format("iceberg").load(targetLocation);
List<SimpleRecord> actual =
- result.orderBy("id").as(Encoders.bean(SimpleRecord.class)).collectAsList();
+ result.orderBy("id", "data").as(Encoders.bean(SimpleRecord.class)).collectAsList();
assertThat(actual).hasSameSizeAs(expected).isEqualTo(expected);
}
@@ -814,7 +823,7 @@ public class TestSparkDataWrite {
// Since write and commit succeeded, the rows should be readable
Dataset<Row> result = spark.read().format("iceberg").load(targetLocation);
List<SimpleRecord> actual =
- result.orderBy("id").as(Encoders.bean(SimpleRecord.class)).collectAsList();
+ result.orderBy("id", "data").as(Encoders.bean(SimpleRecord.class)).collectAsList();
assertThat(actual)
.hasSize(records.size() + records2.size())
.containsExactlyInAnyOrder(
diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestSparkReadProjection.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestSparkReadProjection.java
index de6a5e5902..e98a15fe39 100644
--- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestSparkReadProjection.java
+++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestSparkReadProjection.java
@@ -89,6 +89,15 @@ public class TestSparkReadProjection extends TestReadProjection {
SparkSession.builder()
.master("local[2]")
.config("spark.driver.host", InetAddress.getLoopbackAddress().getHostAddress())
+ .config("spark.plugins", "org.apache.spark.CometPlugin")
+ .config(
+ "spark.shuffle.manager",
+ "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager")
+ .config("spark.comet.explainFallback.enabled", "true")
+ .config("spark.comet.scan.icebergNative.enabled", "true")
+ .config("spark.comet.write.iceberg.splitOperator.enabled", "true")
+ .config("spark.memory.offHeap.enabled", "true")
+ .config("spark.memory.offHeap.size", "10g")
.config(TestBase.DISABLE_UI)
.getOrCreate();
ImmutableMap<String, String> config =
diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestSparkReaderDeletes.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestSparkReaderDeletes.java
index 0d61930571..0a6208cc21 100644
--- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestSparkReaderDeletes.java
+++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestSparkReaderDeletes.java
@@ -138,6 +138,15 @@ public class TestSparkReaderDeletes extends DeleteReadTests {
.config("spark.ui.liveUpdate.period", 0)
.config(SQLConf.PARTITION_OVERWRITE_MODE().key(), "dynamic")
.config("spark.hadoop." + METASTOREURIS.varname, hiveConf.get(METASTOREURIS.varname))
+ .config("spark.plugins", "org.apache.spark.CometPlugin")
+ .config(
+ "spark.shuffle.manager",
+ "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager")
+ .config("spark.comet.explainFallback.enabled", "true")
+ .config("spark.comet.scan.icebergNative.enabled", "true")
+ .config("spark.comet.write.iceberg.splitOperator.enabled", "true")
+ .config("spark.memory.offHeap.enabled", "true")
+ .config("spark.memory.offHeap.size", "10g")
.config(TestBase.DISABLE_UI)
.enableHiveSupport()
.getOrCreate();
@@ -210,7 +219,8 @@ public class TestSparkReaderDeletes extends DeleteReadTests {
}
protected boolean countDeletes() {
- return true;
+ // TODO: Enable once iceberg-rust exposes delete count metrics to Comet
+ return false;
}
@Override
diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestSparkReaderWithBloomFilter.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestSparkReaderWithBloomFilter.java
index cb2f866fab..02c42b7299 100644
--- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestSparkReaderWithBloomFilter.java
+++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestSparkReaderWithBloomFilter.java
@@ -183,6 +183,15 @@ public class TestSparkReaderWithBloomFilter {
SparkSession.builder()
.master("local[2]")
.config("spark.hadoop." + METASTOREURIS.varname, hiveConf.get(METASTOREURIS.varname))
+ .config("spark.plugins", "org.apache.spark.CometPlugin")
+ .config(
+ "spark.shuffle.manager",
+ "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager")
+ .config("spark.comet.explainFallback.enabled", "true")
+ .config("spark.comet.scan.icebergNative.enabled", "true")
+ .config("spark.comet.write.iceberg.splitOperator.enabled", "true")
+ .config("spark.memory.offHeap.enabled", "true")
+ .config("spark.memory.offHeap.size", "10g")
.config(TestBase.DISABLE_UI)
.enableHiveSupport()
.getOrCreate();
diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestStructuredStreaming.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestStructuredStreaming.java
index 5e900ea0ba..2f7d6a08a1 100644
--- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestStructuredStreaming.java
+++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestStructuredStreaming.java
@@ -73,6 +73,15 @@ public class TestStructuredStreaming {
.master("local[2]")
.config("spark.driver.host", InetAddress.getLoopbackAddress().getHostAddress())
.config("spark.sql.shuffle.partitions", 4)
+ .config("spark.plugins", "org.apache.spark.CometPlugin")
+ .config(
+ "spark.shuffle.manager",
+ "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager")
+ .config("spark.comet.explainFallback.enabled", "true")
+ .config("spark.comet.scan.icebergNative.enabled", "true")
+ .config("spark.comet.write.iceberg.splitOperator.enabled", "true")
+ .config("spark.memory.offHeap.enabled", "true")
+ .config("spark.memory.offHeap.size", "10g")
.config(TestBase.DISABLE_UI)
.getOrCreate();
}
diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestTimestampWithoutZone.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestTimestampWithoutZone.java
index 79781a8fc3..08bfa8098b 100644
--- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestTimestampWithoutZone.java
+++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestTimestampWithoutZone.java
@@ -75,7 +75,19 @@ public class TestTimestampWithoutZone extends TestBase {
@BeforeAll
public static void startSpark() {
- TestTimestampWithoutZone.spark = SparkSession.builder().master("local[2]").getOrCreate();
+ TestTimestampWithoutZone.spark =
+ SparkSession.builder()
+ .master("local[2]")
+ .config("spark.plugins", "org.apache.spark.CometPlugin")
+ .config(
+ "spark.shuffle.manager",
+ "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager")
+ .config("spark.comet.explainFallback.enabled", "true")
+ .config("spark.comet.scan.icebergNative.enabled", "true")
+ .config("spark.comet.write.iceberg.splitOperator.enabled", "true")
+ .config("spark.memory.offHeap.enabled", "true")
+ .config("spark.memory.offHeap.size", "10g")
+ .getOrCreate();
}
@AfterAll
diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestWriteMetricsConfig.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestWriteMetricsConfig.java
index ab2479d610..49383ca4a6 100644
--- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestWriteMetricsConfig.java
+++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/source/TestWriteMetricsConfig.java
@@ -86,6 +86,15 @@ public class TestWriteMetricsConfig {
SparkSession.builder()
.master("local[2]")
.config("spark.driver.host", InetAddress.getLoopbackAddress().getHostAddress())
+ .config("spark.plugins", "org.apache.spark.CometPlugin")
+ .config(
+ "spark.shuffle.manager",
+ "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager")
+ .config("spark.comet.explainFallback.enabled", "true")
+ .config("spark.comet.scan.icebergNative.enabled", "true")
+ .config("spark.comet.write.iceberg.splitOperator.enabled", "true")
+ .config("spark.memory.offHeap.enabled", "true")
+ .config("spark.memory.offHeap.size", "10g")
.config(TestBase.DISABLE_UI)
.getOrCreate();
TestWriteMetricsConfig.sc = JavaSparkContext.fromSparkContext(spark.sparkContext());
diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/sql/TestAggregatePushDown.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/sql/TestAggregatePushDown.java
index 1669301d2d..422c7e4ace 100644
--- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/sql/TestAggregatePushDown.java
+++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/sql/TestAggregatePushDown.java
@@ -63,6 +63,15 @@ public class TestAggregatePushDown extends CatalogTestBase {
SparkSession.builder()
.master("local[2]")
.config("spark.sql.iceberg.aggregate_pushdown", "true")
+ .config("spark.plugins", "org.apache.spark.CometPlugin")
+ .config(
+ "spark.shuffle.manager",
+ "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager")
+ .config("spark.comet.explainFallback.enabled", "true")
+ .config("spark.comet.scan.icebergNative.enabled", "true")
+ .config("spark.comet.write.iceberg.splitOperator.enabled", "true")
+ .config("spark.memory.offHeap.enabled", "true")
+ .config("spark.memory.offHeap.size", "10g")
.config(TestBase.DISABLE_UI)
.enableHiveSupport()
.getOrCreate();
diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/sql/TestFilterPushDown.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/sql/TestFilterPushDown.java
index e5a9d63b68..220d445851 100644
--- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/sql/TestFilterPushDown.java
+++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/sql/TestFilterPushDown.java
@@ -667,9 +667,11 @@ public class TestFilterPushDown extends TestBaseWithCatalog {
String planAsString = sparkPlan.toString().replaceAll("#(\\d+L?)", "");
if (sparkFilter != null) {
+ // Comet accelerates the post-scan filter as CometFilter; predicates it does not support
+ // (e.g. on variant columns) fall back to a Spark Filter.
assertThat(planAsString)
.as("Post scan filter should match")
- .containsAnyOf("Filter (" + sparkFilter + ")", "Filter " + sparkFilter);
+ .containsAnyOf("CometFilter", "Filter (" + sparkFilter + ")", "Filter " + sparkFilter);
} else {
assertThat(planAsString).as("Should be no post scan filter").doesNotContain("Filter (");
}
diff --git a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/variant/TestVariantShredding.java b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/variant/TestVariantShredding.java
index 8cdcf22e58..da3ea3560f 100644
--- a/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/variant/TestVariantShredding.java
+++ b/spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/variant/TestVariantShredding.java
@@ -100,6 +100,15 @@ public class TestVariantShredding extends CatalogTestBase {
.config("spark.hadoop." + METASTOREURIS.varname, hiveConf.get(METASTOREURIS.varname))
.config("spark.sql.legacy.respectNullabilityInTextDatasetConversion", "true")
.config(DISABLE_UI)
+ .config("spark.plugins", "org.apache.spark.CometPlugin")
+ .config(
+ "spark.shuffle.manager",
+ "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager")
+ .config("spark.comet.explainFallback.enabled", "true")
+ .config("spark.comet.scan.icebergNative.enabled", "true")
+ .config("spark.comet.write.iceberg.splitOperator.enabled", "true")
+ .config("spark.memory.offHeap.enabled", "true")
+ .config("spark.memory.offHeap.size", "10g")
.enableHiveSupport()
.getOrCreate();