Skip to main content

flow/
expr.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
15//! Small row-batch conversion helper used by stateless streaming.
16
17pub(crate) mod error;
18use common_recordbatch::RecordBatch;
19use datatypes::data_type::DataType;
20use datatypes::prelude::ConcreteDataType;
21use datatypes::vectors::{Helper, VectorRef};
22use itertools::Itertools;
23use snafu::{ResultExt, ensure};
24
25use crate::Error;
26use crate::error::DatatypesSnafu;
27use crate::repr::Row;
28
29#[derive(Debug, Clone)]
30pub struct Batch {
31    batch: Vec<VectorRef>,
32    row_count: usize,
33}
34
35impl TryFrom<RecordBatch> for Batch {
36    type Error = Error;
37    fn try_from(value: RecordBatch) -> Result<Self, Self::Error> {
38        Ok(Self {
39            batch: Helper::try_into_vectors(value.columns()).context(DatatypesSnafu {
40                extra: "failed to convert Arrow array to vector",
41            })?,
42            row_count: value.num_rows(),
43        })
44    }
45}
46
47impl Batch {
48    pub fn try_from_rows_with_types(
49        rows: Vec<Row>,
50        types: &[ConcreteDataType],
51    ) -> Result<Self, error::EvalError> {
52        if rows.is_empty() {
53            return Ok(Self {
54                batch: vec![],
55                row_count: 0,
56            });
57        }
58        let len = rows.len();
59        let mut builders = types
60            .iter()
61            .map(|ty| ty.create_mutable_vector(len))
62            .collect_vec();
63        ensure!(
64            rows.iter().all(|row| row.len() == builders.len()),
65            error::InvalidArgumentSnafu {
66                reason: "row length does not match schema".to_string()
67            }
68        );
69        for row in rows {
70            for (idx, value) in row.iter().enumerate() {
71                builders[idx]
72                    .try_push_value_ref(&value.as_value_ref())
73                    .context(error::DataTypeSnafu {
74                        msg: "failed to convert rows to columns",
75                    })?;
76            }
77        }
78        Ok(Self {
79            batch: builders.into_iter().map(|mut b| b.to_vector()).collect(),
80            row_count: len,
81        })
82    }
83    pub fn batch(&self) -> &[VectorRef] {
84        &self.batch
85    }
86    pub fn row_count(&self) -> usize {
87        self.row_count
88    }
89}