laminar_core/streaming/sink/
mod.rs1use std::sync::Arc;
4
5use arrow::datatypes::SchemaRef;
6use tokio::sync::broadcast;
7
8use super::channel::AsyncConsumer;
9use super::source::{Record, SourceMessage};
10use super::subscription::Subscription;
11
12const DEFAULT_BROADCAST_CAPACITY: usize = 2048;
13
14pub struct Sink<T: Record> {
23 broadcast_tx: broadcast::Sender<SourceMessage<T>>,
24 schema: SchemaRef,
25}
26
27impl<T: Record> Sink<T> {
28 pub(crate) fn new(consumer: AsyncConsumer<SourceMessage<T>>, schema: SchemaRef) -> Self {
29 let (broadcast_tx, _) = broadcast::channel(DEFAULT_BROADCAST_CAPACITY);
30 let tx = broadcast_tx.clone();
31
32 tokio::spawn(async move {
36 drain_loop(consumer, tx).await;
37 });
38
39 Self {
40 broadcast_tx,
41 schema,
42 }
43 }
44
45 #[must_use]
47 pub fn subscribe(&self) -> Subscription<T> {
48 Subscription::new(self.broadcast_tx.subscribe(), Arc::clone(&self.schema))
49 }
50
51 #[must_use]
53 pub fn schema(&self) -> SchemaRef {
54 Arc::clone(&self.schema)
55 }
56
57 #[must_use]
59 pub fn subscriber_count(&self) -> usize {
60 self.broadcast_tx.receiver_count()
61 }
62}
63
64impl<T: Record + std::fmt::Debug> std::fmt::Debug for Sink<T> {
65 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
66 f.debug_struct("Sink")
67 .field("subscribers", &self.subscriber_count())
68 .finish()
69 }
70}
71
72async fn drain_loop<T: Record>(
73 mut consumer: AsyncConsumer<SourceMessage<T>>,
74 tx: broadcast::Sender<SourceMessage<T>>,
75) {
76 while let Ok(msg) = consumer.recv().await {
77 let _ = tx.send(msg);
78 }
79}
80
81#[cfg(test)]
82mod tests;