From 6aee409db881643d9ebd166eb4796ac4c188c3d6 Mon Sep 17 00:00:00 2001 From: Anoop Johnson Date: Wed, 26 Aug 2026 09:15:31 -0700 Subject: [PATCH] perf(arrow): build repeated constant string/binary columns without per-row clones create_primitive_array_repeated built a throwaway `vec![value.clone(); num_rows]` for the Utf8/Binary/LargeBinary/FixedSizeBinary arms before handing it to the array constructor. For strings that clones the value into num_rows separate heap allocations per batch; the binary arms allocate a throwaway intermediate Vec. Stream the single value straight into the Arrow buffer via `from_iter_values` with `std::iter::repeat_n` instead. Benchmarks (release, throwaway harness, old vs new toggled on this file only): - leaf create_primitive_array_repeated (Utf8, 8192-row batch): 23.0 -> 4.1 ns/row (5.6x) - full scan select(x, _partition) over a real manifest + 262k-row parquet file, string-partitioned: 20.9 -> 5.7 ns/row (3.7x) Note that the gain is allocation-counts and scales with how much of the projection is the string partition column. Fixes #3079 --- crates/iceberg/src/arrow/value.rs | 98 +++++++++++++++++++++++++------ 1 file changed, 80 insertions(+), 18 deletions(-) diff --git a/crates/iceberg/src/arrow/value.rs b/crates/iceberg/src/arrow/value.rs index d22e565a4a..a75dc78d20 100644 --- a/crates/iceberg/src/arrow/value.rs +++ b/crates/iceberg/src/arrow/value.rs @@ -864,24 +864,24 @@ pub(crate) fn create_primitive_array_repeated( (DataType::Float64, Some(PrimitiveLiteral::Double(value))) => { Arc::new(Float64Array::from(vec![value.0; num_rows])) } - (DataType::Utf8, Some(PrimitiveLiteral::String(value))) => { - Arc::new(StringArray::from(vec![value.clone(); num_rows])) - } - (DataType::Binary, Some(PrimitiveLiteral::Binary(value))) => { - Arc::new(BinaryArray::from_vec(vec![value; num_rows])) - } - (DataType::LargeBinary, Some(PrimitiveLiteral::Binary(value))) => { - Arc::new(LargeBinaryArray::from_vec(vec![value; num_rows])) - } - (DataType::FixedSizeBinary(len), Some(PrimitiveLiteral::Binary(value))) => { - let repeated: Vec<&[u8]> = vec![value.as_slice(); num_rows]; - Arc::new(FixedSizeBinaryArray::try_from_iter(repeated.into_iter()).map_err(|e| { - Error::new( - ErrorKind::DataInvalid, - format!("Failed to create FixedSizeBinary({len}) array: {e}"), - ) - })?) - } + (DataType::Utf8, Some(PrimitiveLiteral::String(value))) => Arc::new( + StringArray::from_iter_values(std::iter::repeat_n(value.as_str(), num_rows)), + ), + (DataType::Binary, Some(PrimitiveLiteral::Binary(value))) => Arc::new( + BinaryArray::from_iter_values(std::iter::repeat_n(value.as_slice(), num_rows)), + ), + (DataType::LargeBinary, Some(PrimitiveLiteral::Binary(value))) => Arc::new( + LargeBinaryArray::from_iter_values(std::iter::repeat_n(value.as_slice(), num_rows)), + ), + (DataType::FixedSizeBinary(len), Some(PrimitiveLiteral::Binary(value))) => Arc::new( + FixedSizeBinaryArray::try_from_iter(std::iter::repeat_n(value.as_slice(), num_rows)) + .map_err(|e| { + Error::new( + ErrorKind::DataInvalid, + format!("Failed to create FixedSizeBinary({len}) array: {e}"), + ) + })?, + ), (DataType::Time64(TimeUnit::Microsecond), Some(PrimitiveLiteral::Long(value))) => { Arc::new(Time64MicrosecondArray::from(vec![*value; num_rows])) } @@ -1914,4 +1914,66 @@ mod test { assert_eq!(array.data_type(), &target_type); assert_eq!(array.len(), num_rows); } + + #[test] + fn test_create_string_and_binary_arrays_repeated() { + let text = "partition-value-2026"; + let bytes: Vec = vec![0xDE, 0xAD, 0xBE, 0xEF]; + let num_rows = 4; + + let utf8 = create_primitive_array_repeated( + &DataType::Utf8, + &Some(PrimitiveLiteral::String(text.to_string())), + num_rows, + ) + .unwrap(); + let utf8 = utf8.as_any().downcast_ref::().unwrap(); + assert_eq!(utf8.len(), num_rows); + assert!((0..num_rows).all(|i| utf8.value(i) == text)); + + let binary = create_primitive_array_repeated( + &DataType::Binary, + &Some(PrimitiveLiteral::Binary(bytes.clone())), + num_rows, + ) + .unwrap(); + let binary = binary.as_any().downcast_ref::().unwrap(); + assert_eq!(binary.len(), num_rows); + assert!((0..num_rows).all(|i| binary.value(i) == bytes.as_slice())); + + let large = create_primitive_array_repeated( + &DataType::LargeBinary, + &Some(PrimitiveLiteral::Binary(bytes.clone())), + num_rows, + ) + .unwrap(); + let large = large.as_any().downcast_ref::().unwrap(); + assert_eq!(large.len(), num_rows); + assert!((0..num_rows).all(|i| large.value(i) == bytes.as_slice())); + + let fixed = create_primitive_array_repeated( + &DataType::FixedSizeBinary(bytes.len() as i32), + &Some(PrimitiveLiteral::Binary(bytes.clone())), + num_rows, + ) + .unwrap(); + let fixed = fixed + .as_any() + .downcast_ref::() + .unwrap(); + assert_eq!(fixed.len(), num_rows); + assert!((0..num_rows).all(|i| fixed.value(i) == bytes.as_slice())); + } + + #[test] + fn test_create_string_array_repeated_empty() { + // num_rows == 0 must produce an empty (not one-element) array. + let array = create_primitive_array_repeated( + &DataType::Utf8, + &Some(PrimitiveLiteral::String("x".to_string())), + 0, + ) + .unwrap(); + assert_eq!(array.len(), 0); + } }