Skip to main content

ruststream_rdkafka/
subscriber.rs

1//! The subscriber: a stream of Kafka deliveries from one topic subscription.
2
3use std::fmt;
4use std::sync::Arc;
5
6use bytes::Bytes;
7use futures::Stream;
8use rdkafka::Message as _;
9use rdkafka::consumer::StreamConsumer;
10use ruststream::Subscriber;
11
12use crate::convert;
13use crate::error::KafkaError;
14use crate::message::{KafkaMessage, Settlement};
15use crate::topic::Commit;
16use crate::tracker::{CommitTracker, TrackingContext};
17
18/// A consumer-group member on one topic, yielding [`KafkaMessage`] deliveries.
19///
20/// Created by subscribing a [`KafkaTopic`](crate::KafkaTopic) descriptor (or a bare topic name)
21/// through [`KafkaBroker`](crate::KafkaBroker). The subscriber owns a dedicated librdkafka
22/// consumer; dropping it closes the consumer, which leaves the group and (under auto-commit)
23/// commits the final stored position. Under `Commit::Tracked` each in-flight delivery keeps
24/// the consumer alive, so the close happens once the last outstanding message settles or
25/// drops - do not rely on subscriber drop as an immediate group-departure barrier.
26///
27/// Back-pressure: polling the stream is what drives the consumer, so consuming slower simply
28/// fetches slower; librdkafka's own fetch queue bounds (`queued.max.messages.kbytes` and
29/// friends, settable through [`KafkaTopic::config`](crate::KafkaTopic::config)) cap local
30/// buffering.
31pub struct KafkaSubscriber {
32    consumer: Arc<StreamConsumer<TrackingContext>>,
33    topic: String,
34    commit: Commit,
35    tracker: Arc<CommitTracker>,
36}
37
38impl KafkaSubscriber {
39    pub(crate) fn new(
40        consumer: Arc<StreamConsumer<TrackingContext>>,
41        topic: String,
42        commit: Commit,
43        tracker: Arc<CommitTracker>,
44    ) -> Self {
45        Self {
46            consumer,
47            topic,
48            commit,
49            tracker,
50        }
51    }
52
53    /// The topic this subscriber consumes.
54    #[must_use]
55    pub fn topic(&self) -> &str {
56        &self.topic
57    }
58
59    fn map_delivery(&self, delivery: &rdkafka::message::BorrowedMessage<'_>) -> KafkaMessage {
60        let headers = convert::headers_from_message(delivery);
61        let payload = delivery
62            .payload()
63            .map_or_else(Bytes::new, Bytes::copy_from_slice);
64        let settlement = match self.commit {
65            Commit::Auto => Settlement::Advisory,
66            Commit::Tracked => {
67                self.tracker
68                    .delivered(delivery.topic(), delivery.partition(), delivery.offset());
69                Settlement::Tracked {
70                    consumer: Arc::clone(&self.consumer),
71                    tracker: Arc::clone(&self.tracker),
72                }
73            }
74        };
75        KafkaMessage::new(
76            payload,
77            headers,
78            delivery.topic().to_owned(),
79            delivery.partition(),
80            delivery.offset(),
81            delivery.timestamp().to_millis(),
82            settlement,
83        )
84    }
85}
86
87impl fmt::Debug for KafkaSubscriber {
88    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
89        f.debug_struct("KafkaSubscriber")
90            .field("topic", &self.topic)
91            .field("commit", &self.commit)
92            .finish_non_exhaustive()
93    }
94}
95
96impl Subscriber for KafkaSubscriber {
97    type Message = KafkaMessage;
98    type Error = KafkaError;
99
100    /// Streams deliveries as they arrive; the stream yields an error item when the consumer
101    /// fails (it does not end on its own - drop the subscriber to leave the group).
102    ///
103    /// # Cancel safety
104    ///
105    /// Polling is cancel safe (the underlying `recv` is documented cancellation safe, so no
106    /// delivery is lost by dropping the stream between polls), and the stream can be re-created
107    /// by calling `stream` again: deliveries buffer in the consumer, not in the returned stream.
108    fn stream(&mut self) -> impl Stream<Item = Result<Self::Message, Self::Error>> + Send + '_ {
109        futures::stream::unfold(self, |sub| async move {
110            let item = match sub.consumer.recv().await {
111                Ok(delivery) => Ok(sub.map_delivery(&delivery)),
112                Err(err) => Err(KafkaError::consume(err)),
113            };
114            Some((item, sub))
115        })
116    }
117}