Skip to main content

SourceConsumer

Struct SourceConsumer 

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

A Source implementation backed by a Kafka consumer on the source cluster.

Wraps a Consumer and translates each crabka_client_consumer::ConsumerRecord into a ReplicatedRecord that carries the full envelope (topic, partition, offset, timestamp, headers) alongside the raw payload. The connect runtime never sees topic/partition directly — only the ReplicatedRecord value.

Implementations§

Source§

impl SourceConsumer

Source

pub async fn start( bootstrap: &str, group_id: &str, topics: &[String], security: Option<ClientSecurity>, ) -> Result<Self, ConnectError>

Build and start a SourceConsumer subscribed to topics on the cluster at bootstrap, joining group_id.

Offsets reset to earliest (no previously committed offset for the group). Pass security when the source cluster requires authentication/TLS.

§Errors

Returns ConnectError::Backend if the consumer cannot join the group.

Trait Implementations§

Source§

impl Source<(), ReplicatedRecord> for SourceConsumer

Source§

fn poll<'life0, 'async_trait>( &'life0 mut self, ) -> Pin<Box<dyn Future<Output = Result<Option<ConnectRecord<(), ReplicatedRecord>>, ConnectError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Poll the source cluster for the next record.

Returns Ok(None) when the consumer is momentarily caught up (the runtime should back off and retry). Returns Ok(Some(_)) with the next ReplicatedRecord otherwise.

§Errors

Returns ConnectError::Backend if the underlying consumer poll fails.

Source§

fn checkpoint(&self) -> Option<SourceOffset>

Snapshot the current read positions for all partitions seen so far.

Returns None before the first successful poll (nothing to commit yet).

Source§

fn seek<'life0, 'async_trait>( &'life0 mut self, offset: SourceOffset, ) -> Pin<Box<dyn Future<Output = Result<(), ConnectError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Restore the read position from a previously-checkpointed SourceOffset.

The runtime calls this once before the first poll, passing the position loaded from the durable checkpoint store on the target. Each position entry is keyed "<topic>-<partition>"OffsetValue::Long(next_offset) (the value checkpoint wrote: last_consumed + 1). We decode each key back into (topic, partition) and hand the offset to the consumer’s seek.

The consumer holds each seek as pending and materialises it at the top of the first poll that sees the partition assigned — after the group’s post-assignment offset prime, but before any Fetch — so the sought offset is the one fetched. That makes restart resume from the last fully-committed record rather than re-reading the topic from offset 0: no record below the sought offset is re-delivered, and none above it is skipped (no data gap). Delivery remains at-least-once — a crash between a sink flush and the checkpoint save can re-deliver the in-flight batch, but never lose a record.

A malformed key (no -, or a non-integer partition/offset) is skipped with a warning rather than failing the restore: one corrupt entry must not strand recovery for the partitions that decoded cleanly.

§Errors

Returns ConnectError::Backend if the consumer is already closed.

Source§

fn close<'life0, 'async_trait>( &'life0 mut self, ) -> Pin<Box<dyn Future<Output = Result<(), ConnectError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Close the underlying consumer, sending LeaveGroup so a restarted replicator can rejoin the group immediately instead of waiting out the departed member’s session timeout.

§Errors

Returns ConnectError::Backend if the consumer fails to close cleanly.

Source§

fn acknowledge<'life0, 'life1, 'async_trait>( &'life0 mut self, _offset: &'life1 SourceOffset, ) -> Pin<Box<dyn Future<Output = Result<(), ConnectError>> + Send + 'async_trait>>
where 'life0: 'async_trait, 'life1: 'async_trait, Self: 'async_trait,

Acknowledge that offset is durable end-to-end. Read more

Auto Trait Implementations§

Blanket Implementations§

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<'a, T, E> AsTaggedExplicit<'a, E> for T
where T: 'a,

Source§

fn explicit(self, class: Class, tag: u32) -> TaggedParser<'a, Explicit, Self, E>

Source§

impl<'a, T, E> AsTaggedImplicit<'a, E> for T
where T: 'a,

Source§

fn implicit( self, class: Class, constructed: bool, tag: u32, ) -> TaggedParser<'a, Implicit, Self, E>

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> PolicyExt for T
where T: ?Sized,

Source§

fn and<P, B, E>(self, other: P) -> And<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow only if self and other return Action::Follow. Read more
Source§

fn or<P, B, E>(self, other: P) -> Or<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow if either self or other returns Action::Follow. 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<V, T> VZip<V> for T
where V: MultiLane<T>,

Source§

fn vzip(self) -> V

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