pub struct BufferExec { /* private fields */ }Expand description
WARNING: EXPERIMENTAL
Decouples production and consumption of record batches with an internal queue per partition, eagerly filling up the capacity of the queues even before any message is requested.
┌───────────────────────────┐
│ BufferExec │
│ │
│┌────── Partition 0 ──────┐│
││ ┌────┐ ┌────┐││ ┌────┐
──background poll────────▶│ │ │ ├┼┼───────▶ │
││ └────┘ └────┘││ └────┘
│└─────────────────────────┘│
│┌────── Partition 1 ──────┐│
││ ┌────┐ ┌────┐ ┌────┐││ ┌────┐
──background poll─▶│ │ │ │ │ ├┼┼───────▶ │
││ └────┘ └────┘ └────┘││ └────┘
│└─────────────────────────┘│
│ │
│ ... │
│ │
│┌────── Partition N ──────┐│
││ ┌────┐││ ┌────┐
──background poll───────────────▶│ ├┼┼───────▶ │
││ └────┘││ └────┘
│└─────────────────────────┘│
└───────────────────────────┘The capacity is provided in bytes, and for each buffered record batch it will take into account the size reported by RecordBatch::get_array_memory_size.
If a single record batch exceeds the maximum capacity set in the capacity argument, it’s still
allowed to pass in order to not deadlock the buffer.
This is useful for operators that conditionally start polling one of their children only after other child has finished, allowing to perform some early work and accumulating batches in memory so that they can be served immediately when requested.
Implementations§
Source§impl BufferExec
impl BufferExec
Sourcepub fn new(input: Arc<dyn ExecutionPlan>, capacity: usize) -> Self
pub fn new(input: Arc<dyn ExecutionPlan>, capacity: usize) -> Self
Builds a new BufferExec with the provided capacity in bytes.
Sourcepub fn input(&self) -> &Arc<dyn ExecutionPlan> ⓘ
pub fn input(&self) -> &Arc<dyn ExecutionPlan> ⓘ
Returns the input ExecutionPlan of this BufferExec.
Sourcepub fn capacity(&self) -> usize
pub fn capacity(&self) -> usize
Returns the per-partition capacity in bytes for this BufferExec.
Source§impl BufferExec
impl BufferExec
Sourcepub fn try_from_proto(
node: &PhysicalPlanNode,
ctx: &ExecutionPlanDecodeCtx<'_>,
) -> Result<Arc<dyn ExecutionPlan>>
Available on crate feature proto only.
pub fn try_from_proto( node: &PhysicalPlanNode, ctx: &ExecutionPlanDecodeCtx<'_>, ) -> Result<Arc<dyn ExecutionPlan>>
proto only.Reconstruct a BufferExec from its protobuf representation.
The exact inverse of ExecutionPlan::try_to_proto.
Trait Implementations§
Source§impl Clone for BufferExec
impl Clone for BufferExec
Source§fn clone(&self) -> BufferExec
fn clone(&self) -> BufferExec
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 BufferExec
impl Debug for BufferExec
Source§impl DisplayAs for BufferExec
impl DisplayAs for BufferExec
Source§impl ExecutionPlan for BufferExec
impl ExecutionPlan for BufferExec
Source§fn properties(&self) -> &Arc<PlanProperties> ⓘ
fn properties(&self) -> &Arc<PlanProperties> ⓘ
ExecutionPlan, such as output
ordering(s), partitioning information etc. 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 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 supports_limit_pushdown(&self) -> bool
fn supports_limit_pushdown(&self) -> bool
Source§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 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 try_to_proto(
&self,
ctx: &ExecutionPlanEncodeCtx<'_>,
) -> Result<Option<PhysicalPlanNode>>
fn try_to_proto( &self, ctx: &ExecutionPlanEncodeCtx<'_>, ) -> Result<Option<PhysicalPlanNode>>
proto only.Source§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_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 required_input_ordering(&self) -> Vec<Option<OrderingRequirements>>
fn required_input_ordering(&self) -> Vec<Option<OrderingRequirements>>
ExecutionPlan. 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 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 fetch(&self) -> Option<usize>
fn fetch(&self) -> Option<usize>
None means there is no fetch.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 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 moreAuto Trait Implementations§
impl !RefUnwindSafe for BufferExec
impl !UnwindSafe for BufferExec
impl Freeze for BufferExec
impl Send for BufferExec
impl Sync for BufferExec
impl Unpin for BufferExec
impl UnsafeUnpin for BufferExec
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