1use std::time::Duration;
10
11use partitionline::{
12 ConsumerConfig, ConsumerGroup, ConsumerRecords, IsolationLevel, ProduceRecord, Producer,
13 ProducerConfig,
14};
15
16#[tokio::main]
17async fn main() -> partitionline::Result<()> {
18 let bootstrap = std::env::var("KAFKA_BOOTSTRAP").unwrap_or_else(|_| "127.0.0.1:9092".into());
19 let source = std::env::var("KAFKA_TOPIC").unwrap_or_else(|_| "partitionline".into());
20 let dest = std::env::var("KAFKA_OUTPUT_TOPIC").unwrap_or_else(|_| "partitionline-out".into());
21 let group_id = std::env::var("KAFKA_GROUP").unwrap_or_else(|_| "partitionline-eos".into());
22 let transactional_id =
23 std::env::var("KAFKA_TRANSACTIONAL_ID").unwrap_or_else(|_| "partitionline-eos".into());
24
25 let producer = Producer::new(
26 ProducerConfig::bootstrap([bootstrap.clone()])
27 .transactional_id(transactional_id)
28 .linger(Duration::ZERO),
29 )
30 .await?;
31 producer.init_transactions().await?;
32
33 let mut group = ConsumerGroup::join_topics(
34 ConsumerConfig::bootstrap([bootstrap])
35 .isolation(IsolationLevel::ReadCommitted)
36 .max_wait_ms(500),
37 group_id,
38 [source],
39 )
40 .await?;
41
42 loop {
43 let recs = group.poll().await?;
44 if recs.is_empty() {
45 continue;
46 }
47 if let Err(e) = copy_batch(&producer, &group, &recs, &dest).await {
48 drop(producer.abort_transaction().await);
49 group.leave().await?;
50 producer.close().await?;
51 return Err(e);
52 }
53 for rec in &recs {
54 println!("{}-{}@{} -> {dest}", rec.topic, rec.partition, rec.offset);
55 }
56 }
57}
58
59async fn copy_batch(
60 producer: &Producer,
61 group: &ConsumerGroup,
62 recs: &ConsumerRecords,
63 dest: &str,
64) -> partitionline::Result<()> {
65 producer.begin_transaction().await?;
66 for rec in recs {
67 let mut out = ProduceRecord::to(dest);
68 if let Some(key) = rec.key.clone() {
69 out = out.key(key);
70 }
71 if let Some(value) = rec.value.clone() {
72 out = out.value(value);
73 }
74 drop(producer.send(out).await?);
75 }
76 producer
77 .send_offsets_for_group(&group.group_metadata(), recs.next_offsets())
78 .await?;
79 producer.commit_transaction().await?;
80 Ok(())
81}