Describe the bug
Since #5692, Comet routes the Invoke and StaticInvoke expressions it doesn't otherwise support through the JVM codegen dispatcher, instead of falling back to Spark. The predicate of a typed Dataset.filter(lambda) is one of these Invokes. So a typed filter over a native operator now runs in CometFilter, through the dispatcher.
The dispatcher reads sliced boolean input wrong (#6288). Once the batches reaching the filter are sliced, it keeps or drops the wrong rows. In 1.0.0 the filter ran in Spark and returned the right rows, so this is a regression in 1.1.0.
Steps to reproduce
// In a suite extending CometTestBase, default configs
import testImplicits._
withSQLConf(
SQLConf.SHUFFLE_PARTITIONS.key -> "1",
SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false") {
withTempPath { dir =>
spark.range(0, 30000).selectExpr("id", "id % 3 = 0 AS flag").write.parquet(dir.getAbsolutePath)
checkSparkAnswer(
spark.read.parquet(dir.getAbsolutePath)
.groupBy("id", "flag").count()
.as[(Long, Boolean, Long)]
.filter(_._2)
.toDF())
}
}
The aggregate emits more groups than spark.comet.batchSize, and its sliced output batches carry boolean arrays with a non-zero offset into the filter.
Expected behavior
Spark and Comet 1.0.0 return 10,000 rows. Comet 1.1.0-rc1 returns 10,001, with rows misclassified in both directions.
Additional context
#6339 fixes #6288 and is on main, which returns the right rows. #6339 is not on branch-1.1. I cherry-picked it onto 1.1.0-rc1: it applies cleanly and the reproducer passes, with the predicate still dispatched. The fix for 1.1.0 is to backport #6339.
Found by the 1.1.0 regression audit (#6399) and tracked in #6402.
Describe the bug
Since #5692, Comet routes the
InvokeandStaticInvokeexpressions it doesn't otherwise support through the JVM codegen dispatcher, instead of falling back to Spark. The predicate of a typedDataset.filter(lambda)is one of theseInvokes. So a typed filter over a native operator now runs inCometFilter, through the dispatcher.The dispatcher reads sliced boolean input wrong (#6288). Once the batches reaching the filter are sliced, it keeps or drops the wrong rows. In 1.0.0 the filter ran in Spark and returned the right rows, so this is a regression in 1.1.0.
Steps to reproduce
The aggregate emits more groups than
spark.comet.batchSize, and its sliced output batches carry boolean arrays with a non-zero offset into the filter.Expected behavior
Spark and Comet 1.0.0 return 10,000 rows. Comet 1.1.0-rc1 returns 10,001, with rows misclassified in both directions.
Additional context
#6339 fixes #6288 and is on
main, which returns the right rows. #6339 is not onbranch-1.1. I cherry-picked it onto 1.1.0-rc1: it applies cleanly and the reproducer passes, with the predicate still dispatched. The fix for 1.1.0 is to backport #6339.Found by the 1.1.0 regression audit (#6399) and tracked in #6402.