Skip to main content

offsets/
offsets.rs

1//! Print beginning / end offsets, lag, and committed metadata.
2//!
3//! Broker on `KAFKA_BOOTSTRAP`. Topic `KAFKA_TOPIC` (default `partitionline`).
4
5use partitionline::{
6    Consumer, ConsumerConfig, ConsumerGroup, OffsetAndMetadata, ProduceRecord, Producer,
7    ProducerConfig, TopicPartition,
8};
9
10#[tokio::main]
11async fn main() -> partitionline::Result<()> {
12    let bootstrap = std::env::var("KAFKA_BOOTSTRAP").unwrap_or_else(|_| "127.0.0.1:9092".into());
13    let topic = std::env::var("KAFKA_TOPIC").unwrap_or_else(|_| "partitionline".into());
14
15    let producer = Producer::new(
16        ProducerConfig::bootstrap([bootstrap.clone()]).linger(std::time::Duration::ZERO),
17    )
18    .await?;
19    let md = producer
20        .send(ProduceRecord::to(topic.clone()).value(&b"hello"[..]))
21        .await?;
22    println!("produced {}-{}@{}", md.topic, md.partition, md.offset);
23    producer.close().await?;
24
25    let mut consumer = Consumer::connect(bootstrap.clone()).await?;
26    consumer.assign(&topic, 0, 0).await?;
27    let begin = consumer.beginning_offsets([(topic.as_str(), 0)]).await?;
28    let end = consumer.end_offsets([(topic.as_str(), 0)]).await?;
29    let lag = consumer.current_lag((topic.as_str(), 0)).await?;
30    println!("begin={begin:?} end={end:?} lag={lag:?}");
31    consumer.close().await?;
32
33    let mut group = ConsumerGroup::join(
34        ConsumerConfig::bootstrap([bootstrap]).max_wait_ms(500),
35        "partitionline-offsets",
36        topic.clone(),
37    )
38    .await?;
39    let recs = group.poll().await?;
40    if let Some(rec) = recs.first() {
41        let tp = TopicPartition::new(&rec.topic, rec.partition);
42        group
43            .commit_with_metadata([(
44                tp,
45                OffsetAndMetadata::with_metadata(rec.offset + 1, "example"),
46            )])
47            .await?;
48    }
49    for (tp, md) in group.committed().await? {
50        println!("{tp} committed={} meta={}", md.offset, md.metadata);
51    }
52    group.leave().await?;
53    Ok(())
54}