1use 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}