consume_intercept/
consume_intercept.rs1use 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}