ruststream_rdkafka/
subscriber.rs1use 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
18pub 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 #[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 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}