Skip to content

Add support for SampleExec (DataFrame.sample, TABLESAMPLE, randomSplit) #5109

Description

@andygrove

What is the problem the feature request solves?

SampleExec is not supported, so any query using DataFrame.sample, SQL TABLESAMPLE, or DataFrame.randomSplit falls back to Spark and breaks up the native execution block.

Spark lowers all three of these to a single physical operator:

case class SampleExec(lowerBound: Double, upperBound: Double,
                      withReplacement: Boolean, seed: Long, child: SparkPlan)

There are two execution paths (sql/core/src/main/scala/org/apache/spark/sql/execution/basicPhysicalOperators.scala):

Without replacement (the common case; df.sample(0.1) gives lowerBound=0.0, upperBound=0.1): BernoulliCellSampler(lb, ub) seeded with seed + partitionIndex, drawing one XORShiftRandom.nextDouble() per input row and keeping the row when x >= lb && x < ub.

With replacement: PoissonSampler(ub - lb, useGapSamplingIfPossible = false), which is commons-math3 PoissonDistribution.sample() per row (emitting each row 0..n times), reseeded per partition from new java.util.Random(seed).nextLong() drawn partitionIndex + 1 times.

SampleExec is unchanged across Spark 3.4 through 4.2, so no shims are needed.

Describe the potential solution

Add SampleExec as a fully native operator (nativeExecs), per docs/source/contributor-guide/adding_a_new_operator.md:

  • native/proto/src/proto/operator.proto: new Sample message (lower_bound, upper_bound, with_replacement, seed) plus a slot in oneof op_struct
  • spark/src/main/scala/org/apache/comet/serde/operator/CometSampleExec.scala: CometOperatorSerde[SampleExec] and a CometSampleExec case class extending CometUnaryExec, passing through outputPartitioning and outputOrdering
  • spark/src/main/scala/org/apache/comet/rules/CometExecRule.scala: register in nativeExecs
  • common/src/main/scala/org/apache/comet/CometConf.scala: spark.comet.exec.sample.enabled
  • native/core/src/execution/operators/sample.rs: the ExecutionPlan implementation
  • native/core/src/execution/planner.rs: match arm for the new OpStruct

RNG fidelity

Comet compares results exactly against Spark, so it is not enough to sample the right fraction of rows: the operator has to select the same rows.

The non-replacement path is straightforward. XorShiftRandom in native/spark-expr/src/nondetermenistic_funcs/rand.rs is already a bit-exact port of Spark's XORShiftRandom, including the murmur3 hashSeed applied by setSeed, and planner.partition() supplies the partition index (see RandBuilder in native/core/src/execution/expressions/random.rs, which already computes seed.wrapping_add(planner.partition().into())). The operator draws one next_f64() per row into a BooleanArray, filters the batch, and carries RNG state across batches within the stream.

The with-replacement path is a much larger lift. Bit-exact behaviour needs ports of commons-math3 Well19937c and PoissonDistribution.sample() (the Devroye/Kemp hybrid, which consumes Gaussian and exponential draws), plus a java.util.Random LCG for the per-partition seed derivation. There is precedent for this kind of port (SparkMersenneTwister in shuffle.rs, added for shuffle and uuid), but it is a self-contained piece of work.

Suggested phasing

  1. Support withReplacement = false only, returning Unsupported from getSupportLevel for the with-replacement case. Tests over Parquet scans with a fixed seed comparing exact results, plus a randomSplit test to cover a non-zero lowerBound.
  2. Follow-up issue for the with-replacement path (Well19937c and Poisson ports).

Caveat on row order

The RNG stream is consumed one draw per row, in order, within a partition. That is fine when SampleExec sits above a scan, filter, or project. If it sits above an operator where Comet emits rows in a different order than Spark (for example a join or aggregate), an identical RNG stream still selects a different set of rows. This is inherent to the design and already applies to rand(), but it means exact-match tests should be written over deterministic inputs, and other cases should compare sorted results or row counts only.

Additional context

None.

Metadata

Metadata

Assignees

No one assigned

    Type

    No type

    Projects

    No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions