Skip to main content

FragmentExecutor

Trait FragmentExecutor 

Source
pub trait FragmentExecutor: Send + Sync {
    // Required method
    fn execute<'life0, 'life1, 'async_trait>(
        &'life0 self,
        fragment: &'life1 PlanFragment,
        inputs: Vec<FragmentStream>,
        control: FragmentControl,
    ) -> Pin<Box<dyn Future<Output = DistributedResult<FragmentStream>> + Send + 'async_trait>>
       where Self: 'async_trait,
             'life0: 'async_trait,
             'life1: 'async_trait;

    // Provided method
    fn refill_top_k(
        &self,
        fragment: &PlanFragment,
        offset: usize,
        limit: usize,
        control: FragmentControl,
    ) -> DistributedResult<TopKRefill> { ... }
}
Expand description

Executes one fragment on its worker. Server/engine bindings install an implementation behind RemoteFragmentEndpoint; InMemoryFragmentExecutor remains the deterministic reference executor.

Required Methods§

Source

fn execute<'life0, 'life1, 'async_trait>( &'life0 self, fragment: &'life1 PlanFragment, inputs: Vec<FragmentStream>, control: FragmentControl, ) -> Pin<Box<dyn Future<Output = DistributedResult<FragmentStream>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Runs the fragment. inputs carries one resolved stream per FragmentOperator::RemoteExchangeSource operator, in operator order.

Provided Methods§

Source

fn refill_top_k( &self, fragment: &PlanFragment, offset: usize, limit: usize, control: FragmentControl, ) -> DistributedResult<TopKRefill>

Returns the next limit local top-k candidates after offset (the number already returned), with a tightened unseen bound. Executors without scored streams leave the default, which rejects.

Dyn Compatibility§

This trait is dyn compatible.

In older versions of Rust, dyn compatibility was called "object safety".

Implementors§