pub struct PartialSortExec { /* private fields */ }Expand description
Sort execution plan for inputs that are already partially sorted.
This operator takes input ordered by a prefix of the required ordering, and
produces output ordered by the required ordering, emitting rows sooner
(streaming) and using less peak memory than SortExec which must buffer
all rows before producing any output.
PartialSortExec relies on the property that rows with the same sort
prefix are contiguous, so it can sort one prefix group at a time, emitting
completed groups without reading (and buffering) the entire input.
For example, if the required output is (a, b, c), but the input is only
ordered by (a, b), PartialSortExec sorts only within each (a, b)
group to produce output ordered by (a, b, c).
input ordered by a, b output ordered by a, b, c
+---+---+---+ +---+---+---+
| a | b | c | | a | b | c |
+---+---+---+ +---+---+---+
| 0 | 0 | 3 | -- new group --> | 0 | 0 | 1 |
| 0 | 0 | 2 | | 0 | 0 | 2 |
| 0 | 0 | 1 | | 0 | 0 | 3 |
| 0 | 1 | 1 | -- new group --> | 0 | 1 | 1 |
| 0 | 2 | 4 | -- new group --> | 0 | 2 | 0 |
| 0 | 2 | 0 | | 0 | 2 | 4 |
| 1 | 0 | 5 | -- new group --> | 1 | 0 | 5 |
+---+---+---+ +---+---+---+§Buffering and Emitting Rows
PartialSortExec buffers rows only until it can prove a prefix group
will never be seen again, then sorts and emits buffered rows. A group is
guaranteed to never be seen again once a row with a different prefix
value arrives. This relies on the input’s existing ordering guarantees.
Using the example from above, rows accumulate in the in-memory buffer in
batches. As long as the (a, b) prefix keeps repeating, more rows are
buffered.
Buffer
+---+---+---+
| a | b | c |
+---+---+---+
| 0 | 0 | 3 |
| 0 | 0 | 2 |
| 0 | 0 | 1 |
+---+---+---+Once a batch arrives that contains a new (a, b) prefix, e.g. (0, 2):
every buffered row for previous prefixes may be emitted:
Buffer
+---+---+---+
| a | b | c |
+---+---+---+
| 0 | 0 | 3 |
| 0 | 0 | 2 |
| 0 | 0 | 1 |
| 0 | 1 | 1 | <-- first row of new batch, new prefix
| 0 | 2 | 4 | <-- new prefix
| 0 | 2 | 0 |
| 1 | 0 | 5 | <-- last row of new batch, new prefix
+---+---+---+Once known complete, the buffered rows are sorted by the full (a, b, c)
ordering and emitted as a RecordBatch; Any rows from the most recently
seen prefix remain buffered (as more rows with the same prefix may arrive in
future batches.
Emitted <-- fully sorted on (a, b, c)
+---+---+---+
| a | b | c |
+---+---+---+
| 0 | 0 | 1 | <-- completed group
| 0 | 0 | 2 |
| 0 | 0 | 3 |
| 0 | 2 | 0 | <-- completed group
| 0 | 2 | 4 |
| 0 | 1 | 1 | <-- completed group
+---+---+---+
Buffer
+---+---+---+
| a | b | c |
+---+---+---+
| 1 | 0 | 5 | <-- (possibly) in progress group
+---+---+---+Implementations§
Source§impl PartialSortExec
impl PartialSortExec
Sourcepub fn new(
expr: LexOrdering,
input: Arc<dyn ExecutionPlan>,
common_prefix_length: usize,
) -> Self
pub fn new( expr: LexOrdering, input: Arc<dyn ExecutionPlan>, common_prefix_length: usize, ) -> Self
Create a new partial sort execution plan
Sourcepub fn preserve_partitioning(&self) -> bool
pub fn preserve_partitioning(&self) -> bool
Whether this PartialSortExec preserves partitioning of the children
Sourcepub fn with_preserve_partitioning(self, preserve_partitioning: bool) -> Self
pub fn with_preserve_partitioning(self, preserve_partitioning: bool) -> Self
Specify the partitioning behavior of this partial sort exec
If preserve_partitioning is true, sorts each partition
individually, producing one sorted stream for each input partition.
If preserve_partitioning is false, sorts and merges all
input partitions producing a single, sorted partition.
Sourcepub fn with_fetch(self, fetch: Option<usize>) -> Self
pub fn with_fetch(self, fetch: Option<usize>) -> Self
Modify how many rows to include in the result
If None, then all rows will be returned, in sorted order.
If Some, then only the top fetch rows will be returned.
This can reduce the memory pressure required by the sort
operation since rows that are not going to be included
can be dropped.
Sourcepub fn input(&self) -> &Arc<dyn ExecutionPlan> ⓘ
pub fn input(&self) -> &Arc<dyn ExecutionPlan> ⓘ
Input schema
Sourcepub fn expr(&self) -> &LexOrdering
pub fn expr(&self) -> &LexOrdering
Sort expressions
Sourcepub fn fetch(&self) -> Option<usize>
pub fn fetch(&self) -> Option<usize>
If Some(fetch), limits output to only the first “fetch” items
Sourcepub fn common_prefix_length(&self) -> usize
pub fn common_prefix_length(&self) -> usize
Common prefix length
Trait Implementations§
Source§impl Clone for PartialSortExec
impl Clone for PartialSortExec
Source§fn clone(&self) -> PartialSortExec
fn clone(&self) -> PartialSortExec
1.0.0 (const: unstable) · Source§fn clone_from(&mut self, source: &Self)
fn clone_from(&mut self, source: &Self)
source. Read moreSource§impl Debug for PartialSortExec
impl Debug for PartialSortExec
Source§impl DisplayAs for PartialSortExec
impl DisplayAs for PartialSortExec
Source§impl ExecutionPlan for PartialSortExec
impl ExecutionPlan for PartialSortExec
Source§fn name(&self) -> &'static str
fn name(&self) -> &'static str
Source§fn properties(&self) -> &Arc<PlanProperties> ⓘ
fn properties(&self) -> &Arc<PlanProperties> ⓘ
ExecutionPlan, such as output
ordering(s), partitioning information etc. Read moreSource§fn fetch(&self) -> Option<usize>
fn fetch(&self) -> Option<usize>
None means there is no fetch.Source§fn required_input_distribution(&self) -> Vec<Distribution>
fn required_input_distribution(&self) -> Vec<Distribution>
Use input_distribution_requirements
Source§fn input_distribution_requirements(&self) -> InputDistributionRequirements
fn input_distribution_requirements(&self) -> InputDistributionRequirements
Source§fn benefits_from_input_partitioning(&self) -> Vec<bool>
fn benefits_from_input_partitioning(&self) -> Vec<bool>
ExecutionPlan benefits from increased
parallelization at its input for each child. Read moreSource§fn children(&self) -> Vec<&Arc<dyn ExecutionPlan>>
fn children(&self) -> Vec<&Arc<dyn ExecutionPlan>>
ExecutionPlans that act as inputs to this plan.
The returned list will be empty for leaf nodes such as scans, will contain
a single value for unary nodes, or two values for binary nodes (such as
joins).Source§fn apply_expressions(
&self,
f: &mut dyn FnMut(&Arc<dyn PhysicalExpr>) -> Result<TreeNodeRecursion>,
) -> Result<TreeNodeRecursion>
fn apply_expressions( &self, f: &mut dyn FnMut(&Arc<dyn PhysicalExpr>) -> Result<TreeNodeRecursion>, ) -> Result<TreeNodeRecursion>
f to each root expression that this node owns and uses
during execution, either by evaluating it or updating it dynamically. Read moreSource§fn replace_children(
self: Arc<Self>,
children: Vec<Arc<dyn ExecutionPlan>>,
options: ReplaceChildrenOptions,
) -> Result<Arc<dyn ExecutionPlan>>
fn replace_children( self: Arc<Self>, children: Vec<Arc<dyn ExecutionPlan>>, options: ReplaceChildrenOptions, ) -> Result<Arc<dyn ExecutionPlan>>
Source§fn with_new_children(
self: Arc<Self>,
children: Vec<Arc<dyn ExecutionPlan>>,
) -> Result<Arc<dyn ExecutionPlan>>
fn with_new_children( self: Arc<Self>, children: Vec<Arc<dyn ExecutionPlan>>, ) -> Result<Arc<dyn ExecutionPlan>>
Use ExecutionPlan::replace_children with ReplaceChildrenOptions
Source§fn with_new_children_and_same_properties(
self: Arc<Self>,
children: Vec<Arc<dyn ExecutionPlan>>,
) -> Result<Arc<dyn ExecutionPlan>>
fn with_new_children_and_same_properties( self: Arc<Self>, children: Vec<Arc<dyn ExecutionPlan>>, ) -> Result<Arc<dyn ExecutionPlan>>
Use ExecutionPlan::replace_children with ReplaceChildrenOptions
ExecutionPlan::replace_children instead.Source§fn execute(
&self,
partition: usize,
context: Arc<TaskContext>,
) -> Result<SendableRecordBatchStream>
fn execute( &self, partition: usize, context: Arc<TaskContext>, ) -> Result<SendableRecordBatchStream>
Source§fn metrics(&self) -> Option<MetricsSet>
fn metrics(&self) -> Option<MetricsSet>
Metrics for this
ExecutionPlan. If no Metrics are available, return None. Read moreSource§fn child_stats_requests(&self, partition: Option<usize>) -> Vec<ChildStats>
fn child_stats_requests(&self, partition: Option<usize>) -> Vec<ChildStats>
StatisticsContext should resolve
before calling Self::statistics_from_inputs. Read moreSource§fn statistics_from_inputs(
&self,
input_stats: &[Arc<Statistics>],
_args: &StatisticsArgs,
) -> Result<Arc<Statistics>>
fn statistics_from_inputs( &self, input_stats: &[Arc<Statistics>], _args: &StatisticsArgs, ) -> Result<Arc<Statistics>>
ExecutionPlan node,
given pre-computed child statistics. Read moreSource§fn static_name() -> &'static strwhere
Self: Sized,
fn static_name() -> &'static strwhere
Self: Sized,
name but can be called without an instance.Source§fn downcast_delegate(&self) -> Option<&dyn ExecutionPlan>
fn downcast_delegate(&self) -> Option<&dyn ExecutionPlan>
ExecutionPlan downcast identity. Read moreSource§fn check_invariants(&self, check: InvariantLevel) -> Result<()>
fn check_invariants(&self, check: InvariantLevel) -> Result<()>
Source§fn dynamic_expressions_produced(&self) -> Vec<Arc<dyn PhysicalExpr>>
fn dynamic_expressions_produced(&self) -> Vec<Arc<dyn PhysicalExpr>>
Source§fn required_input_ordering(&self) -> Vec<Option<OrderingRequirements>>
fn required_input_ordering(&self) -> Vec<Option<OrderingRequirements>>
ExecutionPlan. Read moreSource§fn maintains_input_order(&self) -> Vec<bool>
fn maintains_input_order(&self) -> Vec<bool>
false if this ExecutionPlan’s implementation may reorder
rows within or between partitions. Read moreSource§fn reset_state(self: Arc<Self>) -> Result<Arc<dyn ExecutionPlan>>
fn reset_state(self: Arc<Self>) -> Result<Arc<dyn ExecutionPlan>>
ExecutionPlan. Read moreSource§fn repartitioned(
&self,
_target_partitions: usize,
_config: &ConfigOptions,
) -> Result<Option<Arc<dyn ExecutionPlan>>>
fn repartitioned( &self, _target_partitions: usize, _config: &ConfigOptions, ) -> Result<Option<Arc<dyn ExecutionPlan>>>
ExecutionPlan to
produce target_partitions partitions. Read moreSource§fn partition_statistics(
&self,
partition: Option<usize>,
) -> Result<Arc<Statistics>>
fn partition_statistics( &self, partition: Option<usize>, ) -> Result<Arc<Statistics>>
Use StatisticsContext::compute instead
ExecutionPlan node. Read moreSource§fn supports_limit_pushdown(&self) -> bool
fn supports_limit_pushdown(&self) -> bool
Source§fn with_fetch(&self, _limit: Option<usize>) -> Option<Arc<dyn ExecutionPlan>>
fn with_fetch(&self, _limit: Option<usize>) -> Option<Arc<dyn ExecutionPlan>>
ExecutionPlan node, if it supports
fetch limits. Returns None otherwise. Read moreSource§fn cardinality_effect(&self) -> CardinalityEffect
fn cardinality_effect(&self) -> CardinalityEffect
Source§fn try_swapping_with_projection(
&self,
_projection: &ProjectionExec,
) -> Result<Option<Arc<dyn ExecutionPlan>>>
fn try_swapping_with_projection( &self, _projection: &ProjectionExec, ) -> Result<Option<Arc<dyn ExecutionPlan>>>
ExecutionPlan. Read moreSource§fn gather_filters_for_pushdown(
&self,
_phase: FilterPushdownPhase,
parent_filters: Vec<Arc<dyn PhysicalExpr>>,
_config: &ConfigOptions,
) -> Result<FilterDescription>
fn gather_filters_for_pushdown( &self, _phase: FilterPushdownPhase, parent_filters: Vec<Arc<dyn PhysicalExpr>>, _config: &ConfigOptions, ) -> Result<FilterDescription>
ExecutionPlan::gather_filters_for_pushdown: Read moreSource§fn handle_child_pushdown_result(
&self,
_phase: FilterPushdownPhase,
child_pushdown_result: ChildPushdownResult,
_config: &ConfigOptions,
) -> Result<FilterPushdownPropagation<Arc<dyn ExecutionPlan>>>
fn handle_child_pushdown_result( &self, _phase: FilterPushdownPhase, child_pushdown_result: ChildPushdownResult, _config: &ConfigOptions, ) -> Result<FilterPushdownPropagation<Arc<dyn ExecutionPlan>>>
Source§fn with_new_state(
&self,
_state: Arc<dyn Any + Send + Sync>,
) -> Option<Arc<dyn ExecutionPlan>>
fn with_new_state( &self, _state: Arc<dyn Any + Send + Sync>, ) -> Option<Arc<dyn ExecutionPlan>>
Source§fn try_pushdown_sort(
&self,
_order: &[PhysicalSortExpr],
) -> Result<SortOrderPushdownResult<Arc<dyn ExecutionPlan>>>
fn try_pushdown_sort( &self, _order: &[PhysicalSortExpr], ) -> Result<SortOrderPushdownResult<Arc<dyn ExecutionPlan>>>
Source§fn with_preserve_order(
&self,
_preserve_order: bool,
) -> Option<Arc<dyn ExecutionPlan>>
fn with_preserve_order( &self, _preserve_order: bool, ) -> Option<Arc<dyn ExecutionPlan>>
ExecutionPlan that is aware of order-sensitivity. Read moreSource§fn try_to_proto(
&self,
_ctx: &ExecutionPlanEncodeCtx<'_>,
) -> Result<Option<PhysicalPlanNode>>
fn try_to_proto( &self, _ctx: &ExecutionPlanEncodeCtx<'_>, ) -> Result<Option<PhysicalPlanNode>>
proto only.Auto Trait Implementations§
impl !RefUnwindSafe for PartialSortExec
impl !UnwindSafe for PartialSortExec
impl Freeze for PartialSortExec
impl Send for PartialSortExec
impl Sync for PartialSortExec
impl Unpin for PartialSortExec
impl UnsafeUnpin for PartialSortExec
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
Source§impl<T> CloneToUninit for Twhere
T: Clone,
impl<T> CloneToUninit for Twhere
T: Clone,
Source§impl<T> IntoEither for T
impl<T> IntoEither for T
Source§fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
self into a Left variant of Either<Self, Self>
if into_left is true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read moreSource§fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
self into a Left variant of Either<Self, Self>
if into_left(&self) returns true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read more