Skip to content

Spark 3.4: global aggregate over a single-split Iceberg scan returns no row when DPP prunes the split #6667

Description

@andygrove

Describe the bug

On Spark 3.4, a global aggregate over a native Iceberg scan returns no row when the table has a single input split and dynamic partition pruning removes that split. Spark returns one row.

Spark 3.4's DataSourceV2ScanExecBase.outputPartitioning reports SinglePartition when the scan has one input partition, so EnsureRequirements leaves out the exchange between the partial and final aggregate. When DPP then filters out the only split, BatchScanExec.inputRDD returns sparkContext.parallelize(empty, 1) so the plan still runs on one partition. CometIcebergNativeScan.serializePartitions treats that ParallelCollectionRDD as "DPP pruned everything" and emits no per-partition data, so the native scan's RDD has 0 partitions and the final aggregate never runs.

It happens with AQE on and off. Spark 3.5 and later are not affected, because DataSourceV2ScanExecBase no longer reports SinglePartition there and the plan gets an exchange. The empty-RDD branch came in with #4215, so 1.0 and 1.1 are affected too.

Steps to reproduce

Spark 3.4, spark.comet.scan.icebergNative.enabled=true, Iceberg hadoop catalog probe_cat:

val dimPath = "/tmp/probe_dim"
spark.sql("CREATE TABLE probe_cat.db.fact (amount INT, store_id INT) USING iceberg PARTITIONED BY (store_id)")
// every row in one partition, so the table is a single data file and a single split
spark.range(100).selectExpr("cast(id as int) as amount", "5 as store_id")
  .coalesce(1).writeTo("probe_cat.db.fact").append()
spark.range(1).selectExpr("99 as store_id", "'X' as country").write.parquet(dimPath)
spark.read.parquet(dimPath).createOrReplaceTempView("probe_dim")

spark.sql("""SELECT /*+ BROADCAST(d) */ count(*) AS c, sum(f.amount) AS s
  FROM probe_cat.db.fact f JOIN probe_dim d ON f.store_id = d.store_id
  WHERE d.country = 'X'""").collect()

Expected behavior

[0,null], as Spark returns. Comet returns no rows. The final plan has HashAggregate(Final) directly over HashAggregate(Partial), with no exchange, above the CometIcebergNativeScan, whose RDD has 0 partitions.

Additional context

Found while reviewing #6268, which does not change this (the scan's execution partition count is the same before and after). Emitting one empty partition in the ParallelCollectionRDD case, to match the partition count Spark's inputRDD has, looks like the fix.

Activity

  1. LinSimon-901101 commented on Oct 8, 2026

    @LinSimon-901101
    Contributor

    take

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions