Skip to main content

datafusion_functions/
binaries.rs

1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements.  See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership.  The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License.  You may obtain a copy of the License at
8//
9//   http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied.  See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18use crate::strings::{ColumnarValueRef, ConcatBuilder};
19use arrow::array::{
20    Array, ArrayDataBuilder, ArrayRef, BinaryViewArray, GenericBinaryArray,
21    OffsetSizeTrait, make_view,
22};
23use arrow_buffer::{ArrowNativeType, Buffer, MutableBuffer, NullBuffer, ScalarBuffer};
24use datafusion_common::{Result, exec_datafusion_err, exec_err, internal_err};
25use std::marker::PhantomData;
26use std::sync::Arc;
27
28pub(crate) struct ConcatGenericBinaryBuilder<O: OffsetSizeTrait + ArrowNativeType> {
29    offsets_buffer: MutableBuffer,
30    value_buffer: MutableBuffer,
31    _phantom: PhantomData<O>,
32}
33pub(crate) type ConcatBinaryBuilder = ConcatGenericBinaryBuilder<i32>;
34pub(crate) type ConcatLargeBinaryBuilder = ConcatGenericBinaryBuilder<i64>;
35
36impl<O: OffsetSizeTrait + ArrowNativeType> ConcatGenericBinaryBuilder<O> {
37    pub fn with_capacity(item_capacity: usize, data_capacity: usize) -> Self {
38        let capacity = item_capacity
39            .checked_add(1)
40            .map(|i| i.saturating_mul(size_of::<O>()))
41            .expect("capacity integer overflow");
42
43        let mut offsets_buffer = MutableBuffer::with_capacity(capacity);
44        // SAFETY: the first offset value is definitely not going to exceed the bounds.
45        unsafe { offsets_buffer.push_unchecked(O::usize_as(0)) };
46        Self {
47            offsets_buffer,
48            value_buffer: MutableBuffer::with_capacity(data_capacity),
49            _phantom: PhantomData,
50        }
51    }
52}
53
54impl<O: OffsetSizeTrait + ArrowNativeType> ConcatBuilder
55    for ConcatGenericBinaryBuilder<O>
56{
57    fn write<const CHECK_VALID: bool>(
58        &mut self,
59        column: &ColumnarValueRef,
60        i: usize,
61    ) -> Result<()> {
62        match column {
63            ColumnarValueRef::Scalar(s) => {
64                self.value_buffer.extend_from_slice(s);
65            }
66            ColumnarValueRef::NullableBinaryArray(array) => {
67                if !CHECK_VALID || array.is_valid(i) {
68                    self.value_buffer.extend_from_slice(array.value(i));
69                }
70            }
71            ColumnarValueRef::NullableLargeBinaryArray(array) => {
72                if !CHECK_VALID || array.is_valid(i) {
73                    self.value_buffer.extend_from_slice(array.value(i));
74                }
75            }
76            ColumnarValueRef::NullableBinaryViewArray(array) => {
77                if !CHECK_VALID || array.is_valid(i) {
78                    self.value_buffer.extend_from_slice(array.value(i));
79                }
80            }
81            ColumnarValueRef::NonNullableBinaryArray(array) => {
82                self.value_buffer.extend_from_slice(array.value(i));
83            }
84            ColumnarValueRef::NonNullableLargeBinaryArray(array) => {
85                self.value_buffer.extend_from_slice(array.value(i));
86            }
87            ColumnarValueRef::NonNullableBinaryViewArray(array) => {
88                self.value_buffer.extend_from_slice(array.value(i));
89            }
90            _ => {
91                return exec_err!(
92                    "concat: unexpected column type for binary builder: {column:?}"
93                );
94            }
95        }
96        Ok(())
97    }
98
99    fn append_offset(&mut self) -> Result<()> {
100        let next_offset: O = O::from_usize(self.value_buffer.len())
101            .ok_or_else(|| exec_datafusion_err!("byte array offset overflow"))?;
102        self.offsets_buffer.push(next_offset);
103        Ok(())
104    }
105
106    /// Finalize the builder into a concrete [`GenericBinaryArray<O>`].
107    ///
108    /// # Errors
109    ///
110    /// Returns an error when:
111    ///
112    /// - the provided `null_buffer` is not the same length as the `offsets_buffer`.
113    fn finish(self, null_buffer: Option<NullBuffer>) -> Result<ArrayRef> {
114        let row_count = self.offsets_buffer.len() / size_of::<O>() - 1;
115        if let Some(ref null_buffer) = null_buffer
116            && null_buffer.len() != row_count
117        {
118            return internal_err!(
119                "Null buffer and offsets buffer must be the same length"
120            );
121        }
122        let array_builder = ArrayDataBuilder::new(GenericBinaryArray::<O>::DATA_TYPE)
123            .len(row_count)
124            .add_buffer(self.offsets_buffer.into())
125            .add_buffer(self.value_buffer.into())
126            .nulls(null_buffer);
127        // SAFETY: all data that was appended was valid and the values
128        // and offsets were created correctly
129        let array_data = unsafe { array_builder.build_unchecked() };
130        let array = GenericBinaryArray::<O>::from(array_data);
131        Ok(Arc::new(array))
132    }
133}
134
135/// Builder used by `concat`/`concat_ws` to assemble a [`BinaryViewArray`] one
136/// row at a time from multiple input columns.
137///
138/// Each row is written via repeated `write` calls (one per input
139/// fragment) followed by a single `append_offset` to commit the row
140/// as a single binary view. The output null buffer is supplied by the caller
141/// at `finish` time, avoiding per-row NULL handling work.
142///
143pub(crate) struct ConcatBinaryViewBuilder {
144    views: Vec<u128>,
145    data: Vec<u8>,
146    block: Vec<u8>,
147}
148
149impl ConcatBinaryViewBuilder {
150    pub fn with_capacity(item_capacity: usize, data_capacity: usize) -> Self {
151        Self {
152            views: Vec::with_capacity(item_capacity),
153            data: Vec::with_capacity(data_capacity),
154            block: vec![],
155        }
156    }
157}
158
159impl ConcatBuilder for ConcatBinaryViewBuilder {
160    fn write<const CHECK_VALID: bool>(
161        &mut self,
162        column: &ColumnarValueRef,
163        i: usize,
164    ) -> Result<()> {
165        match column {
166            ColumnarValueRef::Scalar(s) => {
167                self.block.extend_from_slice(s);
168            }
169            ColumnarValueRef::NullableBinaryArray(array) => {
170                if !CHECK_VALID || array.is_valid(i) {
171                    self.block.extend_from_slice(array.value(i));
172                }
173            }
174            ColumnarValueRef::NullableLargeBinaryArray(array) => {
175                if !CHECK_VALID || array.is_valid(i) {
176                    self.block.extend_from_slice(array.value(i));
177                }
178            }
179            ColumnarValueRef::NullableBinaryViewArray(array) => {
180                if !CHECK_VALID || array.is_valid(i) {
181                    self.block.extend_from_slice(array.value(i));
182                }
183            }
184            ColumnarValueRef::NonNullableBinaryArray(array) => {
185                self.block.extend_from_slice(array.value(i));
186            }
187            ColumnarValueRef::NonNullableLargeBinaryArray(array) => {
188                self.block.extend_from_slice(array.value(i));
189            }
190            ColumnarValueRef::NonNullableBinaryViewArray(array) => {
191                self.block.extend_from_slice(array.value(i));
192            }
193            _ => {
194                return exec_err!(
195                    "concat: unexpected column type for binary view builder: {column:?}"
196                );
197            }
198        }
199        Ok(())
200    }
201
202    /// Finalizes the current row by converting the accumulated data into a
203    /// StringView and appending it to the views buffer.
204    fn append_offset(&mut self) -> Result<()> {
205        let v = &self.block;
206        if v.len() > 12 {
207            let offset: u32 = self
208                .data
209                .len()
210                .try_into()
211                .map_err(|_| exec_datafusion_err!("byte array offset overflow"))?;
212            self.data.extend_from_slice(v);
213            self.views.push(make_view(v, 0, offset));
214        } else {
215            self.views.push(make_view(v, 0, 0));
216        }
217
218        self.block.clear();
219        Ok(())
220    }
221
222    /// Finalize the builder into a concrete [`BinaryViewArray`].
223    ///
224    /// # Errors
225    ///
226    /// Returns an error when:
227    ///
228    /// - the provided `null_buffer` length does not match the row count.
229    fn finish(self, null_buffer: Option<NullBuffer>) -> Result<ArrayRef> {
230        if let Some(ref nulls) = null_buffer
231            && nulls.len() != self.views.len()
232        {
233            return internal_err!(
234                "Null buffer length ({}) must match row count ({})",
235                nulls.len(),
236                self.views.len()
237            );
238        }
239
240        let buffers: Vec<Buffer> = if self.data.is_empty() {
241            vec![]
242        } else {
243            vec![Buffer::from(self.data)]
244        };
245
246        // SAFETY: views were constructed with correct lengths, offsets, and
247        // prefixes.
248        let array = unsafe {
249            BinaryViewArray::new_unchecked(
250                ScalarBuffer::from(self.views),
251                buffers,
252                null_buffer,
253            )
254        };
255        Ok(Arc::new(array))
256    }
257}