Skip to main content

StreamsGroupHeartbeatRequestData

Struct StreamsGroupHeartbeatRequestData 

Source
pub struct StreamsGroupHeartbeatRequestData {
Show 18 fields pub group_id: KafkaString, pub member_id: KafkaString, pub member_epoch: i32, pub endpoint_information_epoch: i32, pub instance_id: Option<KafkaString>, pub rack_id: Option<KafkaString>, pub rebalance_timeout_ms: i32, pub topology: Option<Box<Topology>>, pub active_tasks: Option<Vec<TaskIds>>, pub standby_tasks: Option<Vec<TaskIds>>, pub warmup_tasks: Option<Vec<TaskIds>>, pub process_id: Option<KafkaString>, pub user_endpoint: Option<Box<Endpoint>>, pub client_tags: Option<Vec<KeyValue>>, pub task_offsets: Option<Vec<TaskOffset>>, pub task_end_offsets: Option<Vec<TaskOffset>>, pub shutdown_application: bool, pub _unknown_tagged_fields: Vec<RawTaggedField>,
}

Fields§

§group_id: KafkaString

The group identifier.

§member_id: KafkaString

The member ID generated by the streams consumer. The member ID must be kept during the entire lifetime of the streams consumer process.

§member_epoch: i32

The current member epoch; 0 to join the group; -1 to leave the group; -2 to indicate that the static member will rejoin.

§endpoint_information_epoch: i32

The current endpoint epoch of this client, represents the latest endpoint epoch this client received

§instance_id: Option<KafkaString>

null if not provided or if it didn’t change since the last heartbeat; the instance ID for static membership otherwise.

§rack_id: Option<KafkaString>

null if not provided or if it didn’t change since the last heartbeat; the rack ID of the member otherwise.

§rebalance_timeout_ms: i32

-1 if it didn’t change since the last heartbeat; the maximum time in milliseconds that the coordinator will wait on the member to revoke its tasks otherwise.

§topology: Option<Box<Topology>>

The topology metadata of the streams application. Used to initialize the topology of the group and to check if the topology corresponds to the topology initialized for the group. Only sent when memberEpoch = 0, must be non-empty. Null otherwise.

§active_tasks: Option<Vec<TaskIds>>

Currently owned active tasks for this client. Null if unchanged since last heartbeat.

§standby_tasks: Option<Vec<TaskIds>>

Currently owned standby tasks for this client. Null if unchanged since last heartbeat.

§warmup_tasks: Option<Vec<TaskIds>>

Currently owned warm-up tasks for this client. Null if unchanged since last heartbeat.

§process_id: Option<KafkaString>

Identity of the streams instance that may have multiple consumers. Null if unchanged since last heartbeat.

§user_endpoint: Option<Box<Endpoint>>

User-defined endpoint for Interactive Queries. Null if unchanged since last heartbeat, or if not defined on the client.

§client_tags: Option<Vec<KeyValue>>

Used for rack-aware assignment algorithm. Null if unchanged since last heartbeat.

§task_offsets: Option<Vec<TaskOffset>>

Cumulative changelog offsets for tasks. Only updated when a warm-up task has caught up, and according to the task offset interval. Null if unchanged since last heartbeat.

§task_end_offsets: Option<Vec<TaskOffset>>

Cumulative changelog end-offsets for tasks. Only updated when a warm-up task has caught up, and according to the task offset interval. Null if unchanged since last heartbeat.

§shutdown_application: bool

Whether all Streams clients in the group should shut down.

§_unknown_tagged_fields: Vec<RawTaggedField>

Implementations§

Source§

impl StreamsGroupHeartbeatRequestData

Source

pub fn with_group_id(self, value: KafkaString) -> Self

Source

pub fn with_member_id(self, value: KafkaString) -> Self

Source

pub fn with_member_epoch(self, value: i32) -> Self

Source

pub fn with_endpoint_information_epoch(self, value: i32) -> Self

Source

pub fn with_instance_id(self, value: Option<KafkaString>) -> Self

Source

pub fn with_rack_id(self, value: Option<KafkaString>) -> Self

Source

pub fn with_rebalance_timeout_ms(self, value: i32) -> Self

Source

pub fn with_topology(self, value: Option<Box<Topology>>) -> Self

Source

pub fn with_active_tasks(self, value: Option<Vec<TaskIds>>) -> Self

Source

pub fn with_standby_tasks(self, value: Option<Vec<TaskIds>>) -> Self

Source

pub fn with_warmup_tasks(self, value: Option<Vec<TaskIds>>) -> Self

Source

pub fn with_process_id(self, value: Option<KafkaString>) -> Self

Source

pub fn with_user_endpoint(self, value: Option<Box<Endpoint>>) -> Self

Source

pub fn with_client_tags(self, value: Option<Vec<KeyValue>>) -> Self

Source

pub fn with_task_offsets(self, value: Option<Vec<TaskOffset>>) -> Self

Source

pub fn with_task_end_offsets(self, value: Option<Vec<TaskOffset>>) -> Self

Source

pub fn with_shutdown_application(self, value: bool) -> Self

Source

pub fn read(buf: &mut Bytes, version: i16) -> Result<Self>

Source

pub fn write(&self, buf: &mut BytesMut, version: i16) -> Result<()>

Source

pub fn encoded_len(&self, version: i16) -> Result<usize>

Trait Implementations§

Source§

impl Clone for StreamsGroupHeartbeatRequestData

Source§

fn clone(&self) -> StreamsGroupHeartbeatRequestData

Returns a duplicate of the value. Read more
1.0.0 (const: unstable) · Source§

fn clone_from(&mut self, source: &Self)

Performs copy-assignment from source. Read more
Source§

impl Debug for StreamsGroupHeartbeatRequestData

Source§

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

Formats the value using the given formatter. Read more
Source§

impl Default for StreamsGroupHeartbeatRequestData

Source§

fn default() -> Self

Returns the “default value” for a type. Read more
Source§

impl PartialEq for StreamsGroupHeartbeatRequestData

Source§

fn eq(&self, other: &StreamsGroupHeartbeatRequestData) -> bool

Tests for self and other values to be equal, and is used by ==.
1.0.0 (const: unstable) · Source§

fn ne(&self, other: &Rhs) -> bool

Tests for !=. The default implementation is almost always sufficient, and should not be overridden without very good reason.
Source§

impl StructuralPartialEq for StreamsGroupHeartbeatRequestData

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> CloneToUninit for T
where T: Clone,

Source§

unsafe fn clone_to_uninit(&self, dest: *mut u8)

🔬This is a nightly-only experimental API. (clone_to_uninit)
Performs copy-assignment from self to dest. 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> ToOwned for T
where T: Clone,

Source§

type Owned = T

The resulting type after obtaining ownership.
Source§

fn to_owned(&self) -> T

Creates owned data from borrowed data, usually by cloning. Read more
Source§

fn clone_into(&self, target: &mut T)

Uses borrowed data to replace owned data, usually by cloning. Read more
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.