Skip to main content

ConsumerGroup

Struct ConsumerGroup 

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

Joined Kafka consumer group member.

Implementations§

Source§

impl ConsumerGroup

Source

pub fn group_id(&self) -> &str

Returns the Kafka consumer group ID.

Source

pub fn member_id(&self) -> &str

Returns the broker-assigned group member ID.

Source

pub fn generation_id(&self) -> i32

Returns the current group generation ID.

Source

pub fn metadata(&self) -> ConsumerGroupMetadata

Snapshots the current identity for a fenced transactional offset commit.

Take a fresh snapshot after every rejoin because generation and member identity change during rebalancing.

Source

pub fn assignments(&self) -> &[ConsumerAssignment]

Returns the assigned topic partitions and next offsets.

Source

pub fn position(&self, topic: &str, partition: i32) -> Option<i64>

Returns the next offset for a currently assigned topic partition.

Source

pub fn seek(&mut self, topic: &str, partition: i32, offset: i64) -> Result<()>

Changes the next offset for a currently assigned topic partition.

A later group rejoin restores broker-committed or configured reset offsets for the new assignment.

Source

pub fn pause(&mut self, topic: &str, partition: i32) -> Result<()>

Pauses fetching from a currently assigned topic partition.

Pause state is retained across a rejoin when this member keeps the topic partition.

Source

pub fn resume(&mut self, topic: &str, partition: i32) -> Result<()>

Resumes fetching from a currently assigned topic partition.

Source

pub async fn fetch_watermarks( &mut self, topic: impl Into<String>, partition: i32, ) -> Result<PartitionWatermarks>

Fetches the earliest and latest available offsets for a topic partition.

The partition does not need to be assigned to this group member.

Source

pub async fn leave(self) -> Result<()>

Leaves the consumer group and consumes this member handle.

Stop any separately spawned ConsumerGroupHeartbeat before leaving. Kafka receives both the broker member ID and the configured static instance ID, so the member does not remain active until session expiry.

Source

pub async fn poll(&mut self) -> Result<Vec<ConsumerRecord>>

Sends a heartbeat, polls assigned partitions, and advances in-memory offsets.

Source

pub async fn poll_with_heartbeat( &mut self, heartbeat: &mut ConsumerGroupHeartbeat, ) -> Result<Vec<ConsumerRecord>>

Checks a background heartbeat task before polling assigned partitions.

If the background heartbeat task has completed with a group error that requires a rejoin, this method rejoins the group before polling and replaces the handle with a new task using the same interval. Other background heartbeat errors are returned to the caller. A foreground poll rejoin also replaces the stale task before this method returns. When the task is still running, this behaves like ConsumerGroup::poll.

Source

pub async fn spawn_heartbeat_task( &self, interval: Duration, ) -> Result<ConsumerGroupHeartbeat>

Starts a background heartbeat task for this joined group member.

Source

pub async fn heartbeat(&mut self) -> Result<()>

Sends an explicit heartbeat for this group member.

Source

pub async fn commit_offsets(&mut self) -> Result<()>

Commits the current assignment offsets to the group coordinator.

Trait Implementations§

Source§

impl Debug for ConsumerGroup

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. 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<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> 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