Skip to main content

mito2/read/
stream.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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
34/// All kinds of [`Batch`]es to produce in scanner.
35pub enum ScanBatch {
36    Series(SeriesBatch),
37    RecordBatch(DfRecordBatch),
38}
39
40pub type ScanBatchStream = BoxStream<'static, Result<ScanBatch>>;
41
42/// A stream that takes [`ScanBatch`]es and produces (converts them to) [`RecordBatch`]es.
43pub(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                // Safety: Only flat format returns this batch.
77                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                // Safety: Only flat format returns this batch.
92                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}