Repository navigation
Conversation
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ Coverage Diff @@
## main #25693 +/- ##
==========================================
+ Coverage 82.49% 82.64% +0.15%
==========================================
Files 1140 1147 +7
Lines 437667 445889 +8222
Branches 437667 445889 +8222
==========================================
+ Hits 361039 368513 +7474
- Misses 54833 54987 +154
- Partials 21795 22389 +594 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
|
@sunchao Can you take a look? Thanks! |
sunchao
left a comment
There was a problem hiding this comment.
Thanks, @viirya. The logical-null changes fix both reported wrong-result paths. I ran all 100 matching null-aware join tests on d3ccef0; they passed. Applying only the new tests to parent 2306a4b produced the expected 10 failures, with the 20 pre-existing mark cases passing. Current-head CI has 40 successful and 3 skipped checks.
I found one introduced performance issue in final mark emission, detailed inline. Please bound the logical-null computation to the current output chunk or compute it once for the join before approval.
| probe_side_has_null: bool, | ||
| probe_side_non_empty: bool, | ||
| ) -> ArrayRef { | ||
| let build_key_nulls = build_key_column.logical_nulls(); |
There was a problem hiding this comment.
[P2] Avoid rescanning the full build dictionary for every output batch
emit_unmatched_build_rows calls this helper once per batch_size chunk, but build_key_column is the entire concatenated build column. In Arrow 60, DictionaryArray::logical_nulls() allocates an N-bit bitmap and walks all N keys when the dictionary values contain NULL. This makes final emission do N * ceil(N / batch_size) key inspections. It also runs for correlated joins even though null_indices_bitmap supplies the answer and the new buffer is unused.
I reproduced this with a null-aware CollectLeft LeftMark: dictionary<Int32, Int32> build values [1, NULL], alternating keys [0, 1, ...], an empty probe, and batch size 8192. Both revisions return exactly N non-null FALSE marks. With separate base/head target directories and the locked Arrow 60 dependencies, median times over three unoptimized operator runs were 0.185s vs 1.165s for 524,288 rows, and 0.369s vs 4.274s for 1,048,576 rows. Reversing run order retained the smaller-case slowdown; the equivalent plain Int32 control stayed comparable. These are diagnostic debug-build timings, not release query benchmarks.
Could we compute logical nulls only for the current contiguous LeftMark chunk (adjusting the index), or retain the bitmap once for final emission, and skip this work when the correlated bitmap is present? That preserves the fix without repeatedly scanning the full build side.
There was a problem hiding this comment.
Thanks, this is valid. I reproduced the repeated full-dictionary scan with 524,288 rows and batch size 8192. The diagnostic runtime dropped from about 693ms to 141ms after limiting logical_nulls() to the current contiguous output chunk. Correlated joins now skip this computation and continue using their precomputed bitmap. All 104 null-aware tests and the complete datafusion-physical-plan all-features suite pass.
There was a problem hiding this comment.
Thanks, @viirya. I confirmed that 4d50daa fixes the original flat-dictionary case: the same 524,288-row diagnostic improved from a median 1.152s on d3ccef0 to 0.209s on the new head.
One part of this P2 remains for nested dictionary build keys at the new slice/logical_nulls call. Slicing a DictionaryArray only slices its keys and retains all its values. When those values are themselves a dictionary, Arrow 60 recursively materializes the entire inner dictionary's logical-null bitmap on every output chunk. The outer slice therefore does not bound that inner scan.
I verified a nonempty nested-dictionary join executes and produces the expected TRUE/NULL/FALSE marks. To isolate final emission, I then used an empty probe and Dictionary<Int32, Dictionary<Int32, Int32>> build keys: outer keys 0..N, N inner keys all 0, leaf values [1, NULL, 4], batch size 8192. Both pre-PR base 2306a4b and current head produce exactly N non-null FALSE marks. At 524,288 rows, medians were 0.234s on base versus 0.954s on head; at 1,048,576 rows, 0.471s versus 3.307s. These are three-run unoptimized native-operator diagnostics with identical compiler/locked dependencies and separate target directories, not release SQL benchmarks.
Could we compute logical nulls once on entry to final emission and slice that bitmap for each chunk, accounting for any retained allocation? That avoids repeatedly traversing nested dictionary values as well as the flat keys.
There was a problem hiding this comment.
Thanks, confirmed. Slicing the outer dictionary retains its full child dictionaries, so chunk-local logical_nulls still rescanned nested values. Fixed in 3bc2308 by computing build-key logical nulls once during build collection, retaining only a non-empty bitmap, accounting for newly materialized buffer memory, and indexing the cached bitmap by the original build row during final emission. The 524,288-row nested-dictionary diagnostic dropped from about 561ms to 177ms locally. I also added nested-dictionary correctness coverage across five batch sizes; all 109 null-aware tests and the complete 2,409-test datafusion-physical-plan all-features suite pass.
sunchao
left a comment
There was a problem hiding this comment.
Thanks, @viirya. Re-reviewed 4d50daa with independent source passes and fresh runtime checks. The flat-dictionary build scan from my earlier review is fixed, and the chunk-local null indexing and correlated bypass look correct. All 100 matching null-aware tests pass; current-head CI has 40 successful and 3 skipped checks.
Two performance issues remain: the nested-dictionary part of final emission (updated in the existing thread), and repeated whole-probe logical-null scans on lookup resumption (inline below). Both were reproduced against the pre-PR base with identical result assertions. Please address these before approval. Timings in the comments are scoped to unoptimized native-operator diagnostics.
| } else { | ||
| probe_key_column.null_count() > 0 | ||
| }; | ||
| let probe_has_null = probe_key_column.logical_null_count() > 0; |
There was a problem hiding this comment.
[P2] Record probe logical nulls once per input batch
process_probe_batch calls this on every resumed lookup chunk, not just when a probe batch is first fetched. With dictionary values containing a NULL, Arrow 60's logical_null_count() walks every key in the whole probe batch, even if that NULL dictionary entry is unused. This change therefore adds another full probe-key scan for every lookup chunk.
I reproduced the slowdown with one build value 1 and one native probe batch of 524,288 dictionary keys all 0, dictionary values [1, NULL], and join batch size 8192. Both pre-PR base 2306a4b and head 4d50daa return exactly one TRUE mark. Median unoptimized operator time over three runs increased from 1.610s to 2.348s; the plain Int32 control stayed about 0.057s. Compiler, locked Arrow 60 dependencies, and fixture were identical, with separate target directories. This oversized-batch case is supported through native/custom execution plans; normal DataSourceExec scans split batches, so these are not ordinary SQL performance claims.
Please compute and record the probe summary only on the first chunk (state.offset == (0, None)), as the adjacent correlated bookkeeping already does. Keep the anti-join's shared NULL-hint check outside that guard so another partition can still trigger early exit between chunks.
There was a problem hiding this comment.
Thanks, confirmed. Fixed in 3bc2308 by computing and recording the probe logical-null summary only when state.offset is the initial (0, None). The shared NULL-hint check for LeftAnti remains outside that guard, so another partition can still stop resumed lookup early. The 524,288-row oversized-probe diagnostic dropped from about 1.29s to 910ms with identical result assertions. All 109 null-aware tests and the complete datafusion-physical-plan all-features suite pass.
Which issue does this PR close?
Rationale for this change
Null-aware
LeftMarkjoins use physical null checks for dictionary keys. They miss valid dictionary keys that point toNULLdictionary values and returnFALSEmarks where SQL three-valued logic requiresNULL.What changes are included in this PR?
Null-aware
LeftMarkjoins now use Arrow's logical-null APIs on both probe and build keys. The build-side logical-null bitmap is computed once during build collection, retained with memory accounting, and reused by final output chunks. Probe-side null summaries are computed once per input batch even when hash lookup resumes across several chunks. Correlated joins continue using their precomputed per-row null bitmap.What is the testing strategy for this PR?
The diff adds three
HashJoinExectest functions covering logical dictionary nulls on the build and probe sides, including nested dictionary build keys.#[apply(hash_join_exec_configs)]expands each function across five batch sizes, so Cargo reports 15 new test instances. The probe-side function also checks both probe-partition completion orders inside each instance.Ablation testing confirmed that the original build- and probe-side functions fail at every batch size without the correctness fix. With the fix, all 109 null-aware tests and the complete 2,409-test
datafusion-physical-planall-features suite pass.For 524,288-row debug diagnostics, bounding the original flat-dictionary scan reduced runtime from about 693ms to 141ms. The follow-up nested-dictionary diagnostic dropped from about 561ms to 177ms, and the oversized resumed-probe diagnostic dropped from about 1.29s to 910ms.
Are there any user-facing changes?
Null-aware mark joins over dictionary-encoded keys now return the correct nullable marks. There are no public API changes.