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
impl SourceConsumer
Sourcepub async fn start(
bootstrap: &str,
group_id: &str,
topics: &[String],
security: Option<ClientSecurity>,
) -> Result<Self, ConnectError>
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
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,
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>
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,
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,
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,
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,
offset is durable end-to-end. Read more