apache / apache/datafusion-comet
Fuse OR-connected literal contains predicates on the same input column
- Dominant language
- Scala
- Stars
- 1.3k
- Forks
- 373
- Avg merge
- 2d 4h
- Merged PRs (30d)
- 198
Description
### What is the problem the feature request solves?
Queries sometimes test a string column for any of several literal substrings:
```sql
SELECT id
FROM events
WHERE contains(message, 'timeout')
OR contains(message, 'connection reset')
OR contains(message, 'permission denied')
OR contains(message, 'rate limit');
```
Each individual `Contains` is already native. In the current implementation, [Comet's scalar-needle path](https://github.com/apache/datafusion-comet/blob/451c99963206fa6bf0387239aa12887a16255516/native/spark-expr/src/string_funcs/contains.rs#L75) passes a scalar directly to Arrow's optimized `contains` kernel instead of expanding the needle into a full array. However, the alternatives remain separate searches of the same input. [CometOr](https://github.com/apache/datafusion-comet/blob/451c99963206fa6bf0387239aa12887a16255516/spark/src/main/scala/org/apache/comet/serde/predicates.scala#L102) flattens and balances the OR tree; it does not combine the searches.
For a large string and many absent or late-matching needles, repeated searches can revisit the same bytes and produce intermediate Boolean arrays. A prepared matcher for a fixed set of literal needles could answer “does any needle occur?” with less repeated work. Log/event filtering and text categorization are useful synthetic workload shapes for evaluating this opportunity.
This is a performance hypothesis, not a demonstrated regression or a claim that every pair of `Contains` expressions should be fused. The current single-pattern kernels are optimized, and DataFusion's [OR evaluation](https://github.com/apache/datafusion/blob/55.1.0/datafusion/physical-expr/src/expressions/binary.rs#L540) can avoid work when earlier alternatives match. Preparation cost and early-match behavior must be included in the comparison.
### Describe the potential solution
Add a narrowly scoped optimization for an OR subtree made entirely of eligible literal `Contains` predicates on the same input attribute. Keep the Spark SQL expression unchanged for users; lower the eligible subtree to one native matcher during Comet planning.
#### Initial eligibility
- Every leaf in the selected OR subtree is `Contains(attribute, nonNullStringLiteral)`.
- Every leaf refers to the **same bound input attribute**, not merely the same column name. Restrict the input to an `AttributeReference`/its bound-column representation initially. Do not generalize to repeated casts, UDFs, nested-field expressions, or other computed children.
- All operands use Spark's `UTF8_BINARY` semantics. On Spark versions with collations, retain the existing [CometContains collation check](https://github.com/apache/datafusion-comet/blob/451c99963206fa6bf0387239aa12887a16255516/spark/src/main/scala/org/apache/comet/serde/strings.scala#L365).
- Ordinary expression enablement and native-support checks must already have succeeded for `Or` and every `Contains`. In particular, fusion must not bypass `spark.comet.expression.Contains.enabled=false` or turn a JVM-dispatched collated predicate into a native byte comparison. The existing [expression admission path](https://github.com/apache/datafusion-comet/blob/451c99963206fa6bf0387239aa12887a16255516/spark/src/main/scala/org/apache/comet/serde/QueryPlanSerde.scala#L938) enforces these checks.
- Pattern count and total literal bytes must fall within bounded limits, and the pattern set must meet a benchmark-supported threshold for fusion. Two alternatives are a correctness-eligible shape, not necessarily a profitable one.
Match a complete eligible subtree. Do not collect `Contains` leaves across unrelated OR operands or move them across `AND`, `IF`, `CASE`, or other evaluation boundaries. A surrounding expression must retain its existing evaluation and short-circuit behavior.
#### Implementation approaches
Two reasonable locations are:
1. **Native planning after normal Scala admission:** recognize the serialized/native OR tree whose leaves are the admitted Comet `contains` function with a common bound column and literal needles. Build the prepared matcher in place of that subtree. This avoids bypassing Scala's individual-expression gates and may avoid a new protobuf expression kind.
2. **A dedicated serde lowering:** recognize the restricted shape in `CometOr` and emit a private native expression/function representation. This approach must explicitly preserve the normal admission checks for every consumed leaf, along with explain/coverage attribution and the original fallback path.
Either approach should preserve the original OR implementation when fusion is declined. Avoid repeatedly flattening and compiling overlapping subtrees in a deep OR chain; prepare one matcher for the selected group.
For the native implementation, evaluate a prepared multiple-literal searcher, such as an Aho–Corasick matcher or another suitable byte-oriented algorithm. The algorithm choice should follow benchmarks. Do not translate the literals into user-facing regular-expression semantics.
The matcher should be prepared once per native planned expression and reused across batches. Own the literal data and compiled search state for that expression's lifetime; do not retain input batch buffers or add an unbounded executor-global cache. Bound both accepted pattern data and compiled-state growth, and release the state with its owning plan. An over-limit or recoverably unsupported preparation attempt should use the ordinary OR path rather than introduce a new query failure. This does not require recovering from process-level allocation failure.
#### Required semantics
The result remains Spark's nullable Boolean result:
| Input / pattern set | Result |
| --- | --- |
| Null input with any eligible non-null literals | `NULL` |
| Non-null input with at least one matching literal | `TRUE` |
| Non-null input with no matching literals | `FALSE` |
| Non-null input with an empty-string literal among the alternatives | `TRUE` |
| Null input with an empty-string literal among the alternatives | `NULL` |
Null literals are outside the initial eligibility rule; leave those expressions on the existing path. Empty, duplicate, overlapping, and prefix-related literals must preserve results. An empty literal must not simplify the result to an unconditional `TRUE`, because the input can be null.
Match case-sensitively under `UTF8_BINARY`, including multibyte text and embedded NUL bytes. Do not add Unicode normalization or case folding. Keep row boundaries intact: searching Arrow's concatenated value buffer must not manufacture a match spanning two different rows. Preserve slices, null offsets, and every input representation admitted to the fused path; retain the existing path for representations not supported initially.
#### Validation and acceptance criteria
- [ ] Compare fused, unfused Comet, and Spark results for eligible shapes in both filters and projections. Include nullable input, empty strings/literals, duplicate and overlapping needles, match at the start/end, longer-than-input needles, multibyte text, embedded NULs, and multiple batches.
- [ ] Verify the plan contains one prepared native matcher for an eligible group and that preparation is not repeated per row or batch. Exercise sliced inputs and dictionary-backed inputs if they can reach the admitted path.
- [ ] Verify no fusion for different input attributes, computed/fallible or nondeterministic children, null or nonliteral needles, non-default collations, disabled constituent expressions, and over-limit pattern sets. Include a surrounding conditional to guard against changing evaluation boundaries.
- [ ] Verify repeated execution/plan reuse does not use stale needles and that compiled state is released when its owner is dropped. Test the recoverable preparation-decline path.
- [ ] Benchmark the prepared matcher against the **actual current Comet OR expression**, including its vectorized short-circuit/selection behavior, rather than against an artificial loop that always searches every pattern.
- [ ] Cover 2/4/8/16/32 alternatives, short and long strings, small and large batches, no matches, first-alternative-heavy matches, late-alternative matches, different match positions, common pattern prefixes, and null-heavy input. Measure preparation time, steady-state evaluation, retained matcher memory, and allocation changes; include at least one end-to-end SQL benchmark.
- [ ] Use those results to choose a conservative admission threshold/strategy and provide a way to disable fusion for comparison and troubleshooting. Keep shapes that regress on the existing path. No speedup target is asserted before this measurement.
#### Non-goals
The first implementation does not add a public SQL function, regex alternation, `LIKE ANY`, case-insensitive/collation-aware multi-search, arbitrary expression commoning, dynamic per-row pattern arrays, scan-index integration, or a general predicate-reordering framework. It should fit one focused optimization PR with tests and benchmarks.
### Additional context
- ClickHouse's documented [multiSearchAny](https://clickhouse.com/docs/reference/functions/regular-functions/string-search-functions#multiSearchAny) and [multiSearchAnyUTF8](https://clickhouse.com/docs/reference/functions/regular-functions/string-search-functions#multiSearchAnyUTF8) demonstrate the usefulness of evaluating several literal substring alternatives together. They are an algorithm/API precedent; this proposal retains Spark's Boolean, null, and binary-collation semantics rather than importing another engine's contract.
- DataFusion 55.1.0's [Contains implementation](https://github.com/apache/datafusion/blob/55.1.0/datafusion/functions/src/string/contains.rs#L102) also delegates one needle to Arrow. This proposal targets combining several calls, not replacing a missing single-pattern vector kernel.
- [PR #5322](https://github.com/apache/datafusion-comet/pull/5322) optimizes the scalar-haystack/array-needle case of an individual `Contains`; it does not fuse OR-connected literal needles on an input column. The changes can be evaluated independently.
- The existing [Contains benchmark](https://github.com/apache/datafusion-comet/blob/451c99963206fa6bf0387239aa12887a16255516/native/spark-expr/benches/contains.rs) and [expression optimization guide](https://github.com/apache/datafusion-comet/blob/451c99963206fa6bf0387239aa12887a16255516/docs/source/contributor-guide/optimizing_expressions.md) provide starting points for the benchmark and no-regression work.
Source references above describe Comet main at `451c99963206fa6bf0387239aa12887a16255516`.
Contributor guide
Research direction
Start with docs/source/contributor-guide/optimizing_expressions.md and native/spark-expr/benches/contains.rs, then trace CometOr in spark/src/main/scala/org/apache/comet/serde/predicates.scala and admission in QueryPlanSerde.scala. Compare the existing Contains implementation and tests before choosing the native planning boundary. Done means bounded, eligible OR groups use one reusable matcher, all listed fallback and nullable semantics remain intact, and correctness, plan, lifecycle, regression, and benchmark coverage is provided.
Written by the indexing model from the issue text.
Assessment
- Tech stack
- rust, scala, spark
- Domain
- backend, performance
- Issue type
- Feature
- Difficulty
- 5/5
- Estimated time
- Over a week
- Activity status
- Active
- Clarity
- Clearly specified
- Newbie friendliness
- 35/100