Skip to main content

kip848/
kip848.rs

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