pub struct WorkerPool<W: Worker> {
pub task_name: String,
/* private fields */
}Expand description
A pool of Worker’s that process checkpoints concurrently.
This struct manages a collection of workers that process checkpoints in
parallel. It handles checkpoint distribution, progress tracking, and
graceful shutdown. It can optionally use a Reducer to aggregate and
process worker Messages.
§Examples
§Direct Processing (Without Batching)
use std::sync::Arc;
use async_trait::async_trait;
use iota_data_ingestion_core::{Worker, WorkerPool};
use iota_types::full_checkpoint_content::{CheckpointData, CheckpointTransaction};
struct DirectProcessor {
// generic Database client.
client: Arc<DatabaseClient>,
}
#[async_trait]
impl Worker for DirectProcessor {
type Message = ();
type Error = DatabaseError;
async fn process_checkpoint(
&self,
checkpoint: Arc<CheckpointData>,
) -> Result<Self::Message, Self::Error> {
// extract a particulat transaction we care about.
let tx: CheckpointTransaction = extract_transaction(checkpoint.as_ref());
// store the transaction in our database of choice.
self.client.store_transaction(&tx).await?;
Ok(())
}
}
// configure worker pool for direct processing.
let processor = DirectProcessor {
client: Arc::new(DatabaseClient::new()),
};
let pool = WorkerPool::new(processor, "direct_processing".into(), 5, Default::default());§Batch Processing (With Reducer)
use std::sync::Arc;
use async_trait::async_trait;
use iota_data_ingestion_core::{Reducer, Worker, WorkerPool};
use iota_types::full_checkpoint_content::{CheckpointData, CheckpointTransaction};
// worker that accumulates transactions for batch processing.
struct BatchProcessor;
#[async_trait]
impl Worker for BatchProcessor {
type Message = Vec<CheckpointTransaction>;
type Error = DatabaseError;
async fn process_checkpoint(
&self,
checkpoint: Arc<CheckpointData>,
) -> Result<Self::Message, Self::Error> {
// collect all checkpoint transactions for batch processing.
Ok(checkpoint.transactions.clone())
}
}
// batch reducer for efficient storage.
struct TransactionBatchReducer {
batch_size: usize,
// generic Database client.
client: Arc<DatabaseClient>,
}
#[async_trait]
impl Reducer<BatchProcessor> for TransactionBatchReducer {
async fn commit(&self, batch: &[Vec<CheckpointTransaction>]) -> Result<(), DatabaseError> {
let flattened: Vec<CheckpointTransaction> = batch.iter().flatten().cloned().collect();
// store the transaction batch in the database of choice.
self.client.store_transactions_batch(&flattened).await?;
Ok(())
}
fn should_close_batch(
&self,
batch: &[Vec<CheckpointTransaction>],
_: Option<&Vec<CheckpointTransaction>>,
) -> bool {
batch.iter().map(|b| b.len()).sum::<usize>() >= self.batch_size
}
}
// configure worker pool with batch processing.
let processor = BatchProcessor;
let reducer = TransactionBatchReducer {
batch_size: 1000,
client: Arc::new(DatabaseClient::new()),
};
let pool = WorkerPool::new_with_reducer(
processor,
"batch_processing".into(),
5,
Default::default(),
reducer,
);Fields§
§task_name: StringAn unique name of the WorkerPool task.
Implementations§
Source§impl<W: Worker + 'static> WorkerPool<W>
impl<W: Worker + 'static> WorkerPool<W>
Sourcepub fn new(
worker: W,
task_name: String,
concurrency: usize,
backoff: ExponentialBackoff,
) -> Self
pub fn new( worker: W, task_name: String, concurrency: usize, backoff: ExponentialBackoff, ) -> Self
Creates a new WorkerPool without a reducer.
Sourcepub fn new_with_reducer<R>(
worker: W,
task_name: String,
concurrency: usize,
backoff: ExponentialBackoff,
reducer: R,
) -> Selfwhere
R: Reducer<W> + 'static,
pub fn new_with_reducer<R>(
worker: W,
task_name: String,
concurrency: usize,
backoff: ExponentialBackoff,
reducer: R,
) -> Selfwhere
R: Reducer<W> + 'static,
Creates a new WorkerPool with a reducer.
Sourcepub async fn run(
self,
watermark: CheckpointSequenceNumber,
checkpoint_receiver: Receiver<Arc<CheckpointData>>,
pool_status_sender: Sender<WorkerPoolStatus>,
token: CancellationToken,
)
pub async fn run( self, watermark: CheckpointSequenceNumber, checkpoint_receiver: Receiver<Arc<CheckpointData>>, pool_status_sender: Sender<WorkerPoolStatus>, token: CancellationToken, )
Runs the worker pool main logic.
Auto Trait Implementations§
impl<W> Freeze for WorkerPool<W>
impl<W> !RefUnwindSafe for WorkerPool<W>
impl<W> Send for WorkerPool<W>
impl<W> Sync for WorkerPool<W>
impl<W> Unpin for WorkerPool<W>
impl<W> !UnwindSafe for WorkerPool<W>
Blanket Implementations§
§impl<U> As for U
impl<U> As for U
§fn as_<T>(self) -> Twhere
T: CastFrom<U>,
fn as_<T>(self) -> Twhere
T: CastFrom<U>,
Casts
self to type T. The semantics of numeric casting with the as operator are followed, so <T as As>::as_::<U> can be used in the same way as T as U for numeric conversions. Read more§impl<'a, T, E> AsTaggedExplicit<'a, E> for Twhere
T: 'a,
impl<'a, T, E> AsTaggedExplicit<'a, E> for Twhere
T: 'a,
§impl<'a, T, E> AsTaggedImplicit<'a, E> for Twhere
T: 'a,
impl<'a, T, E> AsTaggedImplicit<'a, E> for Twhere
T: 'a,
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<T> Conv for T
impl<T> Conv for T
§impl<T> FmtForward for T
impl<T> FmtForward for T
§fn fmt_binary(self) -> FmtBinary<Self>where
Self: Binary,
fn fmt_binary(self) -> FmtBinary<Self>where
Self: Binary,
Causes
self to use its Binary implementation when Debug-formatted.§fn fmt_display(self) -> FmtDisplay<Self>where
Self: Display,
fn fmt_display(self) -> FmtDisplay<Self>where
Self: Display,
Causes
self to use its Display implementation when
Debug-formatted.§fn fmt_lower_exp(self) -> FmtLowerExp<Self>where
Self: LowerExp,
fn fmt_lower_exp(self) -> FmtLowerExp<Self>where
Self: LowerExp,
Causes
self to use its LowerExp implementation when
Debug-formatted.§fn fmt_lower_hex(self) -> FmtLowerHex<Self>where
Self: LowerHex,
fn fmt_lower_hex(self) -> FmtLowerHex<Self>where
Self: LowerHex,
Causes
self to use its LowerHex implementation when
Debug-formatted.§fn fmt_octal(self) -> FmtOctal<Self>where
Self: Octal,
fn fmt_octal(self) -> FmtOctal<Self>where
Self: Octal,
Causes
self to use its Octal implementation when Debug-formatted.§fn fmt_pointer(self) -> FmtPointer<Self>where
Self: Pointer,
fn fmt_pointer(self) -> FmtPointer<Self>where
Self: Pointer,
Causes
self to use its Pointer implementation when
Debug-formatted.§fn fmt_upper_exp(self) -> FmtUpperExp<Self>where
Self: UpperExp,
fn fmt_upper_exp(self) -> FmtUpperExp<Self>where
Self: UpperExp,
Causes
self to use its UpperExp implementation when
Debug-formatted.§fn fmt_upper_hex(self) -> FmtUpperHex<Self>where
Self: UpperHex,
fn fmt_upper_hex(self) -> FmtUpperHex<Self>where
Self: UpperHex,
Causes
self to use its UpperHex implementation when
Debug-formatted.§fn fmt_list(self) -> FmtList<Self>where
&'a Self: for<'a> IntoIterator,
fn fmt_list(self) -> FmtList<Self>where
&'a Self: for<'a> IntoIterator,
Formats each item in a sequence. Read more
§impl<T> Instrument for T
impl<T> Instrument for T
§fn instrument(self, span: Span) -> Instrumented<Self>
fn instrument(self, span: Span) -> Instrumented<Self>
§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§impl<T> IntoRequest<T> for T
impl<T> IntoRequest<T> for T
§fn into_request(self) -> Request<T>
fn into_request(self) -> Request<T>
Wrap the input message
T in a Request§impl<T> IntoRequest<T> for T
impl<T> IntoRequest<T> for T
§fn into_request(self) -> Request<T>
fn into_request(self) -> Request<T>
Wrap the input message
T in a tonic::Request§impl<L> LayerExt<L> for L
impl<L> LayerExt<L> for L
§fn named_layer<S>(&self, service: S) -> Layered<<L as Layer<S>>::Service, S>where
L: Layer<S>,
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].§impl<T> Pipe for Twhere
T: ?Sized,
impl<T> Pipe for Twhere
T: ?Sized,
§fn pipe<R>(self, func: impl FnOnce(Self) -> R) -> Rwhere
Self: Sized,
fn pipe<R>(self, func: impl FnOnce(Self) -> R) -> Rwhere
Self: Sized,
Pipes by value. This is generally the method you want to use. Read more
§fn pipe_ref<'a, R>(&'a self, func: impl FnOnce(&'a Self) -> R) -> Rwhere
R: 'a,
fn pipe_ref<'a, R>(&'a self, func: impl FnOnce(&'a Self) -> R) -> Rwhere
R: 'a,
Borrows
self and passes that borrow into the pipe function. Read more§fn pipe_ref_mut<'a, R>(&'a mut self, func: impl FnOnce(&'a mut Self) -> R) -> Rwhere
R: 'a,
fn pipe_ref_mut<'a, R>(&'a mut self, func: impl FnOnce(&'a mut Self) -> R) -> Rwhere
R: 'a,
Mutably borrows
self and passes that borrow into the pipe function. Read more§fn pipe_borrow<'a, B, R>(&'a self, func: impl FnOnce(&'a B) -> R) -> R
fn pipe_borrow<'a, B, R>(&'a self, func: impl FnOnce(&'a B) -> R) -> R
§fn pipe_borrow_mut<'a, B, R>(
&'a mut self,
func: impl FnOnce(&'a mut B) -> R,
) -> R
fn pipe_borrow_mut<'a, B, R>( &'a mut self, func: impl FnOnce(&'a mut B) -> R, ) -> R
§fn pipe_as_ref<'a, U, R>(&'a self, func: impl FnOnce(&'a U) -> R) -> R
fn pipe_as_ref<'a, U, R>(&'a self, func: impl FnOnce(&'a U) -> R) -> R
Borrows
self, then passes self.as_ref() into the pipe function.§fn pipe_as_mut<'a, U, R>(&'a mut self, func: impl FnOnce(&'a mut U) -> R) -> R
fn pipe_as_mut<'a, U, R>(&'a mut self, func: impl FnOnce(&'a mut U) -> R) -> R
Mutably borrows
self, then passes self.as_mut() into the pipe
function.§fn pipe_deref<'a, T, R>(&'a self, func: impl FnOnce(&'a T) -> R) -> R
fn pipe_deref<'a, T, R>(&'a self, func: impl FnOnce(&'a T) -> R) -> R
Borrows
self, then passes self.deref() into the pipe function.§impl<T> Pointable for T
impl<T> Pointable for T
§impl<T> PolicyExt for Twhere
T: ?Sized,
impl<T> PolicyExt for Twhere
T: ?Sized,
§impl<T> Tap for T
impl<T> Tap for T
§fn tap_borrow<B>(self, func: impl FnOnce(&B)) -> Self
fn tap_borrow<B>(self, func: impl FnOnce(&B)) -> Self
Immutable access to the
Borrow<B> of a value. Read more§fn tap_borrow_mut<B>(self, func: impl FnOnce(&mut B)) -> Self
fn tap_borrow_mut<B>(self, func: impl FnOnce(&mut B)) -> Self
Mutable access to the
BorrowMut<B> of a value. Read more§fn tap_ref<R>(self, func: impl FnOnce(&R)) -> Self
fn tap_ref<R>(self, func: impl FnOnce(&R)) -> Self
Immutable access to the
AsRef<R> view of a value. Read more§fn tap_ref_mut<R>(self, func: impl FnOnce(&mut R)) -> Self
fn tap_ref_mut<R>(self, func: impl FnOnce(&mut R)) -> Self
Mutable access to the
AsMut<R> view of a value. Read more§fn tap_deref<T>(self, func: impl FnOnce(&T)) -> Self
fn tap_deref<T>(self, func: impl FnOnce(&T)) -> Self
Immutable access to the
Deref::Target of a value. Read more§fn tap_deref_mut<T>(self, func: impl FnOnce(&mut T)) -> Self
fn tap_deref_mut<T>(self, func: impl FnOnce(&mut T)) -> Self
Mutable access to the
Deref::Target of a value. Read more§fn tap_dbg(self, func: impl FnOnce(&Self)) -> Self
fn tap_dbg(self, func: impl FnOnce(&Self)) -> Self
Calls
.tap() only in debug builds, and is erased in release builds.§fn tap_mut_dbg(self, func: impl FnOnce(&mut Self)) -> Self
fn tap_mut_dbg(self, func: impl FnOnce(&mut Self)) -> Self
Calls
.tap_mut() only in debug builds, and is erased in release
builds.§fn tap_borrow_dbg<B>(self, func: impl FnOnce(&B)) -> Self
fn tap_borrow_dbg<B>(self, func: impl FnOnce(&B)) -> Self
Calls
.tap_borrow() only in debug builds, and is erased in release
builds.§fn tap_borrow_mut_dbg<B>(self, func: impl FnOnce(&mut B)) -> Self
fn tap_borrow_mut_dbg<B>(self, func: impl FnOnce(&mut B)) -> Self
Calls
.tap_borrow_mut() only in debug builds, and is erased in release
builds.§fn tap_ref_dbg<R>(self, func: impl FnOnce(&R)) -> Self
fn tap_ref_dbg<R>(self, func: impl FnOnce(&R)) -> Self
Calls
.tap_ref() only in debug builds, and is erased in release
builds.§fn tap_ref_mut_dbg<R>(self, func: impl FnOnce(&mut R)) -> Self
fn tap_ref_mut_dbg<R>(self, func: impl FnOnce(&mut R)) -> Self
Calls
.tap_ref_mut() only in debug builds, and is erased in release
builds.§fn tap_deref_dbg<T>(self, func: impl FnOnce(&T)) -> Self
fn tap_deref_dbg<T>(self, func: impl FnOnce(&T)) -> Self
Calls
.tap_deref() only in debug builds, and is erased in release
builds.