| diff --git a/gradle/libs.versions.toml b/gradle/libs.versions.toml |
| index db659dc06b..78f5e3440f 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.0.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..c5e31185a9 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,14 @@ 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.memory.offHeap.enabled", "true") |
| + .config("spark.memory.offHeap.size", "10g") |
| .config(TestBase.DISABLE_UI) |
| .enableHiveSupport() |
| .getOrCreate(); |
| 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..350b4a7561 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,14 @@ 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.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..402561a241 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,14 @@ 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.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..004db59f1f 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,14 @@ 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.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..b3fde306c8 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,14 @@ 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.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..62daec8291 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,17 @@ 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.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..325dcb95d1 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,14 @@ 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.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..1292883f55 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,14 @@ 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.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..92166f7b48 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,14 @@ 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.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..0a9a958b2c 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,14 @@ 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.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..3b73dcc014 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,14 @@ 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.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..d696e85139 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,14 @@ 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.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..fae1357a40 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,14 @@ 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.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..043d6596c0 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,18 @@ 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.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..ed47a307c3 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,18 @@ 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.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..abab7798bc 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,14 @@ 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.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..62a7bdc04d 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,14 @@ 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.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..26304e8894 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,14 @@ 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.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..0e83b02f45 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,14 @@ 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.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..5d41bbeb5e 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,14 @@ 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.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..38b6a1f6a4 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,14 @@ 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.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..eba566120e 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,14 @@ 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.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..0dac36653c 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,14 @@ 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.memory.offHeap.enabled", "true") |
| + .config("spark.memory.offHeap.size", "10g") |
| .config(TestBase.DISABLE_UI) |
| .getOrCreate(); |
| } |
| @@ -149,7 +157,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 +227,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 +272,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 +324,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 +366,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 +405,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 +469,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 +822,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..e220f0dd31 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,14 @@ 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.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..108d31b965 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,14 @@ 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.memory.offHeap.enabled", "true") |
| + .config("spark.memory.offHeap.size", "10g") |
| .config(TestBase.DISABLE_UI) |
| .enableHiveSupport() |
| .getOrCreate(); |
| @@ -210,7 +218,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..a161d541cf 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,14 @@ 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.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..c9299d1ebf 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,14 @@ 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.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..748b0a8327 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,18 @@ 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.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..9b2864bb33 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,14 @@ 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.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..f48b4f674a 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,14 @@ 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.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..820f9fc03b 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,14 @@ 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.memory.offHeap.enabled", "true") |
| + .config("spark.memory.offHeap.size", "10g") |
| .enableHiveSupport() |
| .getOrCreate(); |
| |