1use std::collections::VecDeque;
16use std::pin::Pin;
17use std::sync::Arc;
18use std::task::{Context, Poll};
19use std::time::Instant;
20
21use common_error::ext::BoxedError;
22use common_recordbatch::error::ExternalSnafu;
23use common_recordbatch::{DfRecordBatch, RecordBatch};
24use futures::stream::BoxStream;
25use futures::{Stream, StreamExt};
26use snafu::ResultExt;
27
28use crate::cache::CacheStrategy;
29use crate::error::Result;
30use crate::read::flat_projection::FlatProjectionMapper;
31use crate::read::scan_util::PartitionMetrics;
32use crate::read::series_scan::SeriesBatch;
33
34pub enum ScanBatch {
36 Series(SeriesBatch),
37 RecordBatch(DfRecordBatch),
38}
39
40pub type ScanBatchStream = BoxStream<'static, Result<ScanBatch>>;
41
42pub(crate) struct ConvertBatchStream {
44 inner: ScanBatchStream,
45 projection_mapper: Arc<FlatProjectionMapper>,
46 cache_strategy: CacheStrategy,
47 partition_metrics: PartitionMetrics,
48 pending: VecDeque<RecordBatch>,
49}
50
51impl ConvertBatchStream {
52 pub(crate) fn new(
53 inner: ScanBatchStream,
54 projection_mapper: Arc<FlatProjectionMapper>,
55 cache_strategy: CacheStrategy,
56 partition_metrics: PartitionMetrics,
57 ) -> Self {
58 Self {
59 inner,
60 projection_mapper,
61 cache_strategy,
62 partition_metrics,
63 pending: VecDeque::new(),
64 }
65 }
66
67 fn convert(&mut self, batch: ScanBatch) -> common_recordbatch::error::Result<RecordBatch> {
68 match batch {
69 ScanBatch::Series(series) => {
70 debug_assert!(
71 self.pending.is_empty(),
72 "ConvertBatchStream should not convert a new SeriesBatch when pending batches exist"
73 );
74
75 let SeriesBatch::Flat(flat_batch) = series;
76 for batch in flat_batch.batches {
78 self.pending.push_back(
79 self.projection_mapper
80 .convert(&batch, &self.cache_strategy)?,
81 );
82 }
83
84 let output_schema = self.projection_mapper.output_schema();
85 Ok(self
86 .pending
87 .pop_front()
88 .unwrap_or_else(|| RecordBatch::new_empty(output_schema)))
89 }
90 ScanBatch::RecordBatch(df_record_batch) => {
91 self.projection_mapper
93 .convert(&df_record_batch, &self.cache_strategy)
94 }
95 }
96 }
97}
98
99impl Stream for ConvertBatchStream {
100 type Item = common_recordbatch::error::Result<RecordBatch>;
101
102 fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
103 if let Some(batch) = self.pending.pop_front() {
104 return Poll::Ready(Some(Ok(batch)));
105 }
106
107 let batch = futures::ready!(self.inner.poll_next_unpin(cx));
108 let Some(batch) = batch else {
109 return Poll::Ready(None);
110 };
111
112 let record_batch = match batch {
113 Ok(batch) => {
114 let start = Instant::now();
115 let record_batch = self.convert(batch);
116 self.partition_metrics
117 .inc_convert_batch_cost(start.elapsed());
118 record_batch
119 }
120 Err(e) => Err(BoxedError::new(e)).context(ExternalSnafu),
121 };
122 Poll::Ready(Some(record_batch))
123 }
124}