datafusion_functions/
binaries.rs1use 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 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 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 let array_data = unsafe { array_builder.build_unchecked() };
130 let array = GenericBinaryArray::<O>::from(array_data);
131 Ok(Arc::new(array))
132 }
133}
134
135pub(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 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 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 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}