Skip to main content

CoalesceRule

Struct CoalesceRule 

Source
pub struct CoalesceRule { /* private fields */ }
Expand description

Merges partitions whose memory_bytes falls below min_partition_bytes.

When coalescing is beneficial (i.e. the advised group count is smaller than the current partition count), apply rewrites the physical plan by appending a NodeOp::CoalescePartitions node that signals downstream operators to merge the output into target_partitions partitions.

Implementations§

Source§

impl CoalesceRule

Source

pub fn new(min_partition_bytes: u64) -> Self

Create a new CoalesceRule with the given minimum partition byte threshold.

Uses the default target_partition_bytes of 128 MiB and no parallelism floor; see Self::with_min_partitions.

Source

pub fn with_target_partition_bytes(self, target_partition_bytes: u64) -> Self

Set a custom target_partition_bytes (bytes per merged output partition).

Source

pub fn target_partition_bytes(&self) -> u64

Return the configured target_partition_bytes.

Source

pub fn with_min_partitions(self, min_partitions: usize) -> Self

Never coalesce below min_partitions partitions.

Sizing partitions purely by bytes answers “how big should a partition be” and never asks “how many workers are there”. A stage whose whole output is under target_partition_bytes collapses to a single group, so it runs as one task on one core — measured live on TPC-H q2 at SF100, where four stages coalesced to 1 partition and the cluster sat at one busy core per executor with eight of nine slots idle. Bytes were small; the work over them was not, and coalescing cannot see that.

Callers pass the live slot count so the floor tracks the actual cluster. This mirrors Spark’s coalescePartitions.parallelismFirst, which shrinks the advisory partition size for the same reason.

The floor is advisory in one direction only: it never raises the partition count above what the stage already has, because coalescing may only merge.

Source

pub fn min_partitions(&self) -> usize

Return the configured parallelism floor.

Source

pub fn advise(&self, stats: &[RuntimeStats]) -> CoalesceAdvice

Compute coalesce advice from per-partition stats, without modifying the plan.

Partitions are sorted by memory_bytes (ascending) before grouping so that all small partitions cluster together regardless of their original execution order. Without sorting, a large partition sitting between two small ones would prevent them from coalescing (Spark’s AQE sorts before coalescing for the same reason). Each group of small partitions is capped at target_partition_bytes. Large partitions are always singleton groups.

Each group contains the original partition indices (not sorted indices), so callers can map groups back to the original execution order.

Example: [small(0), big(1), small(2)][[0,2], [1]] (2 groups) vs. the old consecutive-only approach: [[0], [1], [2]] (3 groups, no gain)

Trait Implementations§

Source§

impl AqeRule for CoalesceRule

Source§

fn apply( &self, plan: &PhysicalPlan, stats: &[RuntimeStats], ) -> Option<PhysicalPlan>

Compute coalesce advice and, when beneficial, rewrite the plan.

When advise() produces fewer groups than the current partition count, stamps coalesced_partition_count on the plan and appends a NodeOp::CoalescePartitions node carrying the computed target count.

Source§

fn name(&self) -> &str

Short, stable rule name used in explain and diagnostics output.

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Allocation for T
where T: RefUnwindSafe + Send + Sync,

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self>

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self>

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> IntoEither for T

Source§

fn into_either(self, into_left: bool) -> Either<Self, Self>

Converts 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 more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
where F: FnOnce(&Self) -> bool,

Converts 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
Source§

impl<T> IntoRequest<T> for T

Source§

fn into_request(self) -> Request<T>

Wrap the input message T in a tonic::Request
Source§

impl<L> LayerExt<L> for L

Source§

fn named_layer<S>(&self, service: S) -> Layered<<L as Layer<S>>::Service, S>
where L: Layer<S>,

Applies the layer to a service and wraps it in Layered.
Source§

impl<T> Pointable for T

Source§

const ALIGN: usize

The alignment of pointer.
Source§

type Init = T

The type for initializers.
Source§

unsafe fn init(init: <T as Pointable>::Init) -> usize

Initializes a with the given initializer. Read more
Source§

unsafe fn deref<'a>(ptr: usize) -> &'a T

Dereferences the given pointer. Read more
Source§

unsafe fn deref_mut<'a>(ptr: usize) -> &'a mut T

Mutably dereferences the given pointer. Read more
Source§

unsafe fn drop(ptr: usize)

Drops the object pointed to by the given pointer. Read more
Source§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = Infallible

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self>
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self>

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more