Skip to main content

eos/
eos.rs

1//! Consume, produce, and commit offsets in one transaction (exactly-once).
2//!
3//! Broker on `KAFKA_BOOTSTRAP`. Source `KAFKA_TOPIC` (default `partitionline`).
4//! Destination `KAFKA_OUTPUT_TOPIC` (default `partitionline-out`). Group
5//! `KAFKA_GROUP` (default `partitionline-eos`). `KAFKA_TRANSACTIONAL_ID`
6//! defaults to `partitionline-eos`. Auto-commit stays off; offsets go through
7//! [`partitionline::Producer::send_offsets_for_group`].
8
9use 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}