Skip to main content

pause/
pause.rs

1//! Pause and resume assigned partitions without dropping the assignment.
2//!
3//! Broker on `KAFKA_BOOTSTRAP`. Topic `KAFKA_TOPIC` (default `partitionline`).
4
5use partitionline::Consumer;
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 topic = std::env::var("KAFKA_TOPIC").unwrap_or_else(|_| "partitionline".into());
11    let mut consumer = Consumer::connect(bootstrap).await?;
12    consumer.assign_topic(topic, 0).await?;
13    let assigned = consumer.assignment();
14    println!("assigned {assigned:?}");
15    if let Some(tp) = assigned.first().cloned() {
16        consumer.pause([tp.clone()]);
17        println!("paused {:?}", consumer.paused());
18        let recs = consumer.fetch().await?;
19        println!("while paused: {} records", recs.len());
20        consumer.resume([tp]);
21    }
22    let recs = consumer.fetch().await?;
23    println!("after resume: {} records", recs.len());
24    consumer.close().await?;
25    Ok(())
26}