Adaptive filters on morsel - #57
Closed
adriangb wants to merge 80 commits into
Closed
Conversation
This PR implements morsel-driven execution for Parquet files in DataFusion, enabling row-group level work sharing across partitions to mitigate data skew. Key changes: - Introduced `WorkQueue` in `datafusion/datasource/src/file_stream.rs` for shared pool of work. - Added `morselize` method to `FileOpener` trait to allow dynamic splitting of files into morsels. - Implemented `morselize` for `ParquetOpener` to split files into individual row groups. - Cached `ParquetMetaData` in `ParquetMorsel` extensions to avoid redundant I/O. - Modified `FileStream` to support work stealing from the shared queue. - Implemented `Weak` pointer pattern for `WorkQueue` in `FileScanConfig` to support plan re-executability. - Added `MorselizingGuard` to ensure shared state consistency on cancellation. - Added `allow_morsel_driven` configuration option (enabled by default for Parquet). - Implemented row-group pruning during the morselization phase for better efficiency. Tests: - Added `parquet_morsel_driven_execution` test to verify work distribution and re-executability. - Added `parquet_morsel_driven_enabled_by_default` to verify the default configuration. Co-authored-by: Dandandan <163737+Dandandan@users.noreply.github.com>
…en-execution-237164415184908839
These tests don't need morsel-driven execution disabled: - custom_datasource: uses a custom ExecutionPlan, not file-based - partition_statistics: only checks statistics metadata - json_shredding: single-row filtered result is order-independent Also remove leftover query_id getter and unused atomic imports. Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
…tion When morsel-driven execution is enabled, the WorkQueue handles load balancing at runtime, making byte-range file splitting unnecessary. Distribute whole files round-robin across target partitions instead. Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
After morselizing a file into row-group morsels, push them to the front of the shared queue instead of the back. This way the same (or nearby) worker picks up sibling row groups next, keeping I/O sequential within each file and the page cache warm. Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
…en-execution-237164415184908839 # Conflicts: # datafusion/sqllogictest/test_files/dynamic_filter_pushdown_config.slt
After morselizing, push all row-group morsels to the front of the queue and return to Idle, instead of keeping the first morsel and opening it inline. The worker then pulls the first morsel through the normal is_leaf_morsel fast path, keeping the code simpler while preserving I/O locality. Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
Split the single WorkQueue into two internal queues: one for whole files awaiting morselization and one for already-morselized leaf morsels (row groups). Workers drain the morsel queue first, so freshly produced row groups are consumed before the next file is opened. This keeps I/O sequential within each file without needing push-to-front tricks. Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
Stop the time_opening timer before transitioning back to Idle after pushing morsels to the queue. Without this, re-entering Idle would call start() on an already-running timer, triggering an assertion failure. Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
Combines morsel-driven Parquet scan (per-row-group work units) with adaptive filter pushdown (selectivity-based filter placement). Key merge decisions: - morselize() uses combined predicate for coarse pruning (unchanged) - open() uses predicate_conjuncts with selectivity tracker per morsel - Proto field renumbered: filter_pushdown_min_bytes_per_sec 35 -> 42 to avoid collision with allow_morsel_driven (field 35) - SQL test outputs take morsel-driven plan format Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
- Fix proto field collision: filter_pushdown_min_bytes_per_sec 35 -> 42 - Fix open() to use pruning_predicate for row group pruning - Fix morselize() to derive predicate from predicate_conjuncts - Update SLT expected outputs for Optional(DynamicFilter) format - Use slt:ignore for scan_efficiency_ratio (float precision) Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
Member
Author
|
run benchmarks |
Member
Author
|
run benchmark clickbench_partitioned |
Replace per-partition WorkQueue-based morsel execution with a single
shared pipeline that uses buffer_unordered at both the morselize and
open stages, connected to all partitions via a bounded MPMC channel
(async-channel). This decouples I/O concurrency from CPU parallelism.
Pipeline architecture:
files → buffer_unordered(M) morselize → flatten morsels
→ buffer_unordered(N) open → drain batches → MPMC channel
New config options:
- morsel_morselize_concurrency (default 0 = 2×CPUs): concurrent metadata fetches
- morsel_open_concurrency (default 2): concurrent row group opens
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
adriangb
force-pushed
the
adaptive-filters-on-morsel
branch
from
March 2, 2026 21:35
d2095a6 to
8c8a383
Compare
…dered" This reverts commit 40a8307.
adriangb
force-pushed
the
adaptive-filters-on-morsel
branch
from
March 2, 2026 22:22
8c8a383 to
217e391
Compare
|
Thank you for your contribution. Unfortunately, this pull request is stale because it has been open 60 days with no activity. Please remove the stale label or comment or this will be closed in 7 days. |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
No description provided.