use arrow::array::{ArrayRef, OffsetSizeTrait};
use arrow_array::builder::GenericStringBuilder;
use arrow_array::Array;
use arrow_schema::DataType;
use datafusion::physical_plan::ColumnarValue;
use datafusion_common::{cast::as_generic_string_array, DataFusionError, ScalarValue};
use std::fmt::Write;
use std::sync::Arc;
pub fn spark_read_side_padding(args: &[ColumnarValue]) -> Result<ColumnarValue, DataFusionError> {
match args {
[ColumnarValue::Array(array), ColumnarValue::Scalar(ScalarValue::Int32(Some(length)))] => {
match array.data_type() {
DataType::Utf8 => spark_read_side_padding_internal::<i32>(array, *length),
DataType::LargeUtf8 => spark_read_side_padding_internal::<i64>(array, *length),
other => Err(DataFusionError::Internal(format!(
"Unsupported data type {other:?} for function read_side_padding",
))),
}
}
other => Err(DataFusionError::Internal(format!(
"Unsupported arguments {other:?} for function read_side_padding",
))),
}
}
fn spark_read_side_padding_internal<T: OffsetSizeTrait>(
array: &ArrayRef,
length: i32,
) -> Result<ColumnarValue, DataFusionError> {
let string_array = as_generic_string_array::<T>(array)?;
let length = 0.max(length) as usize;
let space_string = " ".repeat(length);
let mut builder =
GenericStringBuilder::<T>::with_capacity(string_array.len(), string_array.len() * length);
for string in string_array.iter() {
match string {
Some(string) => {
let char_len = string.chars().count();
if length <= char_len {
builder.append_value(string);
} else {
builder.write_str(string)?;
builder.append_value(&space_string[char_len..]);
}
}
_ => builder.append_null(),
}
}
Ok(ColumnarValue::Array(Arc::new(builder.finish())))
}