pub struct KafkaControl { /* private fields */ }Expand description
Materialized control for a Kafka source.
Implementations§
Source§impl KafkaControl
impl KafkaControl
pub fn metrics(&self) -> KafkaMetrics
pub fn is_draining(&self) -> bool
Sourcepub fn drain_and_shutdown(&self, timeout: Duration) -> MqResult<()>
pub fn drain_and_shutdown(&self, timeout: Duration) -> MqResult<()>
Stops new emission and waits for already emitted offsets to be committed.
The source iterator completes on its next pull once outstanding offsets are
drained. This is the Datum equivalent of the Alpakka DrainingControl
hook; datum-agent jobs can call this before stopping the materialized
graph.
Sourcepub fn shutdown_now(&self)
pub fn shutdown_now(&self)
Stops the source without claiming uncommitted offsets.
Trait Implementations§
Source§impl Clone for KafkaControl
impl Clone for KafkaControl
Source§fn clone(&self) -> KafkaControl
fn clone(&self) -> KafkaControl
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 moreAuto Trait Implementations§
impl Freeze for KafkaControl
impl RefUnwindSafe for KafkaControl
impl Send for KafkaControl
impl Sync for KafkaControl
impl Unpin for KafkaControl
impl UnsafeUnpin for KafkaControl
impl UnwindSafe for KafkaControl
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
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> Message for T
impl<T> Message for T
Source§fn from_boxed(m: BoxedMessage) -> Result<Self, BoxedDowncastErr>
fn from_boxed(m: BoxedMessage) -> Result<Self, BoxedDowncastErr>
Convert a BoxedMessage to this concrete type
Source§fn box_message(self, pid: &ActorId) -> Result<BoxedMessage, BoxedDowncastErr>
fn box_message(self, pid: &ActorId) -> Result<BoxedMessage, BoxedDowncastErr>
Convert this message to a BoxedMessage