Skip to content

perf: evaluate JVM-dispatched filter conjuncts only on rows passing preceding conjuncts - #11

Merged
msafonov merged 2 commits into
branch-1.1from
joom/1.1-fix-filter-short-circuit
Oct 10, 2026
Merged

msafonov merged 2 commits into
branch-1.1from
joom/1.1-fix-filter-short-circuit

Conversation

@msafonov

Copy link
Copy Markdown
Collaborator

Problem

DataFusion 55's AND evaluates its right side on every row once more than 20% of rows pass the left side (PRE_SELECTION_THRESHOLD). When the right side is a conjunct that runs through the JVM codegen dispatcher (JvmScalarUdf: ScalaUDF, regex family outside the native subset, other dispatched built-ins), every row pays for the JNI round trip, the Arrow export and the Java evaluation. Spark never evaluates it for rows where the preceding conjuncts are false.

Example: dbt uzum_available_products, filter public AND hasActive AND enabledByMerchant AND NOT regexp_like(regexp_replace(lower(name), <subquery>, ''), <subquery>) AND .... The left side passes 31% of rows, so the regex ran on all rows, and Comet was 3.6x slower than vanilla.

Change

CometFilterExec serde, applied only when the condition is deterministic:

  • It flattens the already serialized predicate into its top-level conjuncts (the And proto tree).
  • A conjunct is expensive when its proto contains a JvmScalarUdf node.
  • Each expensive conjunct that is not the first one starts a new segment. Segments are chained as CASE WHEN <previous segments> THEN <segment> END, and conjuncts inside a segment remain a balanced AND.
  • DataFusion's CaseExpr (ExpressionOrExpression) projects only the columns the THEN branch uses, filters them by the WHEN mask and evaluates the dispatcher only on the selected rows. With no ELSE, it does not copy the rejected rows.
  • Semantics: a null result rejects the row just like false, so the filter output is identical. Conjunct order is preserved. Rows where a preceding conjunct is false or null no longer reach the expensive conjunct, which also matches Spark's short-circuit for errors.
  • Only the native filter predicate changes. CometFilterExec.condition (explain, plan stability) and scan dataFilters / pushdown are untouched.
  • Kill switch: spark.comet.exec.filter.shortCircuitJvmDispatch.enabled (default true).

Tests

  • New CometFilterShortCircuitSuite covers the plan rewrite (CASE depth), results equal to Spark with three-valued logic, nested AND/OR, multiple expensive conjuncts, no rewrite for OR, for a single conjunct or without dispatch, no rewrite for nondeterministic conditions, and unchanged scan pushdown. It uses a UDF that throws when it is evaluated on rows the guard rejects: it passes with the rewrite and throws with the switch off.
  • CometExpressionSuite, CometExecSuite, CometRegExpJvmSuite, CometRegexParitySuite, CometCodegenSuite, CometScalaUDFClassLoaderSuite, CometUdfBridgeSuite, CometStringExpressionSuite, CometSqlFileTestSuite: 1113 passed, 0 failed.
  • CometTPCDSV1_4/V2_7 PlanStabilitySuite: 129 passed. spotless clean.

Real data (uzum_available_products filter, 4% sample of product_products_daily_snapshot, 7.82M rows, msaf-spark-comet:joom-1.1-6 with the patched driver classes, 4 executors)

run task-h result (count, sum xxhash64(_id))
vanilla 1.392 2091770, 3934135682968584027
Comet, switch off (current behavior) 5.097 same
Comet, this PR 1.854 same
Comet, manual CASE WHEN a THEN b ELSE false END guard 2.181 same

🤖 Generated with Claude Code

…receding conjuncts

DataFusion's AND evaluates its right side on every row once more than 20% of
rows pass the left side. A conjunct that runs through the JVM codegen
dispatcher (ScalaUDF, regex family outside the native subset, other
dispatched expressions) then costs a JNI round trip and Arrow export for rows
Spark would never evaluate it on.

CometFilterExec now splits the serialized deterministic predicate into its
top-level conjuncts and rewrites `a AND b`, where b contains a JvmScalarUdf,
into `CASE WHEN a THEN b END`. DataFusion's CaseExpr projects the columns b
needs, filters them by a and evaluates b on the selected rows only. Null
results reject rows just like false, so the filter output is unchanged. Scan
pushdown is untouched because only the native filter predicate changes.

Controlled by spark.comet.exec.filter.shortCircuitJvmDispatch.enabled
(default true).

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
@github-actions github-actions Bot added enhancement New feature or request performance labels Oct 10, 2026
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
msafonov added a commit that referenced this pull request Oct 10, 2026
@msafonov
msafonov merged commit c75273e into branch-1.1 Oct 10, 2026
35 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

enhancement New feature or request performance

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant