Skip to main content

cooperative/
cooperative.rs

1//! Join a cooperative-sticky group (KIP-429) and print records.
2//!
3//! `KAFKA_TOPIC` may be a single name or a comma-separated list.
4
5use partitionline::{ConsumerConfig, ConsumerGroup};
6
7#[tokio::main]
8async fn main() -> partitionline::Result<()> {
9    let bootstrap = std::env::var("KAFKA_BOOTSTRAP").unwrap_or_else(|_| "127.0.0.1:9092".into());
10    let group_id = std::env::var("KAFKA_GROUP").unwrap_or_else(|_| "partitionline".into());
11    let topics: Vec<String> = std::env::var("KAFKA_TOPIC")
12        .unwrap_or_else(|_| "partitionline".into())
13        .split(',')
14        .map(str::trim)
15        .filter(|s| !s.is_empty())
16        .map(str::to_string)
17        .collect();
18    let mut group = ConsumerGroup::join_cooperative_sticky_topics(
19        ConsumerConfig::bootstrap([bootstrap]).max_wait_ms(500),
20        group_id,
21        topics,
22    )
23    .await?;
24    loop {
25        let recs = group.poll().await?;
26        if recs.is_empty() {
27            continue;
28        }
29        for rec in &recs {
30            println!("{}-{}@{}", rec.topic, rec.partition, rec.offset);
31        }
32        group.commit().await?;
33    }
34}