pub struct KafkaSource;Expand description
Kafka source entry points.
Implementations§
Source§impl KafkaSource
impl KafkaSource
Sourcepub fn plain(
settings: KafkaConsumerSettings,
subscription: Subscription,
) -> Source<ConsumerRecord, KafkaControl>
pub fn plain( settings: KafkaConsumerSettings, subscription: Subscription, ) -> Source<ConsumerRecord, KafkaControl>
Creates a plain source of owned Kafka records.
Use KafkaSource::committable when records should advance consumer-group
offsets after downstream checkpoints.
Sourcepub fn committable(
settings: KafkaConsumerSettings,
subscription: Subscription,
) -> SourceWithContext<ConsumerRecord, KafkaOffset, KafkaControl>
pub fn committable( settings: KafkaConsumerSettings, subscription: Subscription, ) -> SourceWithContext<ConsumerRecord, KafkaOffset, KafkaControl>
Creates a source whose records carry committable Kafka offsets as context.
Call KafkaOffset::commit only after the
corresponding downstream checkpoint or side effect succeeds.
Sourcepub fn committable_payload_batches(
settings: KafkaConsumerSettings,
subscription: Subscription,
) -> Source<KafkaPayloadBatch, KafkaControl>
pub fn committable_payload_batches( settings: KafkaConsumerSettings, subscription: Subscription, ) -> Source<KafkaPayloadBatch, KafkaControl>
Creates a payload-only batched source with committable partition watermarks.
This is the throughput-oriented counterpart to KafkaSource::committable.
It emits one Datum element per Kafka poll batch, copies payload bytes into a
batch buffer, and exposes a single KafkaBatchOffset for the highest
offset in each touched topic/partition. Downstream code should call
KafkaPayloadBatch::commit only after the whole batch checkpoint succeeds.