pub struct DataSinkExec { /* private fields */ }Expand description
Execution plan for writing record batches to a DataSink
Returns a single row with the number of values written
Implementations§
Source§impl DataSinkExec
impl DataSinkExec
Sourcepub fn new(
input: Arc<dyn ExecutionPlan>,
sink: Arc<dyn DataSink>,
sort_order: Option<LexRequirement>,
) -> Self
pub fn new( input: Arc<dyn ExecutionPlan>, sink: Arc<dyn DataSink>, sort_order: Option<LexRequirement>, ) -> Self
Create a plan to write to sink
Note: DataSinkExec requires its input to have a single partition.
If the input has multiple partitions, the physical optimizer will
automatically insert a Merge-related operator to merge them.
If you construct PhysicalPlan without going through the physical optimizer,
you must ensure that the input has a single partition.
Sourcepub fn input(&self) -> &Arc<dyn ExecutionPlan> ⓘ
pub fn input(&self) -> &Arc<dyn ExecutionPlan> ⓘ
Input execution plan
Sourcepub fn sort_order(&self) -> &Option<LexRequirement>
pub fn sort_order(&self) -> &Option<LexRequirement>
Optional sort order for output data
Sourcepub fn encode_sort_order(
&self,
ctx: &ExecutionPlanEncodeCtx<'_>,
) -> Result<Option<PhysicalSortExprNodeCollection>>
Available on crate feature proto only.
pub fn encode_sort_order( &self, ctx: &ExecutionPlanEncodeCtx<'_>, ) -> Result<Option<PhysicalSortExprNodeCollection>>
proto only.Encode the optional sink ordering for a protobuf plan node.
Sourcepub fn decode_sort_order(
collection: Option<&PhysicalSortExprNodeCollection>,
ctx: &ExecutionPlanDecodeCtx<'_>,
schema: &Schema,
) -> Result<Option<LexRequirement>>
Available on crate feature proto only.
pub fn decode_sort_order( collection: Option<&PhysicalSortExprNodeCollection>, ctx: &ExecutionPlanDecodeCtx<'_>, schema: &Schema, ) -> Result<Option<LexRequirement>>
proto only.Decode the optional sink ordering from a protobuf plan node.
Trait Implementations§
Source§impl Clone for DataSinkExec
impl Clone for DataSinkExec
Source§fn clone(&self) -> DataSinkExec
fn clone(&self) -> DataSinkExec
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 DataSinkExec
impl Debug for DataSinkExec
Source§impl DisplayAs for DataSinkExec
impl DisplayAs for DataSinkExec
Source§impl ExecutionPlan for DataSinkExec
impl ExecutionPlan for DataSinkExec
Source§fn properties(&self) -> &Arc<PlanProperties> ⓘ
fn properties(&self) -> &Arc<PlanProperties> ⓘ
Return a reference to Any that can be used for downcasting
Source§fn execute(
&self,
partition: usize,
context: Arc<TaskContext>,
) -> Result<SendableRecordBatchStream>
fn execute( &self, partition: usize, context: Arc<TaskContext>, ) -> Result<SendableRecordBatchStream>
Execute the plan and return a stream of RecordBatches for
the specified partition.
Source§fn try_to_proto(
&self,
ctx: &ExecutionPlanEncodeCtx<'_>,
) -> Result<Option<PhysicalPlanNode>>
Available on crate feature proto only.
fn try_to_proto( &self, ctx: &ExecutionPlanEncodeCtx<'_>, ) -> Result<Option<PhysicalPlanNode>>
proto only.Delegates protobuf serialization to the underlying sink.
Source§fn name(&self) -> &'static str
fn name(&self) -> &'static str
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 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 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 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 replace_children(
self: Arc<Self>,
children: Vec<Arc<dyn ExecutionPlan>>,
_: ReplaceChildrenOptions,
) -> Result<Arc<dyn ExecutionPlan>>
fn replace_children( self: Arc<Self>, children: Vec<Arc<dyn ExecutionPlan>>, _: 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 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 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 + 'static)>
fn downcast_delegate(&self) -> Option<&(dyn ExecutionPlan + 'static)>
ExecutionPlan downcast identity. Read moreSource§fn check_invariants(&self, check: InvariantLevel) -> Result<(), DataFusionError>
fn check_invariants(&self, check: InvariantLevel) -> Result<(), DataFusionError>
Source§fn dynamic_expressions_produced(&self) -> Vec<Arc<dyn PhysicalExpr>>
fn dynamic_expressions_produced(&self) -> Vec<Arc<dyn PhysicalExpr>>
Source§fn with_new_children_and_same_properties(
self: Arc<Self>,
children: Vec<Arc<dyn ExecutionPlan>>,
) -> Result<Arc<dyn ExecutionPlan>, DataFusionError>
fn with_new_children_and_same_properties( self: Arc<Self>, children: Vec<Arc<dyn ExecutionPlan>>, ) -> Result<Arc<dyn ExecutionPlan>, DataFusionError>
Use ExecutionPlan::replace_children with ReplaceChildrenOptions
ExecutionPlan::replace_children instead.Source§fn reset_state(
self: Arc<Self>,
) -> Result<Arc<dyn ExecutionPlan>, DataFusionError>
fn reset_state( self: Arc<Self>, ) -> Result<Arc<dyn ExecutionPlan>, DataFusionError>
ExecutionPlan. Read moreSource§fn repartitioned(
&self,
_target_partitions: usize,
_config: &ConfigOptions,
) -> Result<Option<Arc<dyn ExecutionPlan>>, DataFusionError>
fn repartitioned( &self, _target_partitions: usize, _config: &ConfigOptions, ) -> Result<Option<Arc<dyn ExecutionPlan>>, DataFusionError>
ExecutionPlan to
produce target_partitions partitions. Read moreSource§fn partition_statistics(
&self,
partition: Option<usize>,
) -> Result<Arc<Statistics>, DataFusionError>
fn partition_statistics( &self, partition: Option<usize>, ) -> Result<Arc<Statistics>, DataFusionError>
Use StatisticsContext::compute instead
ExecutionPlan node. Read moreSource§fn statistics_from_inputs(
&self,
_input_stats: &[Arc<Statistics>],
args: &StatisticsArgs,
) -> Result<Arc<Statistics>, DataFusionError>
fn statistics_from_inputs( &self, _input_stats: &[Arc<Statistics>], args: &StatisticsArgs, ) -> Result<Arc<Statistics>, DataFusionError>
ExecutionPlan node,
given pre-computed child statistics. 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 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 fetch(&self) -> Option<usize>
fn fetch(&self) -> Option<usize>
None means there is no fetch.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>>, DataFusionError>
fn try_swapping_with_projection( &self, _projection: &ProjectionExec, ) -> Result<Option<Arc<dyn ExecutionPlan>>, DataFusionError>
ExecutionPlan. Read moreSource§fn gather_filters_for_pushdown(
&self,
_phase: FilterPushdownPhase,
parent_filters: Vec<Arc<dyn PhysicalExpr>>,
_config: &ConfigOptions,
) -> Result<FilterDescription, DataFusionError>
fn gather_filters_for_pushdown( &self, _phase: FilterPushdownPhase, parent_filters: Vec<Arc<dyn PhysicalExpr>>, _config: &ConfigOptions, ) -> Result<FilterDescription, DataFusionError>
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>>, DataFusionError>
fn handle_child_pushdown_result( &self, _phase: FilterPushdownPhase, child_pushdown_result: ChildPushdownResult, _config: &ConfigOptions, ) -> Result<FilterPushdownPropagation<Arc<dyn ExecutionPlan>>, DataFusionError>
Source§fn with_new_state(
&self,
_state: Arc<dyn Any + Sync + Send>,
) -> Option<Arc<dyn ExecutionPlan>>
fn with_new_state( &self, _state: Arc<dyn Any + Sync + Send>, ) -> Option<Arc<dyn ExecutionPlan>>
Source§fn try_pushdown_sort(
&self,
_order: &[PhysicalSortExpr],
) -> Result<SortOrderPushdownResult<Arc<dyn ExecutionPlan>>, DataFusionError>
fn try_pushdown_sort( &self, _order: &[PhysicalSortExpr], ) -> Result<SortOrderPushdownResult<Arc<dyn ExecutionPlan>>, DataFusionError>
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 DataSinkExec
impl !UnwindSafe for DataSinkExec
impl Freeze for DataSinkExec
impl Send for DataSinkExec
impl Sync for DataSinkExec
impl Unpin for DataSinkExec
impl UnsafeUnpin for DataSinkExec
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