Skip to main content

KafkaConsumer

Struct KafkaConsumer 

Source
#[non_exhaustive]
pub struct KafkaConsumer;
Expand description

The public entry point for constructing a Consumer.

Translates org.apache.kafka.clients.consumer.KafkaConsumer, whose public constructors are the only way a Java user obtains a consumer. Java’s class is a thin delegating wrapper: its constructor calls ConsumerDelegateCreator.create(config, keyDeserializer, valueDeserializer) (KafkaConsumer.java:615), which switches on group.protocol and returns either an AsyncKafkaConsumer or a ClassicKafkaConsumer (ConsumerDelegateCreator.java:57-66); its 55 other methods forward verbatim to that delegate.

Rust keeps only the constructor. The delegate is returned directly as Box<dyn Consumer<K, V>> rather than being stored in a wrapper that re-forwards every method, because Consumer is dyn compatible — so the caller holds exactly what Java’s wrapper would have held, and the 55 forwarding bodies carry no behavior to translate. ConsumerDelegate and ConsumerDelegateCreator are out of scope per consumer-threading.md §20 (“collapses to direct Box::new(AsyncKafkaConsumer)”); this type is where that collapsed creator lives.

Implementations§

Source§

impl KafkaConsumer

Source

pub fn new<K, V>( config: ConsumerConfig, key_deserializer: Box<dyn Deserializer<K>>, value_deserializer: Box<dyn Deserializer<V>>, ) -> Result<Box<dyn Consumer<K, V>>, Error>
where K: Send + Sync + 'static, V: Send + Sync + 'static,

Constructs a new Consumer from a configuration and explicit key/value Deserializers.

For group.protocol=consumer (KIP-848), this returns Box::new(AsyncKafkaConsumer::new(...)?) — the production consumer built end-to-end with SubscriptionState, ConsumerMetadata, NetworkClient, every RequestManager, and a single bg task (ConsumerNetworkThread). For group.protocol=classic, returns Error::unsupported_version per consumer-threading.md §20 (classic protocol deferred to a later milestone).

Java passes deserializers via ConsumerConfig reflection; Rust takes them as explicit Box<dyn> parameters (Phase 1 decision not to translate reflection machinery). The consumer wraps them in Arc<Deserializers<K, V>> internally for sharing with Fetcher/FetchCollector (Shape B).

MockConsumer (Phase 3) does NOT come through this factory — it has its own constructor. The factory is for the production consumer only.

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<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<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

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, 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.
§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

§

fn vzip(self) -> V