This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/paimon-mosaic.git
The following commit(s) were added to refs/heads/main by this push:
new a8817f3 perf(core): speed up reading nullable fixed-width CONST
columns (#75)
a8817f3 is described below
commit a8817f311de51ac8e2f0b88d68395cabccf878aa
Author: jianguotian <[email protected]>
AuthorDate: Wed Aug 19 19:35:06 2026 +0800
perf(core): speed up reading nullable fixed-width CONST columns (#75)
---
core/Cargo.toml | 3 +
core/src/bucket_reader.rs | 876 +++++++++++++++++++++++++++++++++++++++++++---
core/src/reader_tests.rs | 253 +++++++++++++
3 files changed, 1084 insertions(+), 48 deletions(-)
diff --git a/core/Cargo.toml b/core/Cargo.toml
index 8728a66..a2942b7 100644
--- a/core/Cargo.toml
+++ b/core/Cargo.toml
@@ -30,3 +30,6 @@ zstd = "0.13"
arrow-schema = "58"
arrow-array = "58"
arrow-buffer = "58"
+
+[target.'cfg(unix)'.dependencies]
+libc = "0.2"
diff --git a/core/src/bucket_reader.rs b/core/src/bucket_reader.rs
index 5b3cc39..4dbaab4 100644
--- a/core/src/bucket_reader.rs
+++ b/core/src/bucket_reader.rs
@@ -16,6 +16,7 @@
// under the License.
use std::io;
+use std::mem::size_of;
use std::sync::Arc;
use arrow_array::*;
@@ -62,6 +63,59 @@ fn data_variant_for_type(dt: &DataType) -> DataVariant {
}
}
+fn const_values_are_row_aligned(variant: DataVariant) -> bool {
+ match variant {
+ DataVariant::Boolean
+ | DataVariant::Int8
+ | DataVariant::Int16
+ | DataVariant::Int32
+ | DataVariant::Int64
+ | DataVariant::Float32
+ | DataVariant::Float64
+ | DataVariant::TimestampNanos => true,
+ DataVariant::Binary => false,
+ }
+}
+
+const CONST_FILL_ALL_MIN_NON_NULL_DENOMINATOR: usize = 8;
+
+#[derive(Debug, Clone, Copy, PartialEq, Eq)]
+enum ConstFillStrategy {
+ All,
+ NonNullOnly,
+ BulkFillIfAllPagesTouched,
+}
+
+#[derive(Debug, Clone, Copy, PartialEq, Eq)]
+enum ConstOutputLayout {
+ FixedWidth { bytes_per_row: usize },
+ BitPackedBoolean,
+}
+
+impl ConstOutputLayout {
+ fn row_range_for_bytes(
+ self,
+ start_byte: usize,
+ end_byte: usize,
+ num_rows: usize,
+ ) -> (usize, usize) {
+ debug_assert!(start_byte < end_byte);
+ match self {
+ Self::FixedWidth { bytes_per_row } => {
+ debug_assert!(bytes_per_row > 0);
+ (
+ start_byte / bytes_per_row,
+ end_byte.div_ceil(bytes_per_row).min(num_rows),
+ )
+ }
+ Self::BitPackedBoolean => (
+ start_byte.saturating_mul(8).min(num_rows),
+ end_byte.saturating_mul(8).min(num_rows),
+ ),
+ }
+ }
+}
+
#[derive(Debug, Clone)]
enum RawColumnData {
Boolean(Vec<u8>),
@@ -118,6 +172,280 @@ fn invert_bitmap(bitmap: &[u8]) -> Vec<u8> {
bitmap.iter().map(|b| !b).collect()
}
+fn for_each_non_null(null_bitmap: &[u8], num_rows: usize, mut f: impl
FnMut(usize)) {
+ for (byte_index, &nulls) in null_bitmap.iter().enumerate() {
+ let row_base = byte_index * 8;
+ if row_base >= num_rows {
+ break;
+ }
+
+ let rows_in_byte = (num_rows - row_base).min(8);
+ let row_mask = if rows_in_byte == 8 {
+ u8::MAX
+ } else {
+ (1u8 << rows_in_byte) - 1
+ };
+ let mut non_nulls = !nulls & row_mask;
+ while non_nulls != 0 {
+ let bit = non_nulls.trailing_zeros() as usize;
+ f(row_base + bit);
+ non_nulls &= non_nulls - 1;
+ }
+ }
+}
+
+fn row_range_has_non_null(null_bitmap: &[u8], start_row: usize, end_row:
usize) -> bool {
+ if start_row >= end_row {
+ return false;
+ }
+
+ let first_byte = start_row / 8;
+ let last_byte = (end_row - 1) / 8;
+ let first_bit = start_row % 8;
+ let end_bit = end_row % 8;
+
+ if first_byte == last_byte {
+ let end_mask = if end_bit == 0 {
+ u8::MAX
+ } else {
+ (1u8 << end_bit) - 1
+ };
+ let row_mask = end_mask & (u8::MAX << first_bit);
+ return (!null_bitmap[first_byte] & row_mask) != 0;
+ }
+
+ let mut full_start = first_byte;
+ if first_bit != 0 {
+ if (!null_bitmap[first_byte] & (u8::MAX << first_bit)) != 0 {
+ return true;
+ }
+ full_start += 1;
+ }
+
+ let full_end = if end_bit == 0 {
+ last_byte + 1
+ } else {
+ last_byte
+ };
+ if null_bitmap[full_start..full_end]
+ .iter()
+ .any(|&nulls| nulls != u8::MAX)
+ {
+ return true;
+ }
+
+ end_bit != 0 && (!null_bitmap[last_byte] & ((1u8 << end_bit) - 1)) != 0
+}
+
+fn for_each_non_null_row_run(null_bitmap: &[u8], num_rows: usize, mut f: impl
FnMut(usize, usize)) {
+ let mut run_start = None;
+ let mut run_end = 0;
+ for_each_non_null(null_bitmap, num_rows, |row| match run_start {
+ None => {
+ run_start = Some(row);
+ run_end = row + 1;
+ }
+ Some(_) if row == run_end => run_end += 1,
+ Some(_) => {
+ if let Some(start) = run_start.replace(row) {
+ f(start, run_end);
+ }
+ run_end = row + 1;
+ }
+ });
+ if let Some(start) = run_start {
+ f(start, run_end);
+ }
+}
+
+fn system_page_size() -> Option<usize> {
+ #[cfg(unix)]
+ {
+ static PAGE_SIZE: std::sync::OnceLock<Option<usize>> =
std::sync::OnceLock::new();
+ *PAGE_SIZE.get_or_init(|| {
+ // SAFETY: sysconf with _SC_PAGESIZE has no pointer arguments or
caller-owned state.
+ let page_size = unsafe { libc::sysconf(libc::_SC_PAGESIZE) };
+ (page_size > 0).then_some(page_size as usize)
+ })
+ }
+
+ #[cfg(not(unix))]
+ {
+ None
+ }
+}
+
+fn all_output_pages_touched(
+ null_bitmap: &[u8],
+ num_rows: usize,
+ output_addr: usize,
+ output_len_bytes: usize,
+ page_size: usize,
+ layout: ConstOutputLayout,
+) -> bool {
+ if output_len_bytes == 0 || page_size == 0 {
+ return false;
+ }
+ let Some(output_end) = output_addr.checked_add(output_len_bytes) else {
+ return false;
+ };
+
+ let mut page_start = output_addr - output_addr % page_size;
+ loop {
+ let Some(next_page) = page_start.checked_add(page_size) else {
+ return false;
+ };
+ let start_byte = page_start.max(output_addr) - output_addr;
+ let end_byte = next_page.min(output_end) - output_addr;
+ let (start_row, end_row) = layout.row_range_for_bytes(start_byte,
end_byte, num_rows);
+ if !row_range_has_non_null(null_bitmap, start_row, end_row) {
+ return false;
+ }
+
+ if next_page >= output_end {
+ return true;
+ }
+ page_start = next_page;
+ }
+}
+
+fn const_fill_strategy(
+ has_nulls: bool,
+ non_null_count: usize,
+ num_rows: usize,
+) -> ConstFillStrategy {
+ if !has_nulls || non_null_count == num_rows {
+ ConstFillStrategy::All
+ } else if non_null_count <
num_rows.div_ceil(CONST_FILL_ALL_MIN_NON_NULL_DENOMINATOR) {
+ ConstFillStrategy::NonNullOnly
+ } else {
+ ConstFillStrategy::BulkFillIfAllPagesTouched
+ }
+}
+
+trait ConstMaterializeValue: Default + Copy {
+ fn has_default_bit_pattern(self) -> bool;
+}
+
+macro_rules! impl_const_materialize_integer {
+ ($($ty:ty),+ $(,)?) => {
+ $(
+ impl ConstMaterializeValue for $ty {
+ fn has_default_bit_pattern(self) -> bool {
+ self == 0
+ }
+ }
+ )+
+ };
+}
+
+impl_const_materialize_integer!(i8, i16, i32, i64);
+
+impl ConstMaterializeValue for f32 {
+ fn has_default_bit_pattern(self) -> bool {
+ self.to_bits() == 0
+ }
+}
+
+impl ConstMaterializeValue for f64 {
+ fn has_default_bit_pattern(self) -> bool {
+ self.to_bits() == 0
+ }
+}
+
+fn materialize_fixed_const<T: ConstMaterializeValue>(
+ value: T,
+ num_rows: usize,
+ strategy: ConstFillStrategy,
+ null_bitmap: &[u8],
+) -> Vec<T> {
+ if strategy == ConstFillStrategy::All {
+ return vec![value; num_rows];
+ }
+
+ let mut out = vec![T::default(); num_rows];
+ if value.has_default_bit_pattern() {
+ return out;
+ }
+
+ if strategy == ConstFillStrategy::BulkFillIfAllPagesTouched {
+ if let Some(page_size) = system_page_size() {
+ let output_len_bytes = out.len() * size_of::<T>();
+ if all_output_pages_touched(
+ null_bitmap,
+ num_rows,
+ out.as_ptr() as usize,
+ output_len_bytes,
+ page_size,
+ ConstOutputLayout::FixedWidth {
+ bytes_per_row: size_of::<T>(),
+ },
+ ) {
+ out.fill(value);
+ return out;
+ }
+ }
+ }
+
+ match strategy {
+ ConstFillStrategy::All => unreachable!(),
+ ConstFillStrategy::NonNullOnly => {
+ for_each_non_null(null_bitmap, num_rows, |row| out[row] = value);
+ }
+ ConstFillStrategy::BulkFillIfAllPagesTouched => {
+ for_each_non_null_row_run(null_bitmap, num_rows, |start, end| {
+ out[start..end].fill(value);
+ });
+ }
+ }
+ out
+}
+
+fn materialize_boolean_const(
+ value: bool,
+ num_rows: usize,
+ strategy: ConstFillStrategy,
+ null_bitmap: &[u8],
+) -> Vec<u8> {
+ let mut out = vec![0u8; num_rows.div_ceil(8)];
+ if !value {
+ return out;
+ }
+
+ match strategy {
+ ConstFillStrategy::All => out.fill(u8::MAX),
+ ConstFillStrategy::NonNullOnly => {
+ for_each_non_null(null_bitmap, num_rows, |row| {
+ out[row / 8] |= 1 << (row % 8);
+ });
+ }
+ ConstFillStrategy::BulkFillIfAllPagesTouched => {
+ if system_page_size().is_some_and(|page_size| {
+ all_output_pages_touched(
+ null_bitmap,
+ num_rows,
+ out.as_ptr() as usize,
+ out.len(),
+ page_size,
+ ConstOutputLayout::BitPackedBoolean,
+ )
+ }) {
+ out.fill(u8::MAX);
+ } else {
+ for_each_non_null(null_bitmap, num_rows, |row| {
+ out[row / 8] |= 1 << (row % 8);
+ });
+ }
+ }
+ }
+
+ if num_rows & 7 != 0 {
+ let last = out.len() - 1;
+ out[last] &= (1u8 << (num_rows % 8)) - 1;
+ }
+ out
+}
+
fn make_null_buffer(bitmap: Option<Vec<u8>>, num_rows: usize) ->
Option<NullBuffer> {
bitmap.map(|bm| NullBuffer::new(BooleanBuffer::new(Buffer::from_vec(bm),
0, num_rows)))
}
@@ -187,8 +515,34 @@ fn build_array(
dt: &DataType,
null_bitmap: Option<Vec<u8>>,
num_rows: usize,
+ values_are_row_aligned: bool,
) -> io::Result<ArrayRef> {
+ debug_assert!(
+ !values_are_row_aligned
+ || match &data {
+ RawColumnData::Boolean(values) => values.len() ==
num_rows.div_ceil(8),
+ RawColumnData::Int8(values) => values.len() == num_rows,
+ RawColumnData::Int16(values) => values.len() == num_rows,
+ RawColumnData::Int32(values) => values.len() == num_rows,
+ RawColumnData::Int64(values) => values.len() == num_rows,
+ RawColumnData::Float32(values) => values.len() == num_rows,
+ RawColumnData::Float64(values) => values.len() == num_rows,
+ RawColumnData::Binary { .. } => false,
+ RawColumnData::TimestampNanos {
+ millis,
+ nanos_of_milli,
+ } => millis.len() == num_rows && nanos_of_milli.len() ==
num_rows,
+ },
+ "row-aligned CONST data must match the row cardinality"
+ );
+
let null_buf = make_null_buffer(null_bitmap.clone(), num_rows);
+ let no_scatter = None;
+ let scatter_bitmap = if values_are_row_aligned {
+ &no_scatter
+ } else {
+ &null_bitmap
+ };
Ok(match data {
RawColumnData::Boolean(values) => {
@@ -196,15 +550,15 @@ fn build_array(
Arc::new(BooleanArray::new(bool_buf, null_buf))
}
RawColumnData::Int8(values) => {
- let scattered = scatter_fixed(values, &null_bitmap, num_rows);
+ let scattered = scatter_fixed(values, scatter_bitmap, num_rows);
Arc::new(Int8Array::new(ScalarBuffer::from(scattered), null_buf))
}
RawColumnData::Int16(values) => {
- let scattered = scatter_fixed(values, &null_bitmap, num_rows);
+ let scattered = scatter_fixed(values, scatter_bitmap, num_rows);
Arc::new(Int16Array::new(ScalarBuffer::from(scattered), null_buf))
}
RawColumnData::Int32(values) => {
- let scattered = scatter_fixed(values, &null_bitmap, num_rows);
+ let scattered = scatter_fixed(values, scatter_bitmap, num_rows);
match dt {
DataType::Date32 => {
Arc::new(Date32Array::new(ScalarBuffer::from(scattered),
null_buf))
@@ -217,7 +571,7 @@ fn build_array(
}
}
RawColumnData::Int64(values) => {
- let scattered = scatter_fixed(values, &null_bitmap, num_rows);
+ let scattered = scatter_fixed(values, scatter_bitmap, num_rows);
match dt {
DataType::Decimal128(p, s) => {
let i128_values: Vec<i128> = scattered.iter().map(|&v| v
as i128).collect();
@@ -249,16 +603,16 @@ fn build_array(
}
}
RawColumnData::Float32(values) => {
- let scattered = scatter_fixed(values, &null_bitmap, num_rows);
+ let scattered = scatter_fixed(values, scatter_bitmap, num_rows);
Arc::new(Float32Array::new(ScalarBuffer::from(scattered),
null_buf))
}
RawColumnData::Float64(values) => {
- let scattered = scatter_fixed(values, &null_bitmap, num_rows);
+ let scattered = scatter_fixed(values, scatter_bitmap, num_rows);
Arc::new(Float64Array::new(ScalarBuffer::from(scattered),
null_buf))
}
RawColumnData::Binary { offsets, data } => {
let (i32_offsets, out_data) =
- scatter_binary_offsets(offsets, data, &null_bitmap, num_rows);
+ scatter_binary_offsets(offsets, data, scatter_bitmap,
num_rows);
let offset_buf =
OffsetBuffer::new(ScalarBuffer::from(i32_offsets));
match dt {
DataType::Utf8 => Arc::new(StringArray::new(
@@ -301,16 +655,30 @@ fn build_array(
millis,
nanos_of_milli,
} => {
- let millis_scattered = scatter_fixed(millis, &null_bitmap,
num_rows);
- let nanos_scattered = scatter_fixed(nanos_of_milli, &null_bitmap,
num_rows);
+ let millis_scattered = scatter_fixed(millis, scatter_bitmap,
num_rows);
+ let nanos_scattered = scatter_fixed(nanos_of_milli,
scatter_bitmap, num_rows);
match dt {
DataType::Timestamp(TimeUnit::Nanosecond, tz) => {
- let values = millis_scattered
- .into_iter()
- .zip(nanos_scattered)
- .map(|(millis, nanos)|
types::millis_nanos_to_ns(millis, nanos))
- .collect::<io::Result<Vec<_>>>()?;
+ let values = match (values_are_row_aligned,
null_buf.as_ref()) {
+ (true, Some(nulls)) => millis_scattered
+ .into_iter()
+ .zip(nanos_scattered)
+ .enumerate()
+ .map(|(row, (millis, nanos))| {
+ if nulls.is_null(row) {
+ Ok(0)
+ } else {
+ types::millis_nanos_to_ns(millis, nanos)
+ }
+ })
+ .collect::<io::Result<Vec<_>>>()?,
+ _ => millis_scattered
+ .into_iter()
+ .zip(nanos_scattered)
+ .map(|(millis, nanos)|
types::millis_nanos_to_ns(millis, nanos))
+ .collect::<io::Result<Vec<_>>>()?,
+ };
let arr =
TimestampNanosecondArray::new(ScalarBuffer::from(values), null_buf);
Arc::new(if let Some(tz) = tz {
arr.with_timezone(tz.clone())
@@ -802,6 +1170,7 @@ impl BucketReader {
&self.col_types[i],
null_bitmap,
col_rows,
+ self.encodings[i] == ENCODING_CONST &&
const_values_are_row_aligned(variant),
)?);
}
@@ -1073,7 +1442,13 @@ impl ColumnPageReader {
_ => empty_raw_data_for_type(&self.col_type),
};
- build_array(data, &self.col_type, null_bitmap, num_rows)
+ build_array(
+ data,
+ &self.col_type,
+ null_bitmap,
+ num_rows,
+ self.encoding == ENCODING_CONST &&
const_values_are_row_aligned(variant),
+ )
}
}
@@ -1178,6 +1553,12 @@ fn read_all_const(
} else {
num_rows
};
+ // Fixed-width CONST arrays are always materialized at row cardinality, so
build_array can
+ // wrap them without another scatter buffer. Sparse columns write only
valid positions. At
+ // ordinary densities, bulk-fill the whole buffer only when every physical
page covered by the
+ // actual allocation already contains valid rows. Otherwise fill only
contiguous non-null row
+ // runs so null-only pages remain untouched.
+ let fill_strategy = const_fill_strategy(has_nulls, non_null_count,
num_rows);
match variant {
DataVariant::Boolean => {
@@ -1185,42 +1566,48 @@ fn read_all_const(
Value::Boolean(v) => *v,
_ => false,
};
- let mut buf = vec![0u8; num_rows.div_ceil(8)];
- if b {
- for row in 0..num_rows {
- if !has_nulls || !is_null(null_bitmap, row) {
- buf[row / 8] |= 1 << (row % 8);
- }
- }
- }
- Ok(RawColumnData::Boolean(buf))
+ Ok(RawColumnData::Boolean(materialize_boolean_const(
+ b,
+ num_rows,
+ fill_strategy,
+ null_bitmap,
+ )))
}
DataVariant::Int8 => {
let v = match const_value {
Value::TinyInt(x) => *x,
_ => 0,
};
- let mut out = vec![0i8; non_null_count];
- out.fill(v);
- Ok(RawColumnData::Int8(out))
+ Ok(RawColumnData::Int8(materialize_fixed_const(
+ v,
+ num_rows,
+ fill_strategy,
+ null_bitmap,
+ )))
}
DataVariant::Int16 => {
let v = match const_value {
Value::SmallInt(x) => *x,
_ => 0,
};
- let mut out = vec![0i16; non_null_count];
- out.fill(v);
- Ok(RawColumnData::Int16(out))
+ Ok(RawColumnData::Int16(materialize_fixed_const(
+ v,
+ num_rows,
+ fill_strategy,
+ null_bitmap,
+ )))
}
DataVariant::Int32 => {
let v = match const_value {
Value::Integer(x) | Value::Date(x) | Value::Time(x) => *x,
_ => 0,
};
- let mut out = vec![0i32; non_null_count];
- out.fill(v);
- Ok(RawColumnData::Int32(out))
+ Ok(RawColumnData::Int32(materialize_fixed_const(
+ v,
+ num_rows,
+ fill_strategy,
+ null_bitmap,
+ )))
}
DataVariant::Int64 => {
let v = match const_value {
@@ -1230,27 +1617,36 @@ fn read_all_const(
| Value::TimestampMicros(x) => *x,
_ => 0,
};
- let mut out = vec![0i64; non_null_count];
- out.fill(v);
- Ok(RawColumnData::Int64(out))
+ Ok(RawColumnData::Int64(materialize_fixed_const(
+ v,
+ num_rows,
+ fill_strategy,
+ null_bitmap,
+ )))
}
DataVariant::Float32 => {
let v = match const_value {
Value::Float(x) => *x,
_ => 0.0,
};
- let mut out = vec![0.0f32; non_null_count];
- out.fill(v);
- Ok(RawColumnData::Float32(out))
+ Ok(RawColumnData::Float32(materialize_fixed_const(
+ v,
+ num_rows,
+ fill_strategy,
+ null_bitmap,
+ )))
}
DataVariant::Float64 => {
let v = match const_value {
Value::Double(x) => *x,
_ => 0.0,
};
- let mut out = vec![0.0f64; non_null_count];
- out.fill(v);
- Ok(RawColumnData::Float64(out))
+ Ok(RawColumnData::Float64(materialize_fixed_const(
+ v,
+ num_rows,
+ fill_strategy,
+ null_bitmap,
+ )))
}
DataVariant::Binary => {
let bytes = match const_value {
@@ -1274,13 +1670,9 @@ fn read_all_const(
} => (*millis, *nanos_of_milli),
_ => (0, 0),
};
- let mut millis_out = vec![0i64; non_null_count];
- let mut nanos_out = vec![0i32; non_null_count];
- millis_out.fill(m);
- nanos_out.fill(n);
Ok(RawColumnData::TimestampNanos {
- millis: millis_out,
- nanos_of_milli: nanos_out,
+ millis: materialize_fixed_const(m, num_rows, fill_strategy,
null_bitmap),
+ nanos_of_milli: materialize_fixed_const(n, num_rows,
fill_strategy, null_bitmap),
})
}
}
@@ -1636,3 +2028,391 @@ fn read_u64(buf: &[u8], pos: usize) -> u64 {
buf[pos + 7],
])
}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+
+ fn null_bitmap(num_rows: usize, non_null_rows: &[usize]) -> Vec<u8> {
+ let mut bitmap = vec![u8::MAX; num_rows.div_ceil(8)];
+ for &row in non_null_rows {
+ bitmap[row / 8] &= !(1 << (row % 8));
+ }
+ bitmap
+ }
+
+ fn test_page_size() -> usize {
+ system_page_size().unwrap_or(4096)
+ }
+
+ #[test]
+ fn test_sparse_fixed_const_writes_only_non_null_slots() {
+ let num_rows = 32;
+ let bitmap = null_bitmap(num_rows, &[0, 31]);
+ let data = read_all_const(
+ &Value::BigInt(42),
+ num_rows,
+ true,
+ &bitmap,
+ DataVariant::Int64,
+ )
+ .unwrap();
+
+ let RawColumnData::Int64(values) = data else {
+ panic!("expected Int64 CONST data");
+ };
+ assert_eq!(values.len(), num_rows);
+ assert_eq!(values[0], 42);
+ assert_eq!(values[31], 42);
+ assert!(values[1..31].iter().all(|&value| value == 0));
+ }
+
+ #[test]
+ fn test_dense_fixed_const_values_are_row_aligned() {
+ let num_rows = 32;
+ let non_null_rows = (0..24).collect::<Vec<_>>();
+ let bitmap = null_bitmap(num_rows, &non_null_rows);
+ let data = read_all_const(
+ &Value::BigInt(42),
+ num_rows,
+ true,
+ &bitmap,
+ DataVariant::Int64,
+ )
+ .unwrap();
+
+ let RawColumnData::Int64(values) = data else {
+ panic!("expected Int64 CONST data");
+ };
+ assert_eq!(values.len(), num_rows);
+ for row in non_null_rows {
+ assert_eq!(values[row], 42);
+ }
+ }
+
+ #[test]
+ fn test_interleaved_const_chunks_write_only_non_null_runs() {
+ let rows_per_chunk = test_page_size() / size_of::<i64>();
+ let num_rows = rows_per_chunk * 4;
+ let non_null_rows = [0, 2]
+ .into_iter()
+ .flat_map(|chunk| {
+ let start = chunk * rows_per_chunk;
+ start..start + rows_per_chunk / 4
+ })
+ .collect::<Vec<_>>();
+ assert_eq!(
+ non_null_rows.len(),
+ num_rows.div_ceil(CONST_FILL_ALL_MIN_NON_NULL_DENOMINATOR)
+ );
+ let bitmap = null_bitmap(num_rows, &non_null_rows);
+ let data = read_all_const(
+ &Value::BigInt(42),
+ num_rows,
+ true,
+ &bitmap,
+ DataVariant::Int64,
+ )
+ .unwrap();
+
+ let RawColumnData::Int64(values) = data else {
+ panic!("expected Int64 CONST data");
+ };
+ for (row, value) in values.into_iter().enumerate() {
+ let chunk = row / rows_per_chunk;
+ let row_in_chunk = row % rows_per_chunk;
+ let expected = chunk & 1 == 0 && row_in_chunk < rows_per_chunk / 4;
+ assert_eq!(value, if expected { 42 } else { 0 }, "row {row}");
+ }
+ }
+
+ #[test]
+ fn test_clustered_const_at_density_cutoff_fills_only_non_null_runs() {
+ let rows_per_chunk = test_page_size() / size_of::<i64>();
+ let num_rows = rows_per_chunk *
CONST_FILL_ALL_MIN_NON_NULL_DENOMINATOR;
+ let non_null_rows = (0..rows_per_chunk).collect::<Vec<_>>();
+ let bitmap = null_bitmap(num_rows, &non_null_rows);
+ let data = read_all_const(
+ &Value::BigInt(42),
+ num_rows,
+ true,
+ &bitmap,
+ DataVariant::Int64,
+ )
+ .unwrap();
+
+ let RawColumnData::Int64(values) = data else {
+ panic!("expected Int64 CONST data");
+ };
+ assert!(values[..rows_per_chunk].iter().all(|&value| value == 42));
+ assert!(values[rows_per_chunk..].iter().all(|&value| value == 0));
+ }
+
+ #[test]
+ fn test_distributed_const_at_density_cutoff_preserves_valid_values() {
+ let rows_per_chunk = test_page_size() / size_of::<i64>();
+ let num_rows = rows_per_chunk *
CONST_FILL_ALL_MIN_NON_NULL_DENOMINATOR;
+ let non_null_rows = (0..num_rows)
+ .step_by(CONST_FILL_ALL_MIN_NON_NULL_DENOMINATOR)
+ .collect::<Vec<_>>();
+ let bitmap = null_bitmap(num_rows, &non_null_rows);
+ let data = read_all_const(
+ &Value::BigInt(42),
+ num_rows,
+ true,
+ &bitmap,
+ DataVariant::Int64,
+ )
+ .unwrap();
+
+ let RawColumnData::Int64(values) = data else {
+ panic!("expected Int64 CONST data");
+ };
+ assert_eq!(values.len(), num_rows);
+ for row in non_null_rows {
+ assert_eq!(values[row], 42);
+ }
+ }
+
+ #[test]
+ fn
test_output_page_coverage_distinguishes_distributed_and_clustered_values() {
+ let page_size = 4096;
+ let rows_per_page = page_size / size_of::<i64>();
+ let num_rows = rows_per_page * CONST_FILL_ALL_MIN_NON_NULL_DENOMINATOR;
+ let output_addr = 0;
+ let output_len_bytes = num_rows * size_of::<i64>();
+ let layout = ConstOutputLayout::FixedWidth {
+ bytes_per_row: size_of::<i64>(),
+ };
+
+ let distributed_rows = (0..num_rows)
+ .step_by(CONST_FILL_ALL_MIN_NON_NULL_DENOMINATOR)
+ .collect::<Vec<_>>();
+ let distributed_bitmap = null_bitmap(num_rows, &distributed_rows);
+ assert!(all_output_pages_touched(
+ &distributed_bitmap,
+ num_rows,
+ output_addr,
+ output_len_bytes,
+ page_size,
+ layout,
+ ));
+
+ let clustered_rows = (0..rows_per_page).collect::<Vec<_>>();
+ let clustered_bitmap = null_bitmap(num_rows, &clustered_rows);
+ assert!(!all_output_pages_touched(
+ &clustered_bitmap,
+ num_rows,
+ output_addr,
+ output_len_bytes,
+ page_size,
+ layout,
+ ));
+ }
+
+ #[test]
+ fn test_page_coverage_accounts_for_allocation_offset() {
+ let page_size = 4096;
+ let rows_per_chunk = page_size / size_of::<i64>();
+ let num_rows = rows_per_chunk *
CONST_FILL_ALL_MIN_NON_NULL_DENOMINATOR;
+ let non_null_rows = (0..CONST_FILL_ALL_MIN_NON_NULL_DENOMINATOR / 2)
+ .flat_map(|pair| {
+ let even_chunk_start = pair * 2 * rows_per_chunk;
+ let odd_chunk_end = even_chunk_start + 2 * rows_per_chunk;
+ (even_chunk_start..even_chunk_start +
127).chain(std::iter::once(odd_chunk_end - 1))
+ })
+ .collect::<Vec<_>>();
+ assert_eq!(
+ non_null_rows.len(),
+ num_rows.div_ceil(CONST_FILL_ALL_MIN_NON_NULL_DENOMINATOR)
+ );
+ let bitmap = null_bitmap(num_rows, &non_null_rows);
+ for chunk_start in (0..num_rows).step_by(rows_per_chunk) {
+ assert!(row_range_has_non_null(
+ &bitmap,
+ chunk_start,
+ (chunk_start + rows_per_chunk).min(num_rows)
+ ));
+ }
+
+ let layout = ConstOutputLayout::FixedWidth {
+ bytes_per_row: size_of::<i64>(),
+ };
+ let output_len_bytes = num_rows * size_of::<i64>();
+ assert!(all_output_pages_touched(
+ &bitmap,
+ num_rows,
+ 0,
+ output_len_bytes,
+ page_size,
+ layout,
+ ));
+ assert!(!all_output_pages_touched(
+ &bitmap,
+ num_rows,
+ 16,
+ output_len_bytes,
+ page_size,
+ layout,
+ ));
+ }
+
+ #[test]
+ fn test_page_coverage_ignores_trailing_bitmap_bits() {
+ let page_size = 4096;
+ let rows_per_page = page_size / size_of::<i64>();
+ let num_rows = rows_per_page + 1;
+ let non_null_rows = (0..num_rows.div_ceil(8)).collect::<Vec<_>>();
+ let mut bitmap = null_bitmap(num_rows, &non_null_rows);
+ bitmap[num_rows / 8] = 0b0000_0001;
+
+ assert!(!all_output_pages_touched(
+ &bitmap,
+ num_rows,
+ 0,
+ num_rows * size_of::<i64>(),
+ page_size,
+ ConstOutputLayout::FixedWidth {
+ bytes_per_row: size_of::<i64>(),
+ },
+ ));
+ }
+
+ #[test]
+ fn test_boolean_fallback_respects_non_null_and_bitmap_boundaries() {
+ let rows_per_chunk = test_page_size() * 8;
+ let num_rows = rows_per_chunk * 2 + 1;
+ let non_null_rows = (0..num_rows.div_ceil(8)).collect::<Vec<_>>();
+ let mut bitmap = null_bitmap(num_rows, &non_null_rows);
+ bitmap[num_rows / 8] = 0b0000_0001;
+ let data = read_all_const(
+ &Value::Boolean(true),
+ num_rows,
+ true,
+ &bitmap,
+ DataVariant::Boolean,
+ )
+ .unwrap();
+
+ let RawColumnData::Boolean(values) = data else {
+ panic!("expected Boolean CONST data");
+ };
+ for row in 0..num_rows {
+ let value = values[row / 8] & (1 << (row % 8)) != 0;
+ assert_eq!(value, row < non_null_rows.len(), "row {row}");
+ }
+ }
+
+ #[test]
+ fn test_timestamp_nanos_uses_each_buffer_element_width() {
+ let page_size = test_page_size();
+ let millis_rows_per_chunk = page_size / size_of::<i64>();
+ let nanos_rows_per_chunk = page_size / size_of::<i32>();
+ let num_rows = millis_rows_per_chunk *
CONST_FILL_ALL_MIN_NON_NULL_DENOMINATOR;
+ let non_null_rows =
(millis_rows_per_chunk..nanos_rows_per_chunk).collect::<Vec<_>>();
+ let bitmap = null_bitmap(num_rows, &non_null_rows);
+ let data = read_all_const(
+ &Value::TimestampNanos {
+ millis: 42,
+ nanos_of_milli: 7,
+ },
+ num_rows,
+ true,
+ &bitmap,
+ DataVariant::TimestampNanos,
+ )
+ .unwrap();
+
+ let RawColumnData::TimestampNanos {
+ millis,
+ nanos_of_milli,
+ } = data
+ else {
+ panic!("expected TimestampNanos CONST data");
+ };
+ assert!(millis[..millis_rows_per_chunk]
+ .iter()
+ .all(|&value| value == 0));
+ assert!(millis[millis_rows_per_chunk..nanos_rows_per_chunk]
+ .iter()
+ .all(|&value| value == 42));
+ assert!(millis[nanos_rows_per_chunk..]
+ .iter()
+ .all(|&value| value == 0));
+ assert!(nanos_of_milli[..millis_rows_per_chunk]
+ .iter()
+ .all(|&value| value == 0));
+ assert!(nanos_of_milli[millis_rows_per_chunk..nanos_rows_per_chunk]
+ .iter()
+ .all(|&value| value == 7));
+ assert!(nanos_of_milli[nanos_rows_per_chunk..]
+ .iter()
+ .all(|&value| value == 0));
+ }
+
+ #[test]
+ fn test_timestamp_nanos_conversion_ignores_null_hidden_values() {
+ let (min_millis, min_nanos) = types::ns_to_millis_nanos(i64::MIN);
+ let (max_millis, max_nanos) = types::ns_to_millis_nanos(i64::MAX);
+ let data = RawColumnData::TimestampNanos {
+ millis: vec![min_millis, min_millis, max_millis, max_millis],
+ nanos_of_milli: vec![min_nanos, 0, max_nanos, 999_999],
+ };
+ let array = build_array(
+ data,
+ &DataType::Timestamp(TimeUnit::Nanosecond, None),
+ Some(vec![0b0000_0101]),
+ 4,
+ true,
+ )
+ .unwrap();
+ let values = array
+ .as_any()
+ .downcast_ref::<TimestampNanosecondArray>()
+ .unwrap();
+
+ assert_eq!(values.value(0), i64::MIN);
+ assert!(values.is_null(1));
+ assert_eq!(values.value(2), i64::MAX);
+ assert!(values.is_null(3));
+ }
+
+ #[test]
+ fn test_fixed_const_preserves_default_and_negative_zero_values() {
+ let num_rows = test_page_size();
+ let non_null_rows =
+
(0..num_rows.div_ceil(CONST_FILL_ALL_MIN_NON_NULL_DENOMINATOR)).collect::<Vec<_>>();
+ let bitmap = null_bitmap(num_rows, &non_null_rows);
+
+ let zero_data = read_all_const(
+ &Value::SmallInt(0),
+ num_rows,
+ true,
+ &bitmap,
+ DataVariant::Int16,
+ )
+ .unwrap();
+ let RawColumnData::Int16(zero_values) = zero_data else {
+ panic!("expected Int16 CONST data");
+ };
+ assert!(zero_values.iter().all(|&value| value == 0));
+
+ let negative_zero_data = read_all_const(
+ &Value::Float(-0.0),
+ num_rows,
+ true,
+ &bitmap,
+ DataVariant::Float32,
+ )
+ .unwrap();
+ let RawColumnData::Float32(negative_zero_values) = negative_zero_data
else {
+ panic!("expected Float32 CONST data");
+ };
+ assert!(negative_zero_values[..non_null_rows.len()]
+ .iter()
+ .all(|value| value.to_bits() == (-0.0f32).to_bits()));
+ assert!(negative_zero_values[non_null_rows.len()..]
+ .iter()
+ .all(|value| value.to_bits() == 0.0f32.to_bits()));
+ }
+}
diff --git a/core/src/reader_tests.rs b/core/src/reader_tests.rs
index d89c2e7..87cdbb0 100644
--- a/core/src/reader_tests.rs
+++ b/core/src/reader_tests.rs
@@ -815,6 +815,67 @@ fn write_and_read(
(reader, data)
}
+#[test]
+fn test_nullable_const_timestamp_nanos_boundaries_roundtrip() {
+ use crate::reader::Encoding;
+
+ let columns = vec![
+ (
+ "ts_min".to_string(),
+ DataType::Timestamp(TimeUnit::Nanosecond, None),
+ true,
+ ),
+ (
+ "ts_max".to_string(),
+ DataType::Timestamp(TimeUnit::Nanosecond, None),
+ true,
+ ),
+ ];
+ let (min_millis, min_nanos) = crate::types::ns_to_millis_nanos(i64::MIN);
+ let (max_millis, max_nanos) = crate::types::ns_to_millis_nanos(i64::MAX);
+ let rows: Vec<Vec<Value>> = (0..200)
+ .map(|row| {
+ if row < 25 {
+ vec![
+ Value::TimestampNanos {
+ millis: min_millis,
+ nanos_of_milli: min_nanos,
+ },
+ Value::TimestampNanos {
+ millis: max_millis,
+ nanos_of_milli: max_nanos,
+ },
+ ]
+ } else {
+ vec![Value::Null, Value::Null]
+ }
+ })
+ .collect();
+
+ let (reader, _) = write_and_read_paged(columns, &rows);
+ assert!(reader
+ .page_infos(0)
+ .unwrap()
+ .iter()
+ .all(|info| info.encoding == Encoding::Const));
+ let mut rg = reader.row_group_reader(0).unwrap();
+ let batch = rg.read_columns().unwrap();
+ for (column, expected) in [(0, i64::MIN), (1, i64::MAX)] {
+ let values = batch
+ .column(column)
+ .as_any()
+ .downcast_ref::<TimestampNanosecondArray>()
+ .unwrap();
+ assert_eq!(values.len(), 200);
+ assert_eq!(values.null_count(), 175);
+ for row in 0..200 {
+ if !values.is_null(row) {
+ assert_eq!(values.value(row), expected, "column {column}, row
{row}");
+ }
+ }
+ }
+}
+
/// Mirrors Paimon's testSchemaEvolutionTypeWidening.
/// Writes with narrow types (INT, FLOAT, TINYINT) and verifies the reader
/// produces the exact values that a higher layer can widen to
BIGINT/DOUBLE/INT.
@@ -1741,6 +1802,198 @@ fn test_const_encoding_with_nulls() {
}
}
+fn assert_const_encoding_with_nulls_all_primitive_types(batch: &RecordBatch,
num_rows: usize) {
+ let bools = batch_col_bool(batch, "bool");
+ let tiny = batch_col_i8(batch, "tiny");
+ let small = batch_col_i16(batch, "small");
+ let ints = batch_col_i32(batch, "int");
+ let big = batch_col_i64(batch, "big");
+ let floats = batch_col_f32(batch, "float");
+ let doubles = batch_col_f64(batch, "double");
+ let strings = batch_col_string(batch, "string");
+ assert_eq!(
+ strings.value_data().len(),
+ 18 * "constant".len(),
+ "variable-width CONST values should remain compact at null positions"
+ );
+ for i in 0..num_rows {
+ if i % 4 == 0 {
+ assert!(bools.is_null(i));
+ assert!(tiny.is_null(i));
+ assert!(small.is_null(i));
+ assert!(ints.is_null(i));
+ assert!(big.is_null(i));
+ assert!(floats.is_null(i));
+ assert!(doubles.is_null(i));
+ assert!(strings.is_null(i));
+ } else {
+ assert!(bools.value(i));
+ assert_eq!(tiny.value(i), -7);
+ assert_eq!(small.value(i), 1234);
+ assert_eq!(ints.value(i), -56789);
+ assert_eq!(big.value(i), 9_876_543_210);
+ assert_eq!(floats.value(i), 1.25);
+ assert_eq!(doubles.value(i), -12.5);
+ assert_eq!(strings.value(i), "constant");
+ }
+ }
+}
+
+#[test]
+fn test_const_encoding_with_nulls_all_primitive_types() {
+ let columns = vec![
+ ("bool".to_string(), DataType::Boolean, true),
+ ("tiny".to_string(), DataType::Int8, true),
+ ("small".to_string(), DataType::Int16, true),
+ ("int".to_string(), DataType::Int32, true),
+ ("big".to_string(), DataType::Int64, true),
+ ("float".to_string(), DataType::Float32, true),
+ ("double".to_string(), DataType::Float64, true),
+ ("string".to_string(), DataType::Utf8, true),
+ ];
+ let rows: Vec<Vec<Value>> = (0..24)
+ .map(|i| {
+ if i % 4 == 0 {
+ vec![Value::Null; columns.len()]
+ } else {
+ vec![
+ Value::Boolean(true),
+ Value::TinyInt(-7),
+ Value::SmallInt(1234),
+ Value::Integer(-56789),
+ Value::BigInt(9_876_543_210),
+ Value::Float(1.25),
+ Value::Double(-12.5),
+ Value::String(b"constant".to_vec()),
+ ]
+ }
+ })
+ .collect();
+
+ let (reader, _) = write_and_read(columns.clone(), &rows);
+ let mut rg = reader.row_group_reader(0).unwrap();
+ let batch = rg.read_columns().unwrap();
+ assert_const_encoding_with_nulls_all_primitive_types(&batch, rows.len());
+
+ let (reader, _) = write_and_read_paged(columns, &rows);
+ let mut rg = reader.row_group_reader(0).unwrap();
+ let batch = rg.read_columns().unwrap();
+ assert_const_encoding_with_nulls_all_primitive_types(&batch, rows.len());
+}
+
+fn assert_sparse_const_encoding_with_nulls(batch: &RecordBatch, non_null_rows:
&[usize]) {
+ let bools = batch_col_bool(batch, "bool");
+ let tiny = batch_col_i8(batch, "tiny");
+ let small = batch_col_i16(batch, "small");
+ let ints = batch_col_i32(batch, "int");
+ let big = batch_col_i64(batch, "big");
+ let floats = batch_col_f32(batch, "float");
+ let doubles = batch_col_f64(batch, "double");
+ let timestamp_nanos = batch
+ .column_by_name("timestamp_nanos")
+ .unwrap()
+ .as_any()
+ .downcast_ref::<TimestampNanosecondArray>()
+ .unwrap();
+ let expected_timestamp_nanos =
+ crate::types::millis_nanos_to_ns(1_700_000_000_000, 123_456).unwrap();
+
+ for i in 0..batch.num_rows() {
+ if non_null_rows.contains(&i) {
+ assert!(!bools.is_null(i));
+ assert!(!tiny.is_null(i));
+ assert!(!small.is_null(i));
+ assert!(!ints.is_null(i));
+ assert!(!big.is_null(i));
+ assert!(!floats.is_null(i));
+ assert!(!doubles.is_null(i));
+ assert!(!timestamp_nanos.is_null(i));
+ assert!(bools.value(i));
+ assert_eq!(tiny.value(i), -7);
+ assert_eq!(small.value(i), 1234);
+ assert_eq!(ints.value(i), -56_789);
+ assert_eq!(big.value(i), 9_876_543_210);
+ assert_eq!(floats.value(i), 1.25);
+ assert_eq!(doubles.value(i), -12.5);
+ assert_eq!(timestamp_nanos.value(i), expected_timestamp_nanos);
+ } else {
+ assert!(bools.is_null(i));
+ assert!(tiny.is_null(i));
+ assert!(small.is_null(i));
+ assert!(ints.is_null(i));
+ assert!(big.is_null(i));
+ assert!(floats.is_null(i));
+ assert!(doubles.is_null(i));
+ assert!(timestamp_nanos.is_null(i));
+ }
+ }
+}
+
+#[test]
+fn test_sparse_const_encoding_with_nulls() {
+ use crate::reader::Encoding;
+
+ let columns = vec![
+ ("bool".to_string(), DataType::Boolean, true),
+ ("tiny".to_string(), DataType::Int8, true),
+ ("small".to_string(), DataType::Int16, true),
+ ("int".to_string(), DataType::Int32, true),
+ ("big".to_string(), DataType::Int64, true),
+ ("float".to_string(), DataType::Float32, true),
+ ("double".to_string(), DataType::Float64, true),
+ (
+ "timestamp_nanos".to_string(),
+ DataType::Timestamp(TimeUnit::Nanosecond, None),
+ true,
+ ),
+ ];
+ let non_null_rows = [0, 127, 255];
+ let rows: Vec<Vec<Value>> = (0..256)
+ .map(|i| {
+ if non_null_rows.contains(&i) {
+ vec![
+ Value::Boolean(true),
+ Value::TinyInt(-7),
+ Value::SmallInt(1234),
+ Value::Integer(-56_789),
+ Value::BigInt(9_876_543_210),
+ Value::Float(1.25),
+ Value::Double(-12.5),
+ Value::TimestampNanos {
+ millis: 1_700_000_000_000,
+ nanos_of_milli: 123_456,
+ },
+ ]
+ } else {
+ vec![Value::Null; columns.len()]
+ }
+ })
+ .collect();
+ assert!(non_null_rows.len() < rows.len().div_ceil(8));
+
+ let (reader, _) = write_and_read(columns.clone(), &rows);
+ assert!(reader
+ .page_infos(0)
+ .unwrap()
+ .iter()
+ .all(|info| info.encoding == Encoding::Const));
+ let mut rg = reader.row_group_reader(0).unwrap();
+ let batch = rg.read_columns().unwrap();
+ assert_eq!(batch.num_rows(), rows.len());
+ assert_sparse_const_encoding_with_nulls(&batch, &non_null_rows);
+
+ let (reader, _) = write_and_read_paged(columns, &rows);
+ assert!(reader
+ .page_infos(0)
+ .unwrap()
+ .iter()
+ .all(|info| info.encoding == Encoding::Const));
+ let mut rg = reader.row_group_reader(0).unwrap();
+ let batch = rg.read_columns().unwrap();
+ assert_eq!(batch.num_rows(), rows.len());
+ assert_sparse_const_encoding_with_nulls(&batch, &non_null_rows);
+}
+
#[test]
fn test_dict_encoding_with_nulls() {
let columns = vec![("v".to_string(), DataType::Int32, true)];