| 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(); |
| |