Skip to main content

consume_intercept/
consume_intercept.rs

1//! Drop heartbeat records before they reach the caller.
2//!
3//! Broker on `KAFKA_BOOTSTRAP`. Topic `KAFKA_TOPIC` (default `partitionline`).
4
5use partitionline::{Consumer, ConsumerConfig, ConsumerInterceptor, FetchedRecord};
6
7struct SkipHeartbeats;
8
9impl ConsumerInterceptor for SkipHeartbeats {
10    fn on_consume(&self, recs: Vec<FetchedRecord>) -> Vec<FetchedRecord> {
11        recs.into_iter()
12            .filter(|rec| {
13                rec.last_header("kind").and_then(|h| h.value.as_deref()) != Some(b"heartbeat")
14            })
15            .collect()
16    }
17}
18
19#[tokio::main]
20async fn main() -> partitionline::Result<()> {
21    let bootstrap = std::env::var("KAFKA_BOOTSTRAP").unwrap_or_else(|_| "127.0.0.1:9092".into());
22    let topic = std::env::var("KAFKA_TOPIC").unwrap_or_else(|_| "partitionline".into());
23    let mut consumer =
24        Consumer::new(ConsumerConfig::bootstrap([bootstrap]).interceptor(SkipHeartbeats)).await?;
25    consumer.assign_topic(topic, 0).await?;
26    loop {
27        for rec in consumer.fetch().await? {
28            println!(
29                "{}-{}@{} bytes={}",
30                rec.topic,
31                rec.partition,
32                rec.offset,
33                rec.value.as_ref().map(|v| v.len()).unwrap_or(0)
34            );
35        }
36    }
37}