pub struct WorkerCompute { /* private fields */ }Expand description
Runs wasm / container steps on a shared fv_compute::Runtime.
Implementations§
Source§impl KineticsRunner
impl KineticsRunner
pub fn new(runtime: Arc<Runtime>) -> KineticsRunner
Sourcepub fn selector(step: &Value) -> Result<&str, String>
pub fn selector(step: &Value) -> Result<&str, String>
The transform a compute step selects: its ref (or legacy transform) field.
Sourcepub fn ensure_loadable(&self, step: &Value) -> Result<(), String>
pub fn ensure_loadable(&self, step: &Value) -> Result<(), String>
Resolve and load the step’s transform without running it. Call this when a pipeline or stream is configured, so a missing or broken transform fails there rather than per batch.
Sourcepub fn run_sync(
&self,
step: &Value,
inputs: &[(String, Vec<Row>)],
) -> Result<Vec<Row>, String>
pub fn run_sync( &self, step: &Value, inputs: &[(String, Vec<Row>)], ) -> Result<Vec<Row>, String>
Run a compute step synchronously. The transform backends are synchronous (a component
call, or a blocking child process), so this needs no async runtime; the async
ComputeRunner impl is a thin wrapper.
Trait Implementations§
Source§impl Clone for KineticsRunner
impl Clone for KineticsRunner
Source§fn clone(&self) -> KineticsRunner
fn clone(&self) -> KineticsRunner
Returns a duplicate of the value. Read more
1.0.0 (const: unstable) · Source§fn clone_from(&mut self, source: &Self)
fn clone_from(&mut self, source: &Self)
Performs copy-assignment from
source. Read moreSource§impl ComputeRunner for KineticsRunner
impl ComputeRunner for KineticsRunner
fn run<'life0, 'life1, 'life2, 'async_trait>(
&'life0 self,
step: &'life1 Value,
inputs: &'life2 [(String, Vec<Row>)],
) -> Pin<Box<dyn Future<Output = Result<Vec<Row>, String>> + Send + 'async_trait>>where
'life0: 'async_trait,
'life1: 'async_trait,
'life2: 'async_trait,
KineticsRunner: 'async_trait,
Auto Trait Implementations§
impl !RefUnwindSafe for KineticsRunner
impl !UnwindSafe for KineticsRunner
impl Freeze for KineticsRunner
impl Send for KineticsRunner
impl Sync for KineticsRunner
impl Unpin for KineticsRunner
impl UnsafeUnpin for KineticsRunner
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
Mutably borrows from an owned value. Read more
impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
Source§impl<T> CloneToUninit for Twhere
T: Clone,
impl<T> CloneToUninit for Twhere
T: Clone,
Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
Source§fn in_current_span(self) -> Instrumented<Self> ⓘ
fn in_current_span(self) -> Instrumented<Self> ⓘ
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> ⓘ
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 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> ⓘ
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