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-rust.git
The following commit(s) were added to refs/heads/main by this push:
new 0923cbc5 Support BSI indexes and indexed native write paths (#947)
0923cbc5 is described below
commit 0923cbc5befc171d27452db942034732c3bf9295
Author: Jingsong Lee <[email protected]>
AuthorDate: Thu Sep 24 22:03:14 2026 +0800
Support BSI indexes and indexed native write paths (#947)
---
crates/paimon/src/file_index/bsi.rs | 519 ++++++++++++
crates/paimon/src/file_index/bsi/tests.rs | 498 ++++++++++++
.../paimon/src/file_index/file_indexer_factory.rs | 31 +-
crates/paimon/src/file_index/mod.rs | 2 +
crates/paimon/src/table/data_evolution_reader.rs | 48 +-
crates/paimon/src/table/data_evolution_writer.rs | 22 +-
crates/paimon/src/table/data_file_index_writer.rs | 26 +-
.../src/table/data_file_index_writer/tests.rs | 888 ++++++++++++++++++++-
crates/paimon/src/table/data_file_reader.rs | 2 +-
crates/paimon/src/table/kv_file_reader.rs | 35 +-
crates/paimon/src/table/kv_file_writer.rs | 79 +-
crates/paimon/src/table/postpone_file_writer.rs | 107 ++-
crates/paimon/src/table/table_write.rs | 46 +-
13 files changed, 2202 insertions(+), 101 deletions(-)
diff --git a/crates/paimon/src/file_index/bsi.rs
b/crates/paimon/src/file_index/bsi.rs
new file mode 100644
index 00000000..7b71bab9
--- /dev/null
+++ b/crates/paimon/src/file_index/bsi.rs
@@ -0,0 +1,519 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Java V1 `bsi` FileIndex for signed integral, date, time, decimal and
timestamp values.
+//!
+//! Java stores two unsigned bit-sliced indexes: one for nonnegative values and
+//! one for the absolute values of negative values. The outer header and each
+//! slice's min/max use big-endian Java primitive encoding; Roaring bitmaps use
+//! the portable bitmap encoding used by `RoaringBitmap32`.
+
+use std::io::{Cursor, Read};
+
+use bytes::{BufMut, Bytes};
+use roaring::RoaringBitmap;
+
+use crate::common::Options;
+use crate::file_index::file_index_reader::FileIndexReader;
+use crate::file_index::file_index_result::FileIndexResult;
+use crate::file_index::file_index_writer::FileIndexWriter;
+use crate::spec::{DataType, Datum, PredicateOperator};
+use crate::{Error, Result};
+
+const VERSION_1: u8 = 1;
+
+fn invalid(message: impl Into<String>) -> Error {
+ Error::FileIndexFormatInvalid {
+ message: message.into(),
+ }
+}
+
+fn write_error(message: impl Into<String>) -> Error {
+ Error::DataInvalid {
+ message: message.into(),
+ source: None,
+ }
+}
+
+/// The BSI value mapper must agree with Java's `getValueMapper` and must never
+/// silently narrow an integer or decimal. A bad literal disables index
pruning.
+fn mapped_value(data_type: &DataType, datum: &Datum) -> Option<i64> {
+ match (data_type, datum) {
+ (DataType::TinyInt(_), Datum::TinyInt(value)) =>
Some(i64::from(*value)),
+ (DataType::SmallInt(_), Datum::SmallInt(value)) =>
Some(i64::from(*value)),
+ (DataType::Int(_), Datum::Int(value)) => Some(i64::from(*value)),
+ (DataType::BigInt(_), Datum::Long(value)) => Some(*value),
+ (DataType::Date(_), Datum::Date(value)) => Some(i64::from(*value)),
+ (DataType::Time(_), Datum::Time(value)) => Some(i64::from(*value)),
+ (
+ DataType::Decimal(data_type),
+ Datum::Decimal {
+ unscaled,
+ precision,
+ scale,
+ },
+ ) if *precision == data_type.precision() && *scale ==
data_type.scale() => {
+ i64::try_from(*unscaled).ok()
+ }
+ (DataType::Timestamp(data_type), Datum::Timestamp { millis, nanos })
=> {
+ timestamp_value(data_type.precision(), *millis, *nanos)
+ }
+ (
+ DataType::LocalZonedTimestamp(data_type),
+ Datum::LocalZonedTimestamp { millis, nanos },
+ ) => timestamp_value(data_type.precision(), *millis, *nanos),
+ _ => None,
+ }
+}
+
+fn timestamp_value(precision: u32, millis: i64, nanos: i32) -> Option<i64> {
+ if !(0..1_000_000).contains(&nanos) {
+ return None;
+ }
+ if precision <= 3 {
+ Some(millis)
+ } else {
+ millis
+ .checked_mul(1_000)?
+ .checked_add(i64::from(nanos / 1_000))
+ }
+}
+
+fn validate_type(data_type: &DataType) -> Result<()> {
+ match data_type {
+ DataType::TinyInt(_)
+ | DataType::SmallInt(_)
+ | DataType::Int(_)
+ | DataType::BigInt(_)
+ | DataType::Date(_)
+ | DataType::Time(_)
+ | DataType::Timestamp(_)
+ | DataType::LocalZonedTimestamp(_) => Ok(()),
+ DataType::Decimal(_) => Ok(()),
+ _ => Err(Error::Unsupported {
+ message: format!("BSI file index does not support data type
{data_type:?}"),
+ }),
+ }
+}
+
+pub(crate) struct BsiFileIndexWriter {
+ data_type: DataType,
+ values: Vec<Option<i64>>,
+}
+
+impl BsiFileIndexWriter {
+ pub(crate) fn try_new(data_type: DataType, _options: &Options) ->
Result<Self> {
+ validate_type(&data_type)?;
+ Ok(Self {
+ data_type,
+ values: Vec::new(),
+ })
+ }
+}
+
+impl FileIndexWriter for BsiFileIndexWriter {
+ fn write(&mut self, datum: Option<&Datum>) -> Result<()> {
+ if self.values.len() >= i32::MAX as usize {
+ return Err(write_error("BSI row count exceeds Java int range"));
+ }
+ let value = datum
+ .map(|datum| {
+ mapped_value(&self.data_type, datum).ok_or_else(|| {
+ write_error(format!(
+ "Datum {datum:?} does not match BSI type {:?}",
+ self.data_type
+ ))
+ })
+ })
+ .transpose()?;
+ if value == Some(i64::MIN) {
+ // Java's Math.abs(Long.MIN_VALUE) overflows its negative slice.
+ return Err(write_error("BSI cannot encode Long.MIN_VALUE"));
+ }
+ self.values.push(value);
+ Ok(())
+ }
+
+ fn serialized_bytes(&mut self) -> Result<Bytes> {
+ let row_count = i32::try_from(self.values.len())
+ .map_err(|_| write_error("BSI row count exceeds Java int range"))?;
+ let mut positive = Vec::new();
+ let mut negative = Vec::new();
+ for (row, value) in self.values.iter().enumerate() {
+ match value {
+ Some(value) if *value < 0 => negative.push((row as u32,
value.unsigned_abs())),
+ Some(value) => positive.push((row as u32, *value as u64)),
+ None => {}
+ }
+ }
+ let mut bytes = Vec::new();
+ bytes.put_u8(VERSION_1);
+ bytes.put_i32(row_count);
+ bytes.put_u8(u8::from(!positive.is_empty()));
+ if !positive.is_empty() {
+ write_slice_index(&mut bytes, &positive)?;
+ }
+ bytes.put_u8(u8::from(!negative.is_empty()));
+ if !negative.is_empty() {
+ write_slice_index(&mut bytes, &negative)?;
+ }
+ Ok(Bytes::from(bytes))
+ }
+
+ fn empty(&self) -> bool {
+ self.values.is_empty()
+ }
+}
+
+fn write_bitmap(bytes: &mut Vec<u8>, bitmap: &RoaringBitmap) -> Result<()> {
+ bitmap
+ .serialize_into(bytes)
+ .map_err(|error| write_error(format!("Failed to serialize BSI bitmap:
{error}")))
+}
+
+fn write_slice_index(bytes: &mut Vec<u8>, values: &[(u32, u64)]) -> Result<()>
{
+ let max = values.iter().map(|(_, value)| *value).max().unwrap();
+ let width = (u64::BITS - max.leading_zeros()) as usize;
+ let mut existing = RoaringBitmap::new();
+ let mut slices = vec![RoaringBitmap::new(); width];
+ for &(row, mut value) in values {
+ existing.insert(row);
+ while value != 0 {
+ slices[value.trailing_zeros() as usize].insert(row);
+ value &= value - 1;
+ }
+ }
+ bytes.put_u8(VERSION_1);
+ bytes.put_i64(0); // Java's StatsCollectList starts
positiveMin/negativeMin at zero.
+ bytes.put_i64(max as i64);
+ write_bitmap(bytes, &existing)?;
+ bytes.put_i32(width as i32);
+ for slice in slices {
+ write_bitmap(bytes, &slice)?;
+ }
+ Ok(())
+}
+
+fn read_exact<const N: usize>(input: &mut Cursor<&[u8]>, label: &str) ->
Result<[u8; N]> {
+ let mut bytes = [0; N];
+ input
+ .read_exact(&mut bytes)
+ .map_err(|error| invalid(format!("Truncated BSI {label}: {error}")))?;
+ Ok(bytes)
+}
+
+fn read_u8(input: &mut Cursor<&[u8]>, label: &str) -> Result<u8> {
+ Ok(read_exact::<1>(input, label)?[0])
+}
+
+fn read_i32(input: &mut Cursor<&[u8]>, label: &str) -> Result<i32> {
+ Ok(i32::from_be_bytes(read_exact(input, label)?))
+}
+
+fn read_i64(input: &mut Cursor<&[u8]>, label: &str) -> Result<i64> {
+ Ok(i64::from_be_bytes(read_exact(input, label)?))
+}
+
+fn read_bitmap(input: &mut Cursor<&[u8]>, row_count: u32, label: &str) ->
Result<RoaringBitmap> {
+ let bitmap = RoaringBitmap::deserialize_from(input)
+ .map_err(|error| invalid(format!("Invalid BSI {label} bitmap:
{error}")))?;
+ if bitmap.max().is_some_and(|row| row >= row_count) {
+ return Err(invalid(format!(
+ "BSI {label} contains a row past {row_count}"
+ )));
+ }
+ Ok(bitmap)
+}
+
+struct SliceIndex {
+ max: i64,
+ existing: RoaringBitmap,
+ slices: Vec<RoaringBitmap>,
+}
+
+impl SliceIndex {
+ fn read(input: &mut Cursor<&[u8]>, row_count: u32) -> Result<Self> {
+ if read_u8(input, "slice version")? != VERSION_1 {
+ return Err(invalid("Unsupported BSI slice version"));
+ }
+ let min = read_i64(input, "slice min")?;
+ let max = read_i64(input, "slice max")?;
+ if min != 0 || max < 0 {
+ return Err(invalid("Invalid BSI slice min/max"));
+ }
+ let existing = read_bitmap(input, row_count, "existence")?;
+ let width = read_i32(input, "slice count")?;
+ let expected = i64::BITS - max.leading_zeros();
+ if width < 0 || width as u32 != expected {
+ return Err(invalid("BSI slice count does not match max"));
+ }
+ let mut slices = Vec::with_capacity(width as usize);
+ for _ in 0..width {
+ let bitmap = read_bitmap(input, row_count, "slice")?;
+ if !bitmap.is_subset(&existing) {
+ return Err(invalid("BSI slice has rows outside existence
bitmap"));
+ }
+ slices.push(bitmap);
+ }
+ Ok(Self {
+ max,
+ existing,
+ slices,
+ })
+ }
+
+ fn compare(&self, operator: PredicateOperator, value: i64) ->
RoaringBitmap {
+ use PredicateOperator::*;
+ if value < 0 {
+ return match operator {
+ Eq | Lt | LtEq => RoaringBitmap::new(),
+ NotEq | Gt | GtEq => self.existing.clone(),
+ _ => RoaringBitmap::new(),
+ };
+ }
+ if value > self.max {
+ return match operator {
+ NotEq | Lt | LtEq => self.existing.clone(),
+ _ => RoaringBitmap::new(),
+ };
+ }
+ let mut equal = self.existing.clone();
+ let mut less = RoaringBitmap::new();
+ let mut greater = RoaringBitmap::new();
+ for (bit, slice) in self.slices.iter().enumerate().rev() {
+ if (value >> bit) & 1 == 1 {
+ less |= &equal - slice;
+ equal &= slice;
+ } else {
+ greater |= &equal & slice;
+ equal -= slice;
+ }
+ }
+ match operator {
+ Eq => equal,
+ NotEq => &self.existing - &equal,
+ Lt => less,
+ LtEq => &less | &equal,
+ Gt => greater,
+ GtEq => &greater | &equal,
+ _ => RoaringBitmap::new(),
+ }
+ }
+}
+
+pub(crate) struct BsiFileIndexReader {
+ data_type: DataType,
+ row_count: u32,
+ positive: Option<SliceIndex>,
+ negative: Option<SliceIndex>,
+}
+
+impl BsiFileIndexReader {
+ pub(crate) fn try_new(data_type: DataType, serialized: Bytes) ->
Result<Self> {
+ validate_type(&data_type)?;
+ let mut input = Cursor::new(serialized.as_ref());
+ if read_u8(&mut input, "version")? != VERSION_1 {
+ return Err(invalid("Unsupported BSI version"));
+ }
+ let row_count = u32::try_from(read_i32(&mut input, "row count")?)
+ .map_err(|_| invalid("Negative BSI row count"))?;
+ let positive = match read_u8(&mut input, "positive flag")? {
+ 0 => None,
+ 1 => Some(SliceIndex::read(&mut input, row_count)?),
+ _ => return Err(invalid("Invalid BSI positive flag")),
+ };
+ let negative = match read_u8(&mut input, "negative flag")? {
+ 0 => None,
+ 1 => Some(SliceIndex::read(&mut input, row_count)?),
+ _ => return Err(invalid("Invalid BSI negative flag")),
+ };
+ if input.position() as usize != serialized.len() {
+ return Err(invalid("Trailing BSI payload bytes"));
+ }
+ if let (Some(positive), Some(negative)) = (&positive, &negative) {
+ if !positive.existing.is_disjoint(&negative.existing) {
+ return Err(invalid("A BSI row appears in both sign indexes"));
+ }
+ }
+ Ok(Self {
+ data_type,
+ row_count,
+ positive,
+ negative,
+ })
+ }
+
+ fn non_null(&self) -> RoaringBitmap {
+ let mut rows = RoaringBitmap::new();
+ if let Some(positive) = &self.positive {
+ rows |= &positive.existing;
+ }
+ if let Some(negative) = &self.negative {
+ rows |= &negative.existing;
+ }
+ rows
+ }
+
+ fn compare(&self, operator: PredicateOperator, value: i64) ->
RoaringBitmap {
+ use PredicateOperator::*;
+ let positive = self.positive.as_ref();
+ let negative = self.negative.as_ref();
+ if value == i64::MIN {
+ return match operator {
+ Lt | LtEq | Eq => RoaringBitmap::new(),
+ NotEq | Gt | GtEq => self.non_null(),
+ _ => RoaringBitmap::new(),
+ };
+ }
+ if value < 0 {
+ let abs = -value;
+ match operator {
+ Eq => negative.map_or_else(RoaringBitmap::new, |idx|
idx.compare(Eq, abs)),
+ NotEq => {
+ let equal =
+ negative.map_or_else(RoaringBitmap::new, |idx|
idx.compare(Eq, abs));
+ &self.non_null() - &equal
+ }
+ Lt => negative.map_or_else(RoaringBitmap::new, |idx|
idx.compare(Gt, abs)),
+ LtEq => negative.map_or_else(RoaringBitmap::new, |idx|
idx.compare(GtEq, abs)),
+ Gt => {
+ let mut rows =
+ positive.map_or_else(RoaringBitmap::new, |idx|
idx.existing.clone());
+ if let Some(idx) = negative {
+ rows |= idx.compare(Lt, abs);
+ }
+ rows
+ }
+ GtEq => {
+ let mut rows =
+ positive.map_or_else(RoaringBitmap::new, |idx|
idx.existing.clone());
+ if let Some(idx) = negative {
+ rows |= idx.compare(LtEq, abs);
+ }
+ rows
+ }
+ _ => RoaringBitmap::new(),
+ }
+ } else {
+ match operator {
+ Eq => positive.map_or_else(RoaringBitmap::new, |idx|
idx.compare(Eq, value)),
+ NotEq => {
+ let equal =
+ positive.map_or_else(RoaringBitmap::new, |idx|
idx.compare(Eq, value));
+ &self.non_null() - &equal
+ }
+ Lt => {
+ let mut rows =
+ negative.map_or_else(RoaringBitmap::new, |idx|
idx.existing.clone());
+ if let Some(idx) = positive {
+ rows |= idx.compare(Lt, value);
+ }
+ rows
+ }
+ LtEq => {
+ let mut rows =
+ negative.map_or_else(RoaringBitmap::new, |idx|
idx.existing.clone());
+ if let Some(idx) = positive {
+ rows |= idx.compare(LtEq, value);
+ }
+ rows
+ }
+ Gt => positive.map_or_else(RoaringBitmap::new, |idx|
idx.compare(Gt, value)),
+ GtEq => positive.map_or_else(RoaringBitmap::new, |idx|
idx.compare(GtEq, value)),
+ _ => RoaringBitmap::new(),
+ }
+ }
+ }
+}
+
+impl FileIndexReader for BsiFileIndexReader {
+ fn evaluate(
+ &self,
+ _column: &str,
+ _index: usize,
+ _data_type: &DataType,
+ operator: PredicateOperator,
+ literals: &[Datum],
+ ) -> FileIndexResult {
+ use PredicateOperator::*;
+ let rows = match operator {
+ IsNull => {
+ let mut rows = RoaringBitmap::new();
+ rows.insert_range(0..self.row_count);
+ rows -= &self.non_null();
+ rows
+ }
+ IsNotNull => self.non_null(),
+ Eq | NotEq | Lt | LtEq | Gt | GtEq => {
+ let Some(value) = literals
+ .first()
+ .and_then(|datum| mapped_value(&self.data_type, datum))
+ else {
+ return FileIndexResult::Remain;
+ };
+ if self.truncated_timestamp() {
+ return FileIndexResult::Remain;
+ }
+ self.compare(operator, value)
+ }
+ In | NotIn => {
+ if self.truncated_timestamp() {
+ return FileIndexResult::Remain;
+ }
+ let mut equal = RoaringBitmap::new();
+ for literal in literals {
+ let Some(value) = mapped_value(&self.data_type, literal)
else {
+ return FileIndexResult::Remain;
+ };
+ equal |= self.compare(Eq, value);
+ }
+ if operator == NotIn {
+ &self.non_null() - &equal
+ } else {
+ equal
+ }
+ }
+ Between if literals.len() == 2 => {
+ if self.truncated_timestamp() {
+ return FileIndexResult::Remain;
+ }
+ let (Some(lower), Some(upper)) = (
+ mapped_value(&self.data_type, &literals[0]),
+ mapped_value(&self.data_type, &literals[1]),
+ ) else {
+ return FileIndexResult::Remain;
+ };
+ &self.compare(GtEq, lower) & &self.compare(LtEq, upper)
+ }
+ _ => return FileIndexResult::Remain,
+ };
+ FileIndexResult::Selection(rows)
+ }
+}
+
+impl BsiFileIndexReader {
+ fn truncated_timestamp(&self) -> bool {
+ match &self.data_type {
+ DataType::Timestamp(data_type) => data_type.precision() > 6,
+ DataType::LocalZonedTimestamp(data_type) => data_type.precision()
> 6,
+ _ => false,
+ }
+ }
+}
+
+#[cfg(test)]
+mod tests;
diff --git a/crates/paimon/src/file_index/bsi/tests.rs
b/crates/paimon/src/file_index/bsi/tests.rs
new file mode 100644
index 00000000..3a4c992f
--- /dev/null
+++ b/crates/paimon/src/file_index/bsi/tests.rs
@@ -0,0 +1,498 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+use super::*;
+use crate::spec::{
+ BigIntType, BooleanType, DateType, DecimalType, IntType,
LocalZonedTimestampType, SmallIntType,
+ TimeType, TimestampType, TinyIntType,
+};
+use base64::Engine;
+
+fn int_type() -> DataType {
+ DataType::Int(IntType::new())
+}
+
+fn int_reader(values: &[Option<i32>]) -> BsiFileIndexReader {
+ let mut writer = BsiFileIndexWriter::try_new(int_type(),
&Options::new()).unwrap();
+ for value in values {
+ let datum = value.map(Datum::Int);
+ writer.write(datum.as_ref()).unwrap();
+ }
+ BsiFileIndexReader::try_new(int_type(),
writer.serialized_bytes().unwrap()).unwrap()
+}
+
+fn selection(rows: impl IntoIterator<Item = u32>) -> FileIndexResult {
+ FileIndexResult::Selection(rows.into_iter().collect())
+}
+
+fn eval(
+ reader: &BsiFileIndexReader,
+ operator: PredicateOperator,
+ literals: &[i32],
+) -> FileIndexResult {
+ reader.evaluate(
+ "value",
+ 0,
+ &int_type(),
+ operator,
+ &literals.iter().copied().map(Datum::Int).collect::<Vec<_>>(),
+ )
+}
+
+#[test]
+fn test_bsi_signs_nulls_and_comparisons() {
+ let reader = int_reader(&[None, Some(-5), Some(-2), Some(0), Some(3),
Some(8), None]);
+ use PredicateOperator::*;
+ assert_eq!(eval(&reader, IsNull, &[]), selection([0, 6]));
+ assert_eq!(eval(&reader, IsNotNull, &[]), selection([1, 2, 3, 4, 5]));
+ assert_eq!(eval(&reader, Eq, &[-2]), selection([2]));
+ assert_eq!(eval(&reader, Eq, &[0]), selection([3]));
+ assert_eq!(eval(&reader, Eq, &[3]), selection([4]));
+ assert_eq!(eval(&reader, Eq, &[17]), selection([]));
+ assert_eq!(eval(&reader, NotEq, &[-2]), selection([1, 3, 4, 5]));
+ assert_eq!(eval(&reader, Lt, &[-2]), selection([1]));
+ assert_eq!(eval(&reader, LtEq, &[-2]), selection([1, 2]));
+ assert_eq!(eval(&reader, Gt, &[-2]), selection([3, 4, 5]));
+ assert_eq!(eval(&reader, GtEq, &[-2]), selection([2, 3, 4, 5]));
+ assert_eq!(eval(&reader, Lt, &[0]), selection([1, 2]));
+ assert_eq!(eval(&reader, LtEq, &[0]), selection([1, 2, 3]));
+ assert_eq!(eval(&reader, Gt, &[0]), selection([4, 5]));
+ assert_eq!(eval(&reader, GtEq, &[0]), selection([3, 4, 5]));
+ assert_eq!(eval(&reader, Lt, &[9]), selection([1, 2, 3, 4, 5]));
+ assert_eq!(eval(&reader, Gt, &[9]), selection([]));
+ assert_eq!(eval(&reader, LtEq, &[-9]), selection([]));
+ assert_eq!(eval(&reader, GtEq, &[-9]), selection([1, 2, 3, 4, 5]));
+ assert_eq!(eval(&reader, In, &[-5, 8, 20]), selection([1, 5]));
+ assert_eq!(eval(&reader, NotIn, &[-5, 8]), selection([2, 3, 4]));
+ assert_eq!(eval(&reader, Between, &[-2, 3]), selection([2, 3, 4]));
+ assert_eq!(eval(&reader, NotBetween, &[-2, 3]), FileIndexResult::Remain);
+}
+
+#[test]
+fn test_bsi_matches_scalar_for_varied_signed_values() {
+ use PredicateOperator::*;
+ let values = [
+ None,
+ Some(-32768),
+ Some(-256),
+ Some(-65),
+ Some(-8),
+ Some(-1),
+ Some(0),
+ Some(1),
+ Some(7),
+ Some(64),
+ Some(255),
+ Some(32767),
+ ];
+ let reader = int_reader(&values);
+ for literal in [-50000, -32768, -255, -8, -1, 0, 1, 8, 64, 255, 32767,
50000] {
+ for operator in [Eq, NotEq, Lt, LtEq, Gt, GtEq] {
+ let expected: RoaringBitmap = values
+ .iter()
+ .enumerate()
+ .filter_map(|(row, value)| {
+ value.and_then(|value| {
+ let matches = match operator {
+ Eq => value == literal,
+ NotEq => value != literal,
+ Lt => value < literal,
+ LtEq => value <= literal,
+ Gt => value > literal,
+ GtEq => value >= literal,
+ _ => unreachable!(),
+ };
+ matches.then_some(row as u32)
+ })
+ })
+ .collect();
+ assert_eq!(
+ eval(&reader, operator, &[literal]),
+ FileIndexResult::Selection(expected),
+ "{operator:?} {literal}"
+ );
+ }
+ }
+}
+
+#[test]
+fn test_bsi_empty_and_all_null_files() {
+ for values in [&[][..], &[None, None][..]] {
+ let reader = int_reader(values);
+ assert_eq!(
+ eval(&reader, PredicateOperator::IsNull, &[]),
+ selection(0..values.len() as u32)
+ );
+ assert_eq!(
+ eval(&reader, PredicateOperator::IsNotNull, &[]),
+ selection([])
+ );
+ assert_eq!(eval(&reader, PredicateOperator::Eq, &[1]), selection([]));
+ }
+}
+
+#[test]
+fn test_bsi_writer_rejects_unrepresentable_values_and_types() {
+ assert!(matches!(
+ BsiFileIndexWriter::try_new(DataType::Boolean(BooleanType::new()),
&Options::new()),
+ Err(Error::Unsupported { .. })
+ ));
+ let mut wide_decimal = BsiFileIndexWriter::try_new(
+ DataType::Decimal(DecimalType::new(38, 2).unwrap()),
+ &Options::new(),
+ )
+ .unwrap();
+ assert!(matches!(
+ wide_decimal.write(Some(&Datum::Decimal {
+ unscaled: i128::MAX,
+ precision: 38,
+ scale: 2,
+ })),
+ Err(Error::DataInvalid { .. })
+ ));
+ let mut writer = BsiFileIndexWriter::try_new(int_type(),
&Options::new()).unwrap();
+ assert!(matches!(
+ writer.write(Some(&Datum::Long(1))),
+ Err(Error::DataInvalid { .. })
+ ));
+ assert!(writer.empty());
+ let mut long = BsiFileIndexWriter::try_new(
+ DataType::BigInt(crate::spec::BigIntType::new()),
+ &Options::new(),
+ )
+ .unwrap();
+ assert!(matches!(
+ long.write(Some(&Datum::Long(i64::MIN))),
+ Err(Error::DataInvalid { .. })
+ ));
+}
+
+#[test]
+fn test_bsi_decimal_and_timestamp_mapping() {
+ let decimal_type = DataType::Decimal(DecimalType::new(12, 2).unwrap());
+ let mut writer = BsiFileIndexWriter::try_new(decimal_type.clone(),
&Options::new()).unwrap();
+ for value in [-500_i128, 0, 275] {
+ writer
+ .write(Some(&Datum::Decimal {
+ unscaled: value,
+ precision: 12,
+ scale: 2,
+ }))
+ .unwrap();
+ }
+ let reader =
+ BsiFileIndexReader::try_new(decimal_type.clone(),
writer.serialized_bytes().unwrap())
+ .unwrap();
+ assert_eq!(
+ reader.evaluate(
+ "amount",
+ 0,
+ &decimal_type,
+ PredicateOperator::Lt,
+ &[Datum::Decimal {
+ unscaled: 100,
+ precision: 12,
+ scale: 2
+ }]
+ ),
+ selection([0, 1])
+ );
+ let timestamp_type = DataType::Timestamp(TimestampType::new(9).unwrap());
+ let mut writer = BsiFileIndexWriter::try_new(timestamp_type.clone(),
&Options::new()).unwrap();
+ writer
+ .write(Some(&Datum::Timestamp {
+ millis: 1000,
+ nanos: 123_456,
+ }))
+ .unwrap();
+ writer.write(None).unwrap();
+ let reader =
+ BsiFileIndexReader::try_new(timestamp_type.clone(),
writer.serialized_bytes().unwrap())
+ .unwrap();
+ assert_eq!(
+ reader.evaluate("ts", 0, ×tamp_type, PredicateOperator::IsNull,
&[]),
+ selection([1])
+ );
+ assert_eq!(
+ reader.evaluate(
+ "ts",
+ 0,
+ ×tamp_type,
+ PredicateOperator::Eq,
+ &[Datum::Timestamp {
+ millis: 1000,
+ nanos: 123_456
+ }]
+ ),
+ FileIndexResult::Remain
+ );
+
+ let mut invalid_timestamp =
+ BsiFileIndexWriter::try_new(timestamp_type, &Options::new()).unwrap();
+ assert!(matches!(
+ invalid_timestamp.write(Some(&Datum::Timestamp {
+ millis: 1000,
+ nanos: 1_000_000,
+ })),
+ Err(Error::DataInvalid { .. })
+ ));
+ assert!(invalid_timestamp.empty());
+}
+
+#[test]
+fn test_bsi_rejects_bad_headers_without_pruning() {
+ let mut writer = BsiFileIndexWriter::try_new(int_type(),
&Options::new()).unwrap();
+ for value in [-2, 0, 3] {
+ writer.write(Some(&Datum::Int(value))).unwrap();
+ }
+ let bytes = writer.serialized_bytes().unwrap();
+ for length in 0..bytes.len() {
+ assert!(
+ BsiFileIndexReader::try_new(int_type(),
bytes.slice(..length)).is_err(),
+ "truncated length {length}"
+ );
+ }
+ let mut bad_version = bytes.to_vec();
+ bad_version[0] = 2;
+ assert!(BsiFileIndexReader::try_new(int_type(),
Bytes::from(bad_version)).is_err());
+ let mut bad_row_count = bytes.to_vec();
+ bad_row_count[1..5].copy_from_slice(&(-1_i32).to_be_bytes());
+ assert!(BsiFileIndexReader::try_new(int_type(),
Bytes::from(bad_row_count)).is_err());
+ let mut trailing = bytes.to_vec();
+ trailing.push(1);
+ assert!(BsiFileIndexReader::try_new(int_type(),
Bytes::from(trailing)).is_err());
+}
+
+#[test]
+fn test_bsi_java_v1_fixture_round_trip() {
+ // BitSliceIndexBitmapFileIndex.Writer in Java Paimon, for
+ // [null, -5, -2, 0, 3, 8, null] as INT. Keep a genuine Java payload here
+ // so both the outer framing and Roaring bitmap encoding are checked.
+ let java_bytes = base64::engine::general_purpose::STANDARD
+ .decode(concat!(
+
"AQAAAAcBAQAAAAAAAAAAAAAAAAAAAAg6MAAAAQAAAAAAAgAQAAAAAwAEAAUAAAAABDowAAABAAAAAAAAABAAAAAEADowAAAB",
+
"AAAAAAAAABAAAAAEADowAAAAAAAAOjAAAAEAAAAAAAAAEAAAAAUAAQEAAAAAAAAAAAAAAAAAAAAFOjAAAAEAAAAAAAEAEAAAAAEAAgAAAAADOjAAAAEAAAAAAAAAEAAAAAEAOjAAAAEAAAAAAAAAEAAAAAIAOjAAAAEAAAAAAAAAEAAAAAEA"
+ ))
+ .unwrap();
+ let reader = BsiFileIndexReader::try_new(int_type(),
Bytes::from(java_bytes.clone())).unwrap();
+ assert_eq!(eval(&reader, PredicateOperator::Eq, &[-5]), selection([1]));
+ assert_eq!(
+ eval(&reader, PredicateOperator::Gt, &[0]),
+ selection([4, 5])
+ );
+ assert_eq!(
+ eval(&reader, PredicateOperator::IsNull, &[]),
+ selection([0, 6])
+ );
+ let mut writer = BsiFileIndexWriter::try_new(int_type(),
&Options::new()).unwrap();
+ for value in [None, Some(-5), Some(-2), Some(0), Some(3), Some(8), None] {
+ let datum = value.map(Datum::Int);
+ writer.write(datum.as_ref()).unwrap();
+ }
+ assert_eq!(writer.serialized_bytes().unwrap().as_ref(), java_bytes);
+}
+
+#[test]
+fn test_bsi_all_java_numeric_date_time_value_mappers() {
+ let cases: Vec<(DataType, Datum, Datum)> = vec![
+ (
+ DataType::TinyInt(TinyIntType::new()),
+ Datum::TinyInt(-5),
+ Datum::TinyInt(3),
+ ),
+ (
+ DataType::SmallInt(SmallIntType::new()),
+ Datum::SmallInt(-500),
+ Datum::SmallInt(300),
+ ),
+ (int_type(), Datum::Int(-50_000), Datum::Int(30_000)),
+ (
+ DataType::BigInt(BigIntType::new()),
+ Datum::Long(-5_000_000_000),
+ Datum::Long(3_000_000_000),
+ ),
+ (
+ DataType::Date(DateType::new()),
+ Datum::Date(-100),
+ Datum::Date(20_000),
+ ),
+ (
+ DataType::Time(TimeType::new(3).unwrap()),
+ Datum::Time(1),
+ Datum::Time(86_000_000),
+ ),
+ (
+ DataType::Decimal(DecimalType::new(38, 2).unwrap()),
+ Datum::Decimal {
+ unscaled: -500,
+ precision: 38,
+ scale: 2,
+ },
+ Datum::Decimal {
+ unscaled: 300,
+ precision: 38,
+ scale: 2,
+ },
+ ),
+ (
+ DataType::Timestamp(TimestampType::new(6).unwrap()),
+ Datum::Timestamp {
+ millis: -1000,
+ nanos: 0,
+ },
+ Datum::Timestamp {
+ millis: 1000,
+ nanos: 123_000,
+ },
+ ),
+ (
+
DataType::LocalZonedTimestamp(LocalZonedTimestampType::new(3).unwrap()),
+ Datum::LocalZonedTimestamp {
+ millis: -1000,
+ nanos: 0,
+ },
+ Datum::LocalZonedTimestamp {
+ millis: 1000,
+ nanos: 0,
+ },
+ ),
+ ];
+ for (data_type, negative, positive) in cases {
+ let mut writer = BsiFileIndexWriter::try_new(data_type.clone(),
&Options::new()).unwrap();
+ writer.write(Some(&negative)).unwrap();
+ writer.write(None).unwrap();
+ writer.write(Some(&positive)).unwrap();
+ let bytes = writer.serialized_bytes().unwrap();
+ let reader = BsiFileIndexReader::try_new(data_type.clone(),
bytes).unwrap();
+ let compare =
+ |operator, literal: Datum| reader.evaluate("v", 0, &data_type,
operator, &[literal]);
+ assert_eq!(
+ compare(PredicateOperator::Eq, negative.clone()),
+ selection([0])
+ );
+ assert_eq!(
+ compare(PredicateOperator::Eq, positive.clone()),
+ selection([2])
+ );
+ assert_eq!(
+ compare(PredicateOperator::Lt, positive.clone()),
+ selection([0])
+ );
+ assert_eq!(
+ compare(PredicateOperator::Gt, negative.clone()),
+ selection([2])
+ );
+ assert_eq!(
+ reader.evaluate("v", 0, &data_type, PredicateOperator::IsNull,
&[]),
+ selection([1])
+ );
+ assert_eq!(
+ reader.evaluate("v", 0, &data_type, PredicateOperator::NotIn,
&[negative]),
+ selection([2])
+ );
+ }
+}
+
+#[test]
+fn test_bsi_reader_rejects_inconsistent_bitmaps() {
+ let mut writer = BsiFileIndexWriter::try_new(int_type(),
&Options::new()).unwrap();
+ writer.write(Some(&Datum::Int(3))).unwrap();
+ writer.write(Some(&Datum::Int(-2))).unwrap();
+ let good = writer.serialized_bytes().unwrap();
+ let mut bad_flag = good.to_vec();
+ bad_flag[5] = 9;
+ assert!(BsiFileIndexReader::try_new(int_type(),
Bytes::from(bad_flag)).is_err());
+
+ let mut bad_max = good.to_vec();
+ // Outer header (version, row count, positive flag), then slice version,
+ // min (8 bytes), max (8 bytes). Set max below the encoded slice width.
+ bad_max[15..23].copy_from_slice(&0_i64.to_be_bytes());
+ assert!(BsiFileIndexReader::try_new(int_type(),
Bytes::from(bad_max)).is_err());
+
+ let mut bad_min = good.to_vec();
+ bad_min[7..15].copy_from_slice(&(-1_i64).to_be_bytes());
+ assert!(BsiFileIndexReader::try_new(int_type(),
Bytes::from(bad_min)).is_err());
+
+ let mut bad_count = good.to_vec();
+ bad_count[1..5].copy_from_slice(&1_i32.to_be_bytes());
+ assert!(BsiFileIndexReader::try_new(int_type(),
Bytes::from(bad_count)).is_err());
+}
+
+#[test]
+fn test_bsi_bigint_extremes_and_min_literal() {
+ let data_type = DataType::BigInt(BigIntType::new());
+ let mut writer = BsiFileIndexWriter::try_new(data_type.clone(),
&Options::new()).unwrap();
+ for value in [
+ Some(-i64::MAX),
+ Some(-1),
+ None,
+ Some(0),
+ Some(1),
+ Some(i64::MAX),
+ ] {
+ let datum = value.map(Datum::Long);
+ writer.write(datum.as_ref()).unwrap();
+ }
+ let reader =
+ BsiFileIndexReader::try_new(data_type.clone(),
writer.serialized_bytes().unwrap()).unwrap();
+ let evaluate = |operator, values: Vec<i64>| {
+ reader.evaluate(
+ "v",
+ 0,
+ &data_type,
+ operator,
+ &values.into_iter().map(Datum::Long).collect::<Vec<_>>(),
+ )
+ };
+ assert_eq!(
+ evaluate(PredicateOperator::Eq, vec![-i64::MAX]),
+ selection([0])
+ );
+ assert_eq!(
+ evaluate(PredicateOperator::Eq, vec![i64::MAX]),
+ selection([5])
+ );
+ assert_eq!(
+ evaluate(PredicateOperator::Eq, vec![i64::MIN]),
+ selection([])
+ );
+ assert_eq!(
+ evaluate(PredicateOperator::NotEq, vec![i64::MIN]),
+ selection([0, 1, 3, 4, 5])
+ );
+ assert_eq!(
+ evaluate(PredicateOperator::LtEq, vec![i64::MIN]),
+ selection([])
+ );
+ assert_eq!(
+ evaluate(PredicateOperator::Gt, vec![i64::MIN]),
+ selection([0, 1, 3, 4, 5])
+ );
+ assert_eq!(evaluate(PredicateOperator::In, vec![]), selection([]));
+ assert_eq!(
+ evaluate(PredicateOperator::NotIn, vec![]),
+ selection([0, 1, 3, 4, 5])
+ );
+ assert_eq!(
+ evaluate(PredicateOperator::Between, vec![-1, 1]),
+ selection([1, 3, 4])
+ );
+ assert_eq!(
+ reader.evaluate("v", 0, &data_type, PredicateOperator::Eq,
&[Datum::Int(1)],),
+ FileIndexResult::Remain,
+ );
+}
diff --git a/crates/paimon/src/file_index/file_indexer_factory.rs
b/crates/paimon/src/file_index/file_indexer_factory.rs
index 298225c3..11edcf79 100644
--- a/crates/paimon/src/file_index/file_indexer_factory.rs
+++ b/crates/paimon/src/file_index/file_indexer_factory.rs
@@ -21,6 +21,7 @@ use crate::common::Options;
use crate::file_index::bitmap::writer::BitmapFileIndexWriter;
use crate::file_index::bitmap::BitmapFileIndexReader;
use crate::file_index::bloom_filter::{BloomFilterReader, BloomFilterWriter};
+use crate::file_index::bsi::{BsiFileIndexReader, BsiFileIndexWriter};
use crate::file_index::file_index_reader::FileIndexReader;
use crate::file_index::file_index_writer::FileIndexWriter;
use crate::file_index::range_bitmap::writer::RangeBitmapFileIndexWriter;
@@ -31,6 +32,7 @@ use crate::{Error, Result};
pub(crate) const BITMAP_INDEX: &str = "bitmap";
pub(crate) const BLOOM_FILTER_INDEX: &str = "bloom-filter";
pub(crate) const RANGE_BITMAP_INDEX: &str = "range-bitmap";
+pub(crate) const BSI_INDEX: &str = "bsi";
struct FailOpenFileIndexReader;
@@ -41,6 +43,7 @@ enum BuiltinFileIndexer {
Bitmap,
BloomFilter,
RangeBitmap,
+ Bsi,
}
impl BuiltinFileIndexer {
@@ -49,6 +52,7 @@ impl BuiltinFileIndexer {
BITMAP_INDEX => Ok(Self::Bitmap),
BLOOM_FILTER_INDEX => Ok(Self::BloomFilter),
RANGE_BITMAP_INDEX => Ok(Self::RangeBitmap),
+ BSI_INDEX => Ok(Self::Bsi),
_ => Err(Error::Unsupported {
message: format!("Unknown file index identifier:
{identifier}"),
}),
@@ -64,7 +68,7 @@ impl FileIndexerFactory {
pub(crate) fn is_supported(identifier: &str) -> bool {
matches!(
identifier,
- BITMAP_INDEX | BLOOM_FILTER_INDEX | RANGE_BITMAP_INDEX
+ BITMAP_INDEX | BLOOM_FILTER_INDEX | RANGE_BITMAP_INDEX | BSI_INDEX
)
}
@@ -72,7 +76,7 @@ impl FileIndexerFactory {
pub(crate) fn is_write_supported(identifier: &str) -> bool {
matches!(
identifier,
- BITMAP_INDEX | BLOOM_FILTER_INDEX | RANGE_BITMAP_INDEX
+ BITMAP_INDEX | BLOOM_FILTER_INDEX | RANGE_BITMAP_INDEX | BSI_INDEX
)
}
@@ -91,6 +95,9 @@ impl FileIndexerFactory {
BuiltinFileIndexer::RangeBitmap =>
Ok(Box::new(RangeBitmapFileIndexWriter::try_new(
data_type, options,
)?)),
+ BuiltinFileIndexer::Bsi => {
+ Ok(Box::new(BsiFileIndexWriter::try_new(data_type, options)?))
+ }
}
}
@@ -119,6 +126,12 @@ impl FileIndexerFactory {
},
)
}
+ BuiltinFileIndexer::Bsi => {
+ Ok(match BsiFileIndexReader::try_new(data_type, serialized) {
+ Ok(reader) => Box::new(reader),
+ Err(_) => Box::new(FailOpenFileIndexReader),
+ })
+ }
}
}
}
@@ -137,7 +150,12 @@ mod tests {
#[test]
fn test_builtin_writers_track_empty_rows_consistently() {
- for identifier in [BITMAP_INDEX, BLOOM_FILTER_INDEX,
RANGE_BITMAP_INDEX] {
+ for identifier in [
+ BITMAP_INDEX,
+ BLOOM_FILTER_INDEX,
+ RANGE_BITMAP_INDEX,
+ BSI_INDEX,
+ ] {
assert!(FileIndexerFactory::is_write_supported(identifier));
let mut writer =
FileIndexerFactory::create_writer(identifier, int_type(),
&Options::new()).unwrap();
@@ -203,7 +221,12 @@ mod tests {
#[test]
fn test_writer_rejects_mismatched_datum() {
- for identifier in [BITMAP_INDEX, BLOOM_FILTER_INDEX,
RANGE_BITMAP_INDEX] {
+ for identifier in [
+ BITMAP_INDEX,
+ BLOOM_FILTER_INDEX,
+ RANGE_BITMAP_INDEX,
+ BSI_INDEX,
+ ] {
let mut writer =
FileIndexerFactory::create_writer(identifier, int_type(),
&Options::new()).unwrap();
diff --git a/crates/paimon/src/file_index/mod.rs
b/crates/paimon/src/file_index/mod.rs
index 3e5b5132..b9dd9f36 100644
--- a/crates/paimon/src/file_index/mod.rs
+++ b/crates/paimon/src/file_index/mod.rs
@@ -21,6 +21,8 @@
pub(crate) mod bitmap;
#[allow(dead_code)]
pub(crate) mod bloom_filter;
+#[allow(dead_code)]
+pub(crate) mod bsi;
pub(crate) mod evaluator;
mod file_index_format;
#[allow(dead_code)]
diff --git a/crates/paimon/src/table/data_evolution_reader.rs
b/crates/paimon/src/table/data_evolution_reader.rs
index d9e4ad6f..a8da83c5 100644
--- a/crates/paimon/src/table/data_evolution_reader.rs
+++ b/crates/paimon/src/table/data_evolution_reader.rs
@@ -19,17 +19,20 @@ mod blob_fallback;
use super::blob_resolver::BlobReadLimiter;
use super::data_file_reader::{
- append_null_row_id_column, attach_row_id, expand_selected_row_ids,
insert_column_at,
- DataFileReadTiming, DataFileReader,
+ append_null_row_id_column, attach_row_id, expand_selected_row_ids,
+ file_index_selection_to_local_ranges, insert_column_at,
DataFileReadTiming, DataFileReader,
};
use crate::arrow::format::blob::DEFAULT_BLOB_READ_PARALLELISM;
use crate::arrow::format::FilePredicates;
use crate::arrow::format::MosaicPrefetchOptions;
use crate::arrow::{build_target_arrow_schema, ReadBudget};
use crate::deletion_vector::{DeletionVector, DeletionVectorFactory};
+use crate::file_index::evaluator::evaluate_file_index;
+use crate::file_index::file_index_result::FileIndexResult;
use crate::io::FileIO;
use crate::spec::{
- BlobDescriptor, BlobViewStruct, DataField, DataFileMeta, DataType,
Predicate, ROW_ID_FIELD_NAME,
+ BlobDescriptor, BlobViewStruct, CoreOptions, DataField, DataFileMeta,
DataType, Predicate,
+ ROW_ID_FIELD_NAME,
};
use crate::table::dedicated_format_file_writer::is_blob_file_name;
use crate::table::schema_manager::SchemaManager;
@@ -470,6 +473,8 @@ impl DataEvolutionReader {
let push_down_raw_predicates = !self.predicates.is_empty()
&& self.row_id_index.is_none()
&& filter_before_blob_resolution;
+ let read_raw_file_index = push_down_raw_predicates
+ &&
CoreOptions::new(&self.table_options).file_index_read_enabled();
// A managed BLOB file fetches payloads as its batch is decoded.
// Keep batches small and restrict predicate-free file selections
// to the remaining output quota below.
@@ -529,6 +534,43 @@ impl DataEvolutionReader {
let has_row_id = file_meta.first_row_id.is_some();
let mut effective_row_ranges = if has_row_id {
row_ranges.clone() } else { None };
+ if read_raw_file_index {
+ // Independent files can apply a physical row
selection
+ // before decoding. Column-merge groups
cannot: an older
+ // file may hold values still needed by a
newer file.
+ match evaluate_file_index(
+ &self.file_io,
+ split.bucket_path(),
+ &file_meta,
+ &self.table_fields,
+
data_fields.as_deref().unwrap_or(&self.table_fields),
+ &self.predicates,
+ ).await? {
+ FileIndexResult::Skip => continue,
+ FileIndexResult::Selection(selection) => {
+ if let Some(local_ranges) =
+
file_index_selection_to_local_ranges(
+ &selection,
+ file_meta.row_count,
+ )? {
+ let base =
file_meta.first_row_id.unwrap_or(0);
+ let index_ranges = local_ranges
+ .into_iter()
+ .map(|range| RowRange::new(
+ base + range.from(),
+ base + range.to(),
+ ))
+ .collect::<Vec<_>>();
+ effective_row_ranges = Some(match
effective_row_ranges {
+ Some(ref ranges) =>
+
intersect_local_ranges(ranges, &index_ranges),
+ None => index_ranges,
+ });
+ }
+ }
+ FileIndexResult::Remain => {}
+ }
+ }
if self.predicates.is_empty() {
if let Some(left) = remaining {
let selected =
selected_absolute_row_ranges_for_file(
diff --git a/crates/paimon/src/table/data_evolution_writer.rs
b/crates/paimon/src/table/data_evolution_writer.rs
index 9373dcae..98f60f41 100644
--- a/crates/paimon/src/table/data_evolution_writer.rs
+++ b/crates/paimon/src/table/data_evolution_writer.rs
@@ -33,6 +33,7 @@ use crate::spec::{
FileKind, IndexFileMeta, IndexManifest, PartitionComputer, Snapshot,
EMPTY_BINARY_ROW,
};
use crate::table::commit_message::CommitMessage;
+use crate::table::data_file_index_writer::FileIndexOptions;
use crate::table::data_file_writer::DataFileWriter;
use crate::table::index_file_path::IndexFileLocation;
use crate::table::source::data_evolution_anchor_file;
@@ -1000,6 +1001,7 @@ struct PartialWriteSet {
write_columns: Vec<String>,
column_indices: Vec<usize>,
schema: Arc<arrow_schema::Schema>,
+ file_index_options: Option<Arc<FileIndexOptions>>,
}
impl DataEvolutionPartialWriter {
@@ -1019,7 +1021,20 @@ impl DataEvolutionPartialWriter {
let partition_keys: Vec<String> = schema.partition_keys().to_vec();
let fields = schema.fields();
- let write_sets = Self::build_write_sets(&write_columns, fields,
&core_options)?;
+ let mut write_sets = Self::build_write_sets(&write_columns, fields,
&core_options)?;
+ let file_index_options = FileIndexOptions::parse(schema.options(),
fields)?;
+ for write_set in &mut write_sets {
+ write_set.file_index_options = file_index_options
+ .as_ref()
+ .and_then(|options|
options.project_to_fields(&write_set.write_fields))
+ .map(Arc::new);
+ if write_set.kind == PartialFileKind::Vector &&
write_set.file_index_options.is_some() {
+ return Err(crate::Error::Unsupported {
+ message: "FileIndex generation does not support dedicated
VECTOR partial files"
+ .to_string(),
+ });
+ }
+ }
let partition_computer = PartitionComputer::new(
&partition_keys,
fields,
@@ -1088,6 +1103,7 @@ impl DataEvolutionPartialWriter {
write_columns: normal_columns,
column_indices: normal_indices,
schema,
+ file_index_options: None,
});
}
@@ -1102,6 +1118,7 @@ impl DataEvolutionPartialWriter {
write_columns: vector_columns,
column_indices: vector_indices,
schema,
+ file_index_options: None,
});
}
@@ -1161,7 +1178,8 @@ impl DataEvolutionPartialWriter {
Some(0), // file_source: APPEND
Some(first_row_id),
Some(write_set.write_columns.clone()),
- );
+ )
+ .with_file_index(write_set.file_index_options.clone());
self.writers.insert(key.clone(), writer);
}
diff --git a/crates/paimon/src/table/data_file_index_writer.rs
b/crates/paimon/src/table/data_file_index_writer.rs
index 532d5a77..440764f7 100644
--- a/crates/paimon/src/table/data_file_index_writer.rs
+++ b/crates/paimon/src/table/data_file_index_writer.rs
@@ -34,7 +34,7 @@ struct IndexColumnOptions {
indexes: BTreeMap<String, Options>,
}
-/// Validated top-level column indexes, enabled explicitly by ordinary append
writes.
+/// Validated top-level column indexes for the logical table schema.
#[derive(Clone)]
pub(super) struct FileIndexOptions {
columns: Vec<IndexColumnOptions>,
@@ -176,6 +176,30 @@ impl FileIndexOptions {
.collect::<Result<Vec<_>>>()?;
Ok(DataFileIndexWriter { columns })
}
+
+ /// Keep indexes for fields physically present in a partial-column file.
+ /// A partial file may have a different column order from the table schema;
+ /// its index positions must refer to that file's RecordBatch layout.
+ pub(super) fn project_to_fields(&self, fields: &[DataField]) ->
Option<Self> {
+ let columns = self
+ .columns
+ .iter()
+ .filter_map(|column| {
+ fields
+ .iter()
+ .position(|field| field.id() == column.field.id())
+ .map(|position| IndexColumnOptions {
+ field: column.field.clone(),
+ position,
+ indexes: column.indexes.clone(),
+ })
+ })
+ .collect::<Vec<_>>();
+ (!columns.is_empty()).then_some(Self {
+ columns,
+ in_manifest_threshold: self.in_manifest_threshold,
+ })
+ }
}
struct IndexColumn {
diff --git a/crates/paimon/src/table/data_file_index_writer/tests.rs
b/crates/paimon/src/table/data_file_index_writer/tests.rs
index 15615aa9..8bc6cbdd 100644
--- a/crates/paimon/src/table/data_file_index_writer/tests.rs
+++ b/crates/paimon/src/table/data_file_index_writer/tests.rs
@@ -19,7 +19,7 @@ use super::*;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;
-use arrow_array::Int32Array;
+use arrow_array::{Int32Array, StringArray};
use arrow_schema::{DataType as ArrowType, Field, Schema as ArrowSchema};
use futures::TryStreamExt;
use opendal::{services::MemoryConfig, Operator};
@@ -31,6 +31,7 @@ use crate::io::{FileIO, FileIOBuilder, FileIOProvider};
use crate::spec::{
BooleanType, DataFileMeta, Datum, IntType, Predicate, PredicateBuilder,
Schema, TableSchema,
};
+use crate::table::data_evolution_writer::DataEvolutionPartialWriter;
use crate::table::{DataSplitBuilder, SchemaManager, Table};
fn schema(options: &[(&str, &str)]) -> Schema {
@@ -147,7 +148,14 @@ async fn evaluate(
#[tokio::test]
async fn test_file_index_append_commit_reload_and_rolling() {
- for identifier in ["bitmap", "bloom-filter", "both", "range-bitmap",
"all"] {
+ for identifier in [
+ "bitmap",
+ "bloom-filter",
+ "both",
+ "range-bitmap",
+ "bsi",
+ "all",
+ ] {
for rolling in [false, true] {
for threshold in ["0 B", "1 MB"] {
let mut options = vec![
@@ -168,6 +176,9 @@ async fn test_file_index_append_commit_reload_and_rolling()
{
options.push(("file-index.range-bitmap.columns", "id,
value"));
options.push(("file-index.range-bitmap.id.chunk-size",
"0b"));
}
+ if matches!(identifier, "bsi" | "all") {
+ options.push(("file-index.bsi.columns", "id, value"));
+ }
let table = table(memory_io(), schema(&options)).await;
let builder = table.new_write_builder();
let mut writer = builder.new_write().unwrap();
@@ -305,8 +316,6 @@ async fn
test_file_index_skips_unsupported_identifier_groups() {
vec!["bitmap", "bloom-filter"],
] {
let mut options = vec![
- ("file-index.bsi.columns", "id"),
- ("file-index.bsi.id.version", "upstream-specific"),
("file-index.future-index.columns", "missing[nested]"),
(
"file-index.future-index.missing[nested].version",
@@ -371,10 +380,7 @@ fn test_file_index_range_bitmap_enables_generation() {
#[test]
fn test_file_index_skips_unsupported_options_without_columns() {
- let schema = schema(&[
- ("file-index.bsi.id.version", "upstream-specific"),
- ("file-index.future-index.version", "upstream-specific"),
- ]);
+ let schema = schema(&[("file-index.future-index.version",
"upstream-specific")]);
assert!(FileIndexOptions::parse(schema.options(), schema.fields())
.unwrap()
.is_none());
@@ -476,8 +482,8 @@ async fn
test_file_index_invalid_configuration_fails_before_writing() {
vec![("file-index.in-manifest-threshold", "9223372036854775807 TB")],
];
for mut options in cases {
- options.push(("file-index.bsi.columns", "id"));
- options.push(("file-index.bsi.id.version", "upstream-specific"));
+ options.push(("file-index.future-index.columns", "id"));
+ options.push(("file-index.future-index.id.version",
"upstream-specific"));
let table = table(memory_io(), schema(&options)).await;
assert!(
table.new_write_builder().new_write().is_err(),
@@ -501,39 +507,401 @@ async fn
test_file_index_invalid_configuration_fails_before_writing() {
}
#[tokio::test]
-async fn test_file_index_rejects_unsupported_table_write_modes() {
+async fn test_file_index_data_evolution_base_file() {
for identifier in ["bitmap", "range-bitmap"] {
- for schema in [
- Schema::builder()
- .column("id", crate::spec::DataType::Int(IntType::new()))
- .primary_key(["id"])
- .option("bucket", "1")
- .option(format!("file-index.{identifier}.columns"), "id")
- .build()
- .unwrap(),
- Schema::builder()
- .column("id", crate::spec::DataType::Int(IntType::new()))
- .option("data-evolution.enabled", "true")
- .option("row-tracking.enabled", "true")
- .option(format!("file-index.{identifier}.columns"), "id")
- .build()
- .unwrap(),
- ] {
- let table = table(memory_io(), schema).await;
- let error = match table.new_write_builder().new_write() {
- Ok(_) => panic!("unsupported write mode must reject index
generation"),
- Err(error) => error,
- };
- assert!(
- error
- .to_string()
- .contains("FileIndex generation supports ordinary append
writes only"),
- "{error}"
- );
+ let schema = Schema::builder()
+ .column("id", crate::spec::DataType::Int(IntType::new()))
+ .column("value", crate::spec::DataType::Int(IntType::new()))
+ .option("data-evolution.enabled", "true")
+ .option("row-tracking.enabled", "true")
+ .option(format!("file-index.{identifier}.columns"), "id")
+ .build()
+ .unwrap();
+ let table = table(memory_io(), schema).await;
+ let builder = table.new_write_builder();
+ let mut writer = builder.new_write().unwrap();
+ writer
+ .write_arrow_batch(&pk_batch(&[3, 1, 2], &[30, 10, 20]))
+ .await
+ .unwrap();
+ let messages = writer.prepare_commit().await.unwrap();
+ assert_eq!(messages[0].new_files.len(), 1);
+ assert!(messages[0].new_files[0].embedded_index.is_some());
+ builder.new_commit().commit(messages).await.unwrap();
+ let plan = table.new_read_builder().new_scan().plan().await.unwrap();
+ let file = &plan.splits()[0].data_files()[0];
+ assert_eq!(file.first_row_id, Some(0));
+ let predicate = PredicateBuilder::new(table.schema().fields())
+ .equal("id", Datum::Int(1))
+ .unwrap();
+ assert_eq!(
+ evaluate(
+ &table,
+ plan.splits()[0].bucket_path(),
+ file,
+ predicate.clone()
+ )
+ .await,
+ FileIndexResult::Selection([1].into_iter().collect())
+ );
+ assert_eq!(
+ query(&table, true, Some(predicate)).await,
+ vec![(Some(1), Some(10))]
+ );
+ }
+}
+
+#[tokio::test]
+async fn
test_file_index_data_evolution_partial_file_projects_column_positions() {
+ for threshold in ["0 B", "1 MB"] {
+ let schema = Schema::builder()
+ .column("id", crate::spec::DataType::Int(IntType::new()))
+ .column("value", crate::spec::DataType::Int(IntType::new()))
+ .option("data-evolution.enabled", "true")
+ .option("row-tracking.enabled", "true")
+ .option("file-index.bitmap.columns", "id,value")
+ .option("file-index.in-manifest-threshold", threshold)
+ .build()
+ .unwrap();
+ let table = table(memory_io(), schema).await;
+ let builder = table.new_write_builder();
+ let mut base_writer = builder.new_write().unwrap();
+ base_writer
+ .write_arrow_batch(&pk_batch(&[1, 2, 3], &[10, 20, 30]))
+ .await
+ .unwrap();
+ builder
+ .new_commit()
+ .commit(base_writer.prepare_commit().await.unwrap())
+ .await
+ .unwrap();
+
+ let mut partial_writer =
+ DataEvolutionPartialWriter::new(&table,
vec!["value".to_string()]).unwrap();
+ let partial_batch = RecordBatch::try_new(
+ Arc::new(ArrowSchema::new(vec![Field::new(
+ "value",
+ ArrowType::Int32,
+ true,
+ )])),
+ vec![Arc::new(Int32Array::from(vec![30, 11, 20]))],
+ )
+ .unwrap();
+ partial_writer
+ .write_partial_batch(
+ crate::spec::EMPTY_BINARY_ROW.to_serialized_bytes(),
+ 0,
+ 0,
+ 1,
+ partial_batch,
+ )
+ .await
+ .unwrap();
+ let messages = partial_writer.prepare_commit().await.unwrap();
+ assert_eq!(messages.len(), 1);
+ assert_eq!(messages[0].new_files.len(), 1);
+ let file = &messages[0].new_files[0];
+ assert_eq!(
+ file.write_cols.as_deref(),
+ Some(["value".to_string()].as_slice())
+ );
+ assert_eq!(file.embedded_index.is_some(), threshold != "0 B");
+ assert_eq!(file.extra_files.len(), usize::from(threshold == "0 B"));
+ let bucket_dir = format!("{}/bucket-0", table.location());
+ let predicate_builder = PredicateBuilder::new(table.schema().fields());
+ assert_eq!(
+ evaluate(
+ &table,
+ &bucket_dir,
+ file,
+ predicate_builder.equal("value", Datum::Int(11)).unwrap(),
+ )
+ .await,
+ FileIndexResult::Selection([1].into_iter().collect())
+ );
+ assert_eq!(
+ evaluate(
+ &table,
+ &bucket_dir,
+ file,
+ predicate_builder.equal("id", Datum::Int(1)).unwrap(),
+ )
+ .await,
+ FileIndexResult::Remain
+ );
+ builder.new_commit().commit(messages).await.unwrap();
+ assert_eq!(
+ query(&table, true, None).await,
+ vec![
+ (Some(1), Some(30)),
+ (Some(2), Some(11)),
+ (Some(3), Some(20))
+ ]
+ );
+ }
+}
+
+fn primary_key_schema(options: &[(&str, &str)]) -> Schema {
+ let mut builder = Schema::builder()
+ .column("id", crate::spec::DataType::Int(IntType::new()))
+ .column("value", crate::spec::DataType::Int(IntType::new()))
+ .primary_key(["id"])
+ .option("bucket", "1")
+ .option("target-file-size", "128 MB");
+ for (key, value) in options {
+ builder = builder.option(*key, *value);
+ }
+ builder.build().unwrap()
+}
+
+async fn indexed_pk_table(options: &[(&str, &str)]) -> Table {
+ table(memory_io(), primary_key_schema(options)).await
+}
+
+fn pk_batch(ids: &[i32], values: &[i32]) -> RecordBatch {
+ assert_eq!(ids.len(), values.len());
+ batch(
+ ids.iter().copied().map(Some).collect(),
+ values.iter().copied().map(Some).collect(),
+ )
+}
+
+async fn pk_query(table: &Table, enabled: bool, predicate: Option<Predicate>)
-> Vec<(i32, i32)> {
+ query(table, enabled, predicate)
+ .await
+ .into_iter()
+ .map(|(id, value)| (id.unwrap(), value.unwrap()))
+ .collect()
+}
+
+#[tokio::test]
+async fn test_pk_file_index_follows_sorted_deduplicated_rows() {
+ for engine in ["deduplicate", "first-row"] {
+ for threshold in ["0 B", "1 MB"] {
+ for index_type in ["bitmap", "range-bitmap", "bsi"] {
+ let index_option = match index_type {
+ "bitmap" => "file-index.bitmap.columns",
+ "range-bitmap" => "file-index.range-bitmap.columns",
+ _ => "file-index.bsi.columns",
+ };
+ let table = indexed_pk_table(&[
+ ("merge-engine", engine),
+ (index_option, "id,value"),
+ ("file-index.in-manifest-threshold", threshold),
+ ])
+ .await;
+ let builder = table.new_write_builder();
+ let mut writer = builder.new_write().unwrap();
+ // Arrival order differs from file order and both keys repeat.
+ writer
+ .write_arrow_batch(&pk_batch(&[3, 1, 2], &[30, 10, 20]))
+ .await
+ .unwrap();
+ writer
+ .write_arrow_batch(&pk_batch(&[1, 3], &[11, 31]))
+ .await
+ .unwrap();
+ let mut messages = writer.prepare_commit().await.unwrap();
+ if engine == "first-row" {
+ // First-row batch scans expose materialized files, not L0.
+ for message in &mut messages {
+ for file in &mut message.new_files {
+ file.level = 1;
+ }
+ }
+ }
+ assert_eq!(messages.len(), 1);
+ assert_eq!(messages[0].new_files.len(), 1);
+ let file = &messages[0].new_files[0];
+ assert_eq!(file.row_count, 3);
+ assert_eq!(file.embedded_index.is_some(), threshold != "0 B");
+ assert_eq!(file.extra_files.len(), usize::from(threshold == "0
B"));
+ for path in file.collect_files(&format!("{}/bucket-0",
table.location())) {
+ assert!(table.file_io().exists(&path).await.unwrap(),
"{path}");
+ }
+ builder.new_commit().commit(messages).await.unwrap();
+
+ let expected = if engine == "first-row" {
+ vec![(1, 10), (2, 20), (3, 30)]
+ } else {
+ vec![(1, 11), (2, 20), (3, 31)]
+ };
+ assert_eq!(pk_query(&table, false, None).await, expected);
+ assert_eq!(pk_query(&table, true, None).await, expected);
+ let plan =
table.new_read_builder().new_scan().plan().await.unwrap();
+ let predicates =
PredicateBuilder::new(table.schema().fields());
+ let file = &plan.splits()[0].data_files()[0];
+ let id_one = evaluate(
+ &table,
+ plan.splits()[0].bucket_path(),
+ file,
+ predicates.equal("id", Datum::Int(1)).unwrap(),
+ )
+ .await;
+ assert_eq!(
+ id_one,
+ FileIndexResult::Selection([0].into_iter().collect())
+ );
+ let last_value = if engine == "first-row" { 30 } else { 31 };
+ let value_predicate = predicates.equal("value",
Datum::Int(last_value)).unwrap();
+ assert_eq!(
+ evaluate(
+ &table,
+ plan.splits()[0].bucket_path(),
+ file,
+ value_predicate.clone()
+ )
+ .await,
+ FileIndexResult::Selection([2].into_iter().collect()),
+ "engine={engine}, index={index_type},
threshold={threshold}"
+ );
+ assert_eq!(
+ pk_query(&table, true, Some(value_predicate)).await,
+ vec![(3, last_value)]
+ );
+ }
+ }
+ }
+}
+
+#[tokio::test]
+async fn test_pk_file_index_abort_removes_data_and_sidecar() {
+ for threshold in ["0 B", "1 MB"] {
+ let table = indexed_pk_table(&[
+ ("file-index.bitmap.columns", "id,value"),
+ ("file-index.in-manifest-threshold", threshold),
+ ])
+ .await;
+ let builder = table.new_write_builder();
+ let mut writer = builder.new_write().unwrap();
+ writer
+ .write_arrow_batch(&pk_batch(&[2, 1], &[20, 10]))
+ .await
+ .unwrap();
+ let messages = writer.prepare_commit().await.unwrap();
+ let file = &messages[0].new_files[0];
+ let paths = file.collect_files(&format!("{}/bucket-0",
table.location()));
+ assert_eq!(paths.len(), if threshold == "0 B" { 2 } else { 1 });
+ for path in &paths {
+ assert!(table.file_io().exists(path).await.unwrap());
+ }
+ builder.new_commit().abort(&messages).await.unwrap();
+ for path in &paths {
+ assert!(!table.file_io().exists(path).await.unwrap());
+ }
+ }
+}
+
+#[tokio::test]
+async fn test_postpone_file_index_preserves_arrival_order_and_rolls() {
+ for threshold in ["0 B", "1 MB"] {
+ for rolling in [false, true] {
+ let table = indexed_pk_table(&[
+ ("bucket", "-2"),
+ ("target-file-size", if rolling { "1 B" } else { "128 MB" }),
+ ("file-index.bitmap.columns", "id,value"),
+ ("file-index.in-manifest-threshold", threshold),
+ ])
+ .await;
+ let builder = table.new_write_builder();
+ let mut writer = builder.new_write().unwrap();
+ writer
+ .write_arrow_batch(&pk_batch(&[3, 1, 3], &[30, 10, 31]))
+ .await
+ .unwrap();
+ writer
+ .write_arrow_batch(&pk_batch(&[2, 1], &[20, 11]))
+ .await
+ .unwrap();
+ let messages = writer.prepare_commit().await.unwrap();
+ assert_eq!(messages.len(), 1);
+ assert_eq!(messages[0].bucket, crate::spec::POSTPONE_BUCKET);
+ assert_eq!(messages[0].new_files.len(), if rolling { 2 } else { 1
});
+ let predicates = PredicateBuilder::new(table.schema().fields());
+ let bucket_dir = format!("{}/bucket-postpone", table.location());
+ for (position, file) in messages[0].new_files.iter().enumerate() {
+ assert_eq!(file.embedded_index.is_some(), threshold != "0 B");
+ assert_eq!(file.extra_files.len(), usize::from(threshold == "0
B"));
+ let expected_id_3: roaring::RoaringBitmap = if position == 0 {
+ [0, 2].into_iter().collect()
+ } else {
+ roaring::RoaringBitmap::new()
+ };
+ let expected_id_1: roaring::RoaringBitmap = if rolling {
+ [1].into_iter().collect()
+ } else {
+ [1, 4].into_iter().collect()
+ };
+ let result = evaluate(
+ &table,
+ &bucket_dir,
+ file,
+ predicates.equal("id", Datum::Int(3)).unwrap(),
+ )
+ .await;
+ assert_eq!(result, FileIndexResult::Selection(expected_id_3));
+ let result = evaluate(
+ &table,
+ &bucket_dir,
+ file,
+ predicates.equal("id", Datum::Int(1)).unwrap(),
+ )
+ .await;
+ assert_eq!(result, FileIndexResult::Selection(expected_id_1));
+ for path in file.collect_files(&bucket_dir) {
+ assert!(table.file_io().exists(&path).await.unwrap(),
"{path}");
+ }
+ }
+ builder.new_commit().abort(&messages).await.unwrap();
+ for file in &messages[0].new_files {
+ for path in file.collect_files(&bucket_dir) {
+ assert!(!table.file_io().exists(&path).await.unwrap(),
"{path}");
+ }
+ }
}
}
}
+#[tokio::test]
+async fn test_postpone_file_index_omits_value_kind_column() {
+ let table = indexed_pk_table(&[
+ ("bucket", "-2"),
+ ("file-index.bitmap.columns", "id,value"),
+ ("file-index.in-manifest-threshold", "1 MB"),
+ ])
+ .await;
+ let write_batch = RecordBatch::try_new(
+ Arc::new(ArrowSchema::new(vec![
+ Field::new("id", ArrowType::Int32, true),
+ Field::new("value", ArrowType::Int32, true),
+ Field::new("_VALUE_KIND", ArrowType::Int8, false),
+ ])),
+ vec![
+ Arc::new(Int32Array::from(vec![Some(3), Some(1), Some(3)])),
+ Arc::new(Int32Array::from(vec![Some(30), Some(10), Some(31)])),
+ Arc::new(arrow_array::Int8Array::from(vec![0, 3, 2])),
+ ],
+ )
+ .unwrap();
+ let mut writer = table.new_write_builder().new_write().unwrap();
+ writer.write_arrow_batch(&write_batch).await.unwrap();
+ let messages = writer.prepare_commit().await.unwrap();
+ let file = &messages[0].new_files[0];
+ let predicate = PredicateBuilder::new(table.schema().fields())
+ .equal("value", Datum::Int(31))
+ .unwrap();
+ assert_eq!(
+ evaluate(
+ &table,
+ &format!("{}/bucket-postpone", table.location()),
+ file,
+ predicate,
+ )
+ .await,
+ FileIndexResult::Selection([2].into_iter().collect())
+ );
+}
+
#[tokio::test]
async fn test_file_index_uses_partition_bucket_file_row_order() {
let schema = Schema::builder()
@@ -742,6 +1110,450 @@ async fn
test_file_index_prunes_without_opening_data_file() {
}
}
+#[tokio::test]
+async fn test_pk_file_index_prunes_before_sort_merge_opens_data() {
+ for threshold in ["0 B", "1 MB"] {
+ for identifier in ["bitmap", "bsi"] {
+ let storage = StorageProbe::new(0);
+ let schema = primary_key_schema(&[
+ (
+ if identifier == "bitmap" {
+ "file-index.bitmap.columns"
+ } else {
+ "file-index.bsi.columns"
+ },
+ "id",
+ ),
+ ("file-index.in-manifest-threshold", threshold),
+ ]);
+ let table = table(storage.io(), schema).await;
+ let builder = table.new_write_builder();
+ let mut writer = builder.new_write().unwrap();
+ writer
+ .write_arrow_batch(&pk_batch(&[3, 1], &[30, 10]))
+ .await
+ .unwrap();
+ builder
+ .new_commit()
+ .commit(writer.prepare_commit().await.unwrap())
+ .await
+ .unwrap();
+ let predicate = PredicateBuilder::new(table.schema().fields())
+ .equal("id", Datum::Int(2))
+ .unwrap();
+ let mut read_builder = table.new_read_builder();
+ read_builder.with_filter(predicate.clone());
+ let (_, trace) =
read_builder.new_scan().plan_with_trace().await.unwrap();
+ assert_eq!(trace.final_files, 1, "stats must retain the file");
+ storage.data_accesses.store(0, Ordering::SeqCst);
+ assert!(pk_query(&table, true, Some(predicate.clone()))
+ .await
+ .is_empty());
+ assert_eq!(storage.data_accesses.load(Ordering::SeqCst), 0);
+ assert!(pk_query(&table, false, Some(predicate)).await.is_empty());
+ assert!(storage.data_accesses.load(Ordering::SeqCst) > 0);
+ }
+ }
+}
+
+#[tokio::test]
+async fn test_data_evolution_file_index_prunes_independent_file() {
+ for threshold in ["0 B", "1 MB"] {
+ let storage = StorageProbe::new(0);
+ let schema = Schema::builder()
+ .column("id", crate::spec::DataType::Int(IntType::new()))
+ .column("value", crate::spec::DataType::Int(IntType::new()))
+ .option("data-evolution.enabled", "true")
+ .option("row-tracking.enabled", "true")
+ .option("file-index.bsi.columns", "id")
+ .option("file-index.in-manifest-threshold", threshold)
+ .build()
+ .unwrap();
+ let table = table(storage.io(), schema).await;
+ let builder = table.new_write_builder();
+ let mut writer = builder.new_write().unwrap();
+ writer
+ .write_arrow_batch(&pk_batch(&[1, 3], &[10, 30]))
+ .await
+ .unwrap();
+ builder
+ .new_commit()
+ .commit(writer.prepare_commit().await.unwrap())
+ .await
+ .unwrap();
+ let predicate = PredicateBuilder::new(table.schema().fields())
+ .equal("id", Datum::Int(2))
+ .unwrap();
+ let mut read_builder = table.new_read_builder();
+ read_builder.with_filter(predicate.clone());
+ let (_, trace) =
read_builder.new_scan().plan_with_trace().await.unwrap();
+ assert_eq!(trace.final_files, 1, "stats must retain the file");
+ storage.data_accesses.store(0, Ordering::SeqCst);
+ assert!(query(&table, true, Some(predicate.clone()))
+ .await
+ .is_empty());
+ assert_eq!(storage.data_accesses.load(Ordering::SeqCst), 0);
+ assert!(query(&table, false, Some(predicate)).await.is_empty());
+ assert!(storage.data_accesses.load(Ordering::SeqCst) > 0);
+ }
+}
+
+#[tokio::test]
+async fn test_pk_file_index_keeps_old_versions_for_non_key_filter() {
+ for identifier in ["bitmap", "bsi"] {
+ let table = indexed_pk_table(&[
+ (
+ if identifier == "bitmap" {
+ "file-index.bitmap.columns"
+ } else {
+ "file-index.bsi.columns"
+ },
+ "id,value",
+ ),
+ ("file-index.in-manifest-threshold", "0 B"),
+ ])
+ .await;
+ let builder = table.new_write_builder();
+ let mut old_writer = builder.new_write().unwrap();
+ old_writer
+ .write_arrow_batch(&pk_batch(&[1, 2], &[10, 20]))
+ .await
+ .unwrap();
+ builder
+ .new_commit()
+ .commit(old_writer.prepare_commit().await.unwrap())
+ .await
+ .unwrap();
+ let mut new_writer = builder.new_write().unwrap();
+ new_writer
+ .write_arrow_batch(&pk_batch(&[1, 3], &[11, 30]))
+ .await
+ .unwrap();
+ builder
+ .new_commit()
+ .commit(new_writer.prepare_commit().await.unwrap())
+ .await
+ .unwrap();
+ let predicates = PredicateBuilder::new(table.schema().fields());
+ for (predicate, expected) in [
+ (predicates.equal("value", Datum::Int(10)).unwrap(), vec![]),
+ (
+ predicates.equal("value", Datum::Int(11)).unwrap(),
+ vec![(1, 11)],
+ ),
+ (
+ predicates.equal("id", Datum::Int(1)).unwrap(),
+ vec![(1, 11)],
+ ),
+ (
+ predicates.equal("id", Datum::Int(2)).unwrap(),
+ vec![(2, 20)],
+ ),
+ ] {
+ assert_eq!(
+ pk_query(&table, true, Some(predicate.clone())).await,
+ expected
+ );
+ assert_eq!(pk_query(&table, false, Some(predicate)).await,
expected);
+ }
+ }
+}
+
+#[tokio::test]
+async fn test_pk_file_index_excludes_input_changelog_files() {
+ for threshold in ["0 B", "1 MB"] {
+ let table = indexed_pk_table(&[
+ ("changelog-producer", "input"),
+ ("file-index.bsi.columns", "id,value"),
+ ("file-index.in-manifest-threshold", threshold),
+ ])
+ .await;
+ assert_eq!(
+ table
+ .schema()
+ .options()
+ .get("changelog-producer")
+ .map(String::as_str),
+ Some("input")
+ );
+ let builder = table.new_write_builder();
+ let mut writer = builder.new_write().unwrap();
+ writer
+ .write_arrow_batch(&pk_batch(&[2, 1, 1], &[20, 10, 11]))
+ .await
+ .unwrap();
+ let messages = writer.prepare_commit().await.unwrap();
+ assert_eq!(messages.len(), 1);
+ assert_eq!(messages[0].new_files.len(), 1);
+ assert_eq!(messages[0].new_changelog_files.len(), 1);
+ let data = &messages[0].new_files[0];
+ let changelog = &messages[0].new_changelog_files[0];
+ assert_eq!(data.row_count, 2);
+ assert_eq!(changelog.row_count, 3);
+ assert_eq!(data.embedded_index.is_some(), threshold != "0 B");
+ assert_eq!(data.extra_files.len(), usize::from(threshold == "0 B"));
+ assert!(changelog.embedded_index.is_none());
+ assert!(changelog.extra_files.is_empty());
+ let bucket_dir = format!("{}/bucket-0", table.location());
+ let predicate = PredicateBuilder::new(table.schema().fields())
+ .equal("id", Datum::Int(1))
+ .unwrap();
+ assert_eq!(
+ evaluate(&table, &bucket_dir, data, predicate.clone()).await,
+ FileIndexResult::Selection([0].into_iter().collect())
+ );
+ assert_eq!(
+ evaluate(&table, &bucket_dir, changelog, predicate).await,
+ FileIndexResult::Remain
+ );
+ builder.new_commit().commit(messages).await.unwrap();
+ assert_eq!(pk_query(&table, true, None).await, vec![(1, 11), (2, 20)]);
+ }
+}
+
+#[tokio::test]
+async fn test_pk_file_index_crosses_sorted_chunk_boundary() {
+ let table = indexed_pk_table(&[
+ ("file-index.bsi.columns", "id,value"),
+ ("file-index.in-manifest-threshold", "0 B"),
+ ])
+ .await;
+ let ids: Vec<i32> = (0..5000).rev().collect();
+ let values: Vec<i32> = ids.iter().map(|id| id * 2).collect();
+ let builder = table.new_write_builder();
+ let mut writer = builder.new_write().unwrap();
+ writer
+ .write_arrow_batch(&pk_batch(&ids, &values))
+ .await
+ .unwrap();
+ let messages = writer.prepare_commit().await.unwrap();
+ assert_eq!(messages[0].new_files.len(), 1);
+ let file = &messages[0].new_files[0];
+ assert_eq!(file.row_count, 5000);
+ let bucket_dir = format!("{}/bucket-0", table.location());
+ let predicates = PredicateBuilder::new(table.schema().fields());
+ for id in [0, 4095, 4096, 4999] {
+ assert_eq!(
+ evaluate(
+ &table,
+ &bucket_dir,
+ file,
+ predicates.equal("id", Datum::Int(id)).unwrap(),
+ )
+ .await,
+ FileIndexResult::Selection([id as u32].into_iter().collect())
+ );
+ }
+ builder.new_commit().commit(messages).await.unwrap();
+ for id in [0, 4095, 4096, 4999] {
+ let predicate = predicates.equal("id", Datum::Int(id)).unwrap();
+ assert_eq!(
+ pk_query(&table, true, Some(predicate)).await,
+ vec![(id, id * 2)]
+ );
+ }
+}
+
+#[tokio::test]
+async fn test_dynamic_bucket_indexed_commit_keeps_hash_and_changelog_indexes()
{
+ let schema = Schema::builder()
+ .column(
+ "pt",
+
crate::spec::DataType::VarChar(crate::spec::VarCharType::string_type()),
+ )
+ .column("id", crate::spec::DataType::Int(IntType::new()))
+ .column("value", crate::spec::DataType::Int(IntType::new()))
+ .partition_keys(["pt"])
+ .primary_key(["pt", "id"])
+ .option("changelog-producer", "input")
+ .option("file-index.bsi.columns", "id")
+ .option("file-index.in-manifest-threshold", "0 B")
+ .build()
+ .unwrap();
+ let table = table(memory_io(), schema).await;
+ let input = RecordBatch::try_new(
+ Arc::new(ArrowSchema::new(vec![
+ Field::new("pt", ArrowType::Utf8, true),
+ Field::new("id", ArrowType::Int32, true),
+ Field::new("value", ArrowType::Int32, true),
+ ])),
+ vec![
+ Arc::new(StringArray::from(vec!["a", "a"])),
+ Arc::new(Int32Array::from(vec![1, 2])),
+ Arc::new(Int32Array::from(vec![10, 20])),
+ ],
+ )
+ .unwrap();
+ let builder = table.new_write_builder();
+ let mut writer = builder.new_write().unwrap();
+ writer.write_arrow_batch(&input).await.unwrap();
+ let messages = writer.prepare_commit().await.unwrap();
+ assert_eq!(messages.len(), 1);
+ assert_eq!(messages[0].new_files.len(), 1);
+ assert_eq!(messages[0].new_changelog_files.len(), 1);
+ assert_eq!(messages[0].new_index_files.len(), 1);
+ assert_eq!(messages[0].new_index_files[0].index_type, "HASH");
+ assert_eq!(messages[0].new_files[0].extra_files.len(), 1);
+ assert!(messages[0].new_changelog_files[0].extra_files.is_empty());
+ builder.new_commit().commit(messages).await.unwrap();
+ let mut read_builder = table.new_read_builder();
+ read_builder.with_projection(&["id", "value"]).unwrap();
+ let plan = read_builder.new_scan().plan().await.unwrap();
+ let batches: Vec<RecordBatch> = read_builder
+ .new_read()
+ .unwrap()
+ .to_arrow(plan.splits())
+ .unwrap()
+ .try_collect()
+ .await
+ .unwrap();
+ assert_eq!(
+ rows(&batches),
+ vec![(Some(1), Some(10)), (Some(2), Some(20))]
+ );
+}
+
+#[tokio::test]
+async fn test_postpone_indexed_partition_sidecars_are_aborted() {
+ let schema = Schema::builder()
+ .column(
+ "pt",
+
crate::spec::DataType::VarChar(crate::spec::VarCharType::string_type()),
+ )
+ .column("id", crate::spec::DataType::Int(IntType::new()))
+ .column("value", crate::spec::DataType::Int(IntType::new()))
+ .partition_keys(["pt"])
+ .primary_key(["pt", "id"])
+ .option("bucket", "-2")
+ .option("file-index.bsi.columns", "id")
+ .option("file-index.in-manifest-threshold", "0 B")
+ .build()
+ .unwrap();
+ let table = table(memory_io(), schema).await;
+ let input = RecordBatch::try_new(
+ Arc::new(ArrowSchema::new(vec![
+ Field::new("pt", ArrowType::Utf8, true),
+ Field::new("id", ArrowType::Int32, true),
+ Field::new("value", ArrowType::Int32, true),
+ ])),
+ vec![
+ Arc::new(StringArray::from(vec!["a", "b", "a", "b"])),
+ Arc::new(Int32Array::from(vec![3, 2, 1, 4])),
+ Arc::new(Int32Array::from(vec![30, 20, 10, 40])),
+ ],
+ )
+ .unwrap();
+ let builder = table.new_write_builder();
+ let mut writer = builder.new_write().unwrap();
+ writer.write_arrow_batch(&input).await.unwrap();
+ let messages = writer.prepare_commit().await.unwrap();
+ assert_eq!(messages.len(), 2);
+ let computer = crate::spec::PartitionComputer::new(
+ table.schema().partition_keys(),
+ table.schema().fields(),
+ "__DEFAULT_PARTITION__",
+ false,
+ )
+ .unwrap();
+ let mut selected = 0;
+ let predicates = PredicateBuilder::new(table.schema().fields());
+ for message in &messages {
+ assert_eq!(message.bucket, crate::spec::POSTPONE_BUCKET);
+ assert_eq!(message.new_files.len(), 1);
+ let file = &message.new_files[0];
+ assert_eq!(file.extra_files.len(), 1);
+ let partition =
crate::spec::BinaryRow::from_serialized_bytes(&message.partition).unwrap();
+ let partition_path =
computer.generate_partition_path(&partition).unwrap();
+ let bucket_dir = format!("{}/{partition_path}/bucket-postpone",
table.location());
+ let result = evaluate(
+ &table,
+ &bucket_dir,
+ file,
+ predicates.equal("id", Datum::Int(3)).unwrap(),
+ )
+ .await;
+ if result.remain() {
+ assert_eq!(
+ result,
+ FileIndexResult::Selection([0].into_iter().collect())
+ );
+ selected += 1;
+ }
+ for path in file.collect_files(&bucket_dir) {
+ assert!(table.file_io().exists(&path).await.unwrap());
+ }
+ }
+ assert_eq!(selected, 1);
+ builder.new_commit().abort(&messages).await.unwrap();
+ for message in &messages {
+ let partition =
crate::spec::BinaryRow::from_serialized_bytes(&message.partition).unwrap();
+ let partition_path =
computer.generate_partition_path(&partition).unwrap();
+ let bucket_dir = format!("{}/{partition_path}/bucket-postpone",
table.location());
+ for file in &message.new_files {
+ for path in file.collect_files(&bucket_dir) {
+ assert!(!table.file_io().exists(&path).await.unwrap());
+ }
+ }
+ }
+}
+
+#[tokio::test]
+async fn test_partial_column_index_remaps_reordered_fields() {
+ let schema = Schema::builder()
+ .column("id", crate::spec::DataType::Int(IntType::new()))
+ .column("a", crate::spec::DataType::Int(IntType::new()))
+ .column("b", crate::spec::DataType::Int(IntType::new()))
+ .option("data-evolution.enabled", "true")
+ .option("row-tracking.enabled", "true")
+ .option("file-index.bitmap.columns", "a,b")
+ .build()
+ .unwrap();
+ let table = table(memory_io(), schema).await;
+ let mut writer =
+ DataEvolutionPartialWriter::new(&table, vec!["b".to_string(),
"a".to_string()]).unwrap();
+ let partial_batch = RecordBatch::try_new(
+ Arc::new(ArrowSchema::new(vec![
+ Field::new("b", ArrowType::Int32, true),
+ Field::new("a", ArrowType::Int32, true),
+ ])),
+ vec![
+ Arc::new(Int32Array::from(vec![100, 200, 300])),
+ Arc::new(Int32Array::from(vec![10, 20, 30])),
+ ],
+ )
+ .unwrap();
+ writer
+ .write_partial_batch(
+ crate::spec::EMPTY_BINARY_ROW.to_serialized_bytes(),
+ 0,
+ 0,
+ 0,
+ partial_batch,
+ )
+ .await
+ .unwrap();
+ let messages = writer.prepare_commit().await.unwrap();
+ let file = &messages[0].new_files[0];
+ assert_eq!(
+ file.write_cols.as_ref().unwrap(),
+ &vec!["b".to_string(), "a".to_string()]
+ );
+ let bucket_dir = format!("{}/bucket-0", table.location());
+ let predicates = PredicateBuilder::new(table.schema().fields());
+ for (column, value, row) in [("a", 20, 1), ("b", 300, 2)] {
+ assert_eq!(
+ evaluate(
+ &table,
+ &bucket_dir,
+ file,
+ predicates.equal(column, Datum::Int(value)).unwrap(),
+ )
+ .await,
+ FileIndexResult::Selection([row].into_iter().collect())
+ );
+ }
+}
+
#[tokio::test]
async fn test_range_bitmap_append_range_pruning() {
for format in ["parquet", "row"] {
diff --git a/crates/paimon/src/table/data_file_reader.rs
b/crates/paimon/src/table/data_file_reader.rs
index a4a18839..1b175246 100644
--- a/crates/paimon/src/table/data_file_reader.rs
+++ b/crates/paimon/src/table/data_file_reader.rs
@@ -1161,7 +1161,7 @@ const MAX_FILE_INDEX_ROW_RANGES: usize = 65_536;
/// Convert a bitmap into contiguous ranges without visiting every selected
row.
/// `None` means the bitmap is too fragmented to materialize safely and callers
/// must preserve other restrictions and rely on the residual predicate.
-fn file_index_selection_to_local_ranges(
+pub(super) fn file_index_selection_to_local_ranges(
selection: &RoaringBitmap,
row_count: i64,
) -> crate::Result<Option<Vec<RowRange>>> {
diff --git a/crates/paimon/src/table/kv_file_reader.rs
b/crates/paimon/src/table/kv_file_reader.rs
index 174b6542..2382646f 100644
--- a/crates/paimon/src/table/kv_file_reader.rs
+++ b/crates/paimon/src/table/kv_file_reader.rs
@@ -25,7 +25,7 @@
//!
//! Reference: Java Paimon `SortMergeReaderWithMinHeap`.
-use super::data_file_reader::DataFileReader;
+use super::data_file_reader::{file_index_selection_to_local_ranges,
DataFileReader};
use super::sort_merge::{
AggregateMergeFunction, DeduplicateMergeFunction, FirstRowMergeFunction,
MergeFunction,
PartialUpdateMergeFunction, SortMergeReaderBuilder,
@@ -33,6 +33,8 @@ use super::sort_merge::{
use crate::arrow::format::MosaicPrefetchOptions;
use crate::arrow::{build_target_arrow_schema, ReadBudget};
use crate::deletion_vector::DeletionVectorFactory;
+use crate::file_index::evaluator::evaluate_file_index;
+use crate::file_index::file_index_result::FileIndexResult;
use crate::io::FileIO;
use crate::spec::{
BigIntType, CoreOptions, DataField, DataFileMeta, DataType as
PaimonDataType, MergeEngine,
@@ -566,6 +568,8 @@ impl KeyValueFileReader {
let config = self.config;
let table_schema_id = config.table_schema_id;
let pushdown_predicates = self.pushdown_predicates;
+ let file_index_read_enabled =
+ CoreOptions::new(&config.table_options).file_index_read_enabled();
#[cfg(test)]
let input_batch_sizes = self.input_batch_sizes;
@@ -639,6 +643,7 @@ impl KeyValueFileReader {
let run_table_fields = config.table_fields.clone();
let run_primary_keys = config.primary_keys.clone();
let run_file_io = file_io.clone();
+ let run_index_predicates = pushdown_predicates.clone();
let deletion_files_by_split =
deletion_files_by_split.clone();
let run_stream: ArrowRecordBatchStream =
Box::pin(try_stream! {
for MergeFile { split, file: file_meta } in files {
@@ -658,6 +663,32 @@ impl KeyValueFileReader {
let key_names =
file_key_names.as_deref().unwrap_or(&run_primary_keys);
let data_schema_fields =
key_value_data_schema_fields(file_fields,
key_names)?;
+ // Only the PK-only predicate projection may
run before
+ // sort-merge. A value index can match an old
version of
+ // a key, so it must not prune merge inputs
here.
+ let row_ranges = if file_index_read_enabled {
+ match evaluate_file_index(
+ &run_file_io,
+ split.bucket_path(),
+ &file_meta,
+ &run_table_fields,
+ file_fields,
+ &run_index_predicates,
+ ).await? {
+ FileIndexResult::Skip => continue,
+ FileIndexResult::Selection(selection)
+ if split.row_ranges().is_none()
+ &&
file_meta.first_row_id.is_none() => {
+
file_index_selection_to_local_ranges(
+ &selection,
+ file_meta.row_count,
+ )?
+ }
+ _ => split.row_ranges().map(|ranges|
ranges.to_vec()),
+ }
+ } else {
+ split.row_ranges().map(|ranges|
ranges.to_vec())
+ };
let deletion_file = deletion_files_by_split
.get(&(Arc::as_ptr(&split) as usize))
.and_then(|files|
files.get(&file_meta.file_name))
@@ -674,7 +705,7 @@ impl KeyValueFileReader {
data_fields,
data_schema_fields,
deletion_vector,
- split.row_ranges().map(|ranges|
ranges.to_vec()),
+ row_ranges,
)?;
while let Some(batch) =
file_stream.next().await {
yield batch?;
diff --git a/crates/paimon/src/table/kv_file_writer.rs
b/crates/paimon/src/table/kv_file_writer.rs
index 3dc42928..9965277b 100644
--- a/crates/paimon/src/table/kv_file_writer.rs
+++ b/crates/paimon/src/table/kv_file_writer.rs
@@ -34,10 +34,11 @@ use crate::io::FileIO;
use crate::resource::{MemoryReservation, ResourceContext};
use crate::spec::stats::{compute_column_stats, BinaryTableStats};
use crate::spec::{
- bucket_path_under, extract_datum_from_arrow, AggregationConfig,
BinaryRowBuilder, CoreOptions,
- DataField, DataFileMeta, DataType, MergeEngine, PartialUpdateConfig,
RowKind,
- SEQUENCE_NUMBER_FIELD_NAME, VALUE_KIND_FIELD_NAME,
+ bucket_path_under, data_file_to_file_index_file_name,
extract_datum_from_arrow,
+ AggregationConfig, BinaryRowBuilder, CoreOptions, DataField, DataFileMeta,
DataType,
+ MergeEngine, PartialUpdateConfig, RowKind, SEQUENCE_NUMBER_FIELD_NAME,
VALUE_KIND_FIELD_NAME,
};
+use crate::table::data_file_index_writer::FileIndexOptions;
use crate::table::prepared_files::PreparedFiles;
use crate::table::sort_merge::{AggregateMergeFunction, BufferedBatch,
MergeRow};
use crate::Result;
@@ -100,6 +101,8 @@ pub(crate) struct KeyValueWriteConfig {
/// Merge engine for deduplication.
pub merge_engine: MergeEngine,
pub deletion_vectors_enabled: bool,
+ /// File indexes follow the sorted, merged data rows and never changelog
rows.
+ pub file_index_options: Option<Arc<FileIndexOptions>>,
}
struct IndexedFileWrite<'a> {
@@ -423,6 +426,15 @@ impl KeyValueFileWriter {
let last_row = indices.value(indices.len() - 1) as usize;
let min_key = self.extract_key_binary_row(batch, first_row)?;
let max_key = self.extract_key_binary_row(batch, last_row)?;
+ let mut file_index = if write.is_changelog {
+ None
+ } else {
+ self.config
+ .file_index_options
+ .as_ref()
+ .map(|options| options.create_writer())
+ .transpose()?
+ };
let physical_schema = build_physical_schema(&user_schema);
let file_name = format!(
@@ -537,6 +549,25 @@ impl KeyValueFileWriter {
let _ = self.file_io.delete_file(&file_path).await;
return Err(error);
}
+ if let Some(index) = file_index.as_mut() {
+ // The index positions must match the physical file after PK
+ // sorting and flush-time merging. Its field positions refer
+ // to the logical value schema, so omit the two KV metadata
+ // columns from the same output chunk used by the file writer.
+ let logical_indices =
(2..chunk_batch.num_columns()).collect::<Vec<_>>();
+ let index_result = chunk_batch
+ .project(&logical_indices)
+ .map_err(|error| crate::Error::DataInvalid {
+ message: format!("Failed to project KV index values:
{error}"),
+ source: None,
+ })
+ .and_then(|logical_batch| index.write(&logical_batch));
+ if let Err(error) = index_result {
+ let _ = writer.close().await;
+ let _ = self.file_io.delete_file(&file_path).await;
+ return Err(error);
+ }
+ }
}
let write_result = writer.close().await?;
@@ -580,7 +611,7 @@ impl KeyValueFileWriter {
&self.config.primary_key_types,
)?;
- Ok(DataFileMeta {
+ let mut meta = DataFileMeta {
file_name,
file_size,
row_count: indices.len() as i64,
@@ -602,7 +633,44 @@ impl KeyValueFileWriter {
first_row_id: None,
write_cols: None,
column_max_sequence_numbers: None,
- })
+ };
+ if let Some(index) = file_index {
+ let index_result = index.serialize();
+ let bytes = match index_result {
+ Ok(bytes) => bytes,
+ Err(error) => {
+ let _ = self.file_io.delete_file(&file_path).await;
+ return Err(error);
+ }
+ };
+ let threshold = self
+ .config
+ .file_index_options
+ .as_ref()
+ .expect("file index writer must have options")
+ .in_manifest_threshold;
+ if bytes.len() as u64 > threshold as u64 {
+ let name = data_file_to_file_index_file_name(&meta.file_name);
+ let index_path = format!("{bucket_dir}/{name}");
+ let output = match self.file_io.new_output(&index_path) {
+ Ok(output) => output,
+ Err(error) => {
+ let _ = self.file_io.delete_file(&file_path).await;
+ return Err(error);
+ }
+ };
+ if let Err(error) = output.write(bytes).await {
+ let _ = self.file_io.delete_file(&index_path).await;
+ let _ = self.file_io.delete_file(&file_path).await;
+ return Err(error);
+ }
+ meta.extra_files.push(name);
+ } else {
+ meta.embedded_index = Some(bytes.to_vec());
+ }
+ }
+
+ Ok(meta)
}
fn indexed_delete_row_count(batch: &RecordBatch, indices: &UInt32Array) ->
Result<i64> {
@@ -1066,6 +1134,7 @@ mod tests {
sequence_field_indices: vec![1],
merge_engine,
deletion_vectors_enabled: false,
+ file_index_options: None,
}
}
diff --git a/crates/paimon/src/table/postpone_file_writer.rs
b/crates/paimon/src/table/postpone_file_writer.rs
index ef879fa2..b932de33 100644
--- a/crates/paimon/src/table/postpone_file_writer.rs
+++ b/crates/paimon/src/table/postpone_file_writer.rs
@@ -28,7 +28,11 @@ use crate::arrow::format::{create_format_writer,
with_write_resources, FormatFil
use crate::io::FileIO;
use crate::resource::ResourceContext;
use crate::spec::stats::BinaryTableStats;
-use crate::spec::{bucket_path_under, DataFileMeta, EMPTY_SERIALIZED_ROW,
VALUE_KIND_FIELD_NAME};
+use crate::spec::{
+ bucket_path_under, data_file_to_file_index_file_name, DataFileMeta,
EMPTY_SERIALIZED_ROW,
+ VALUE_KIND_FIELD_NAME,
+};
+use crate::table::data_file_index_writer::{DataFileIndexWriter,
FileIndexOptions};
use crate::table::kv_file_writer::build_physical_schema;
use crate::Result;
use arrow_array::{Int64Array, Int8Array, RecordBatch};
@@ -49,6 +53,7 @@ pub(crate) struct PostponeWriteConfig {
pub file_format: String,
/// Data file name prefix: `"data--u-{commitUser}-s-{writeId}-w-"`.
pub data_file_prefix: String,
+ pub file_index_options: Option<Arc<FileIndexOptions>>,
}
/// Writer for postpone bucket mode (`bucket = -2`).
@@ -61,6 +66,7 @@ pub(crate) struct PostponeFileWriter {
config: PostponeWriteConfig,
next_sequence_number: i64,
current_writer: Option<Box<dyn FormatFileWriter>>,
+ current_index: Option<DataFileIndexWriter>,
current_file_name: Option<String>,
current_row_count: i64,
/// Sequence number at which the current file started.
@@ -81,6 +87,7 @@ impl PostponeFileWriter {
config,
next_sequence_number: 0,
current_writer: None,
+ current_index: None,
current_file_name: None,
current_row_count: 0,
current_file_start_seq: 0,
@@ -98,6 +105,14 @@ impl PostponeFileWriter {
}
pub(crate) async fn write(&mut self, batch: &RecordBatch) -> Result<()> {
+ let result = self.write_batch(batch).await;
+ if result.is_err() && self.config.file_index_options.is_some() {
+ self.abort().await;
+ }
+ result
+ }
+
+ async fn write_batch(&mut self, batch: &RecordBatch) -> Result<()> {
if batch.num_rows() == 0 {
return Ok(());
}
@@ -146,6 +161,19 @@ impl PostponeFileWriter {
.unwrap()
.write(&physical_batch)
.await?;
+ if let Some(index) = &mut self.current_index {
+ // FileIndex field positions refer to the user schema, while the
+ // physical file starts with two KV metadata columns. Project the
+ // exact rows written to this file in their arrival order.
+ let logical_indices =
(2..physical_batch.num_columns()).collect::<Vec<_>>();
+ let logical_batch =
physical_batch.project(&logical_indices).map_err(|error| {
+ crate::Error::DataInvalid {
+ message: format!("Failed to project postpone index values:
{error}"),
+ source: None,
+ }
+ })?;
+ index.write(&logical_batch)?;
+ }
self.next_sequence_number = end_seq + 1;
self.current_row_count += num_rows as i64;
@@ -169,6 +197,7 @@ impl PostponeFileWriter {
if let Some(writer) = self.current_writer.take() {
let _ = writer.close().await;
}
+ self.current_index = None;
self.current_file_name = None;
while self.in_flight_closes.join_next().await.is_some() {}
for path in self.created_paths.drain(..) {
@@ -178,6 +207,14 @@ impl PostponeFileWriter {
}
pub(crate) async fn prepare_commit(&mut self) -> Result<Vec<DataFileMeta>>
{
+ let result = self.finish().await;
+ if result.is_err() && self.config.file_index_options.is_some() {
+ self.abort().await;
+ }
+ result
+ }
+
+ async fn finish(&mut self) -> Result<Vec<DataFileMeta>> {
self.close_current_file().await?;
while let Some(result) = self.in_flight_closes.join_next().await {
let meta = result.map_err(|e| crate::Error::DataInvalid {
@@ -197,6 +234,18 @@ impl PostponeFileWriter {
None => return,
};
let file_name = self.current_file_name.take().unwrap();
+ let index = self.current_index.take();
+ let file_io = self.file_io.clone();
+ let bucket_dir = bucket_path_under(
+ &self.config.table_location,
+ &self.config.partition_path,
+ self.config.bucket,
+ );
+ let threshold = self
+ .config
+ .file_index_options
+ .as_ref()
+ .map(|options| options.in_manifest_threshold);
let row_count = self.current_row_count;
let min_seq = self.current_file_start_seq;
let max_seq = self.next_sequence_number - 1;
@@ -208,7 +257,7 @@ impl PostponeFileWriter {
self.in_flight_closes.spawn(async move {
let file_size = writer.close().await?.file_size as i64;
- Ok(build_meta(
+ let mut meta = build_meta(
file_name,
file_size,
row_count,
@@ -216,11 +265,19 @@ impl PostponeFileWriter {
max_seq,
schema_id,
creation_time,
- ))
+ );
+ write_index(index, threshold, &file_io, &bucket_dir, &mut
meta).await?;
+ Ok(meta)
});
}
async fn open_new_file(&mut self, user_schema: arrow_schema::SchemaRef) ->
Result<()> {
+ let index = self
+ .config
+ .file_index_options
+ .as_ref()
+ .map(|options| options.create_writer())
+ .transpose()?;
let file_name = format!(
"{}{}-{}.{}",
self.config.data_file_prefix,
@@ -237,6 +294,12 @@ impl PostponeFileWriter {
let physical_schema = build_physical_schema(&user_schema);
let file_path = format!("{bucket_dir}/{file_name}");
self.created_paths.push(file_path.clone());
+ if index.is_some() {
+ self.created_paths.push(format!(
+ "{bucket_dir}/{}",
+ data_file_to_file_index_file_name(&file_name)
+ ));
+ }
let output = self.file_io.new_output(&file_path)?;
let writer = create_format_writer(
&output,
@@ -249,6 +312,7 @@ impl PostponeFileWriter {
)
.await?;
self.current_writer = Some(with_write_resources(writer,
self.resources.as_ref()));
+ self.current_index = index;
self.current_file_name = Some(file_name);
self.current_row_count = 0;
self.current_file_start_seq = self.next_sequence_number;
@@ -262,6 +326,7 @@ impl PostponeFileWriter {
None => return Ok(()),
};
let file_name = self.current_file_name.take().unwrap();
+ let index = self.current_index.take();
let row_count = self.current_row_count;
self.current_row_count = 0;
let file_size = writer.close().await?.file_size as i64;
@@ -269,7 +334,7 @@ impl PostponeFileWriter {
let min_seq = self.current_file_start_seq;
let max_seq = self.next_sequence_number - 1;
- let meta = build_meta(
+ let mut meta = build_meta(
file_name,
file_size,
row_count,
@@ -278,11 +343,45 @@ impl PostponeFileWriter {
self.config.schema_id,
self.current_file_creation_time,
);
+ let bucket_dir = bucket_path_under(
+ &self.config.table_location,
+ &self.config.partition_path,
+ self.config.bucket,
+ );
+ let threshold = self
+ .config
+ .file_index_options
+ .as_ref()
+ .map(|options| options.in_manifest_threshold);
+ write_index(index, threshold, &self.file_io, &bucket_dir, &mut
meta).await?;
self.written_files.push(meta);
Ok(())
}
}
+async fn write_index(
+ index: Option<DataFileIndexWriter>,
+ threshold: Option<i64>,
+ file_io: &FileIO,
+ bucket_dir: &str,
+ meta: &mut DataFileMeta,
+) -> Result<()> {
+ if let Some(index) = index {
+ let bytes = index.serialize()?;
+ if bytes.len() as u64 > threshold.expect("index has threshold") as u64
{
+ let name = data_file_to_file_index_file_name(&meta.file_name);
+ file_io
+ .new_output(&format!("{bucket_dir}/{name}"))?
+ .write(bytes)
+ .await?;
+ meta.extra_files.push(name);
+ } else {
+ meta.embedded_index = Some(bytes.to_vec());
+ }
+ }
+ Ok(())
+}
+
fn build_meta(
file_name: String,
file_size: i64,
diff --git a/crates/paimon/src/table/table_write.rs
b/crates/paimon/src/table/table_write.rs
index 1c3da1b9..01567d67 100644
--- a/crates/paimon/src/table/table_write.rs
+++ b/crates/paimon/src/table/table_write.rs
@@ -398,14 +398,11 @@ impl TableWrite {
let file_index_options = FileIndexOptions::parse(schema.options(),
schema.fields())?;
if file_index_options.is_some()
- && (has_primary_keys
- || has_blob_fields
- || has_dedicated_vector_fields
- || !blob_view_fields.is_empty()
- || core_options.data_evolution_enabled())
+ && (has_blob_fields || has_dedicated_vector_fields ||
!blob_view_fields.is_empty())
{
return Err(crate::Error::Unsupported {
- message: "FileIndex generation supports ordinary append writes
only; primary-key, data-evolution and dedicated Blob/Vector writes are not
supported".to_string(),
+ message: "FileIndex generation does not support dedicated
Blob/Vector writes"
+ .to_string(),
});
}
@@ -943,9 +940,6 @@ impl TableWrite {
pub async fn prepare_commit(&mut self) -> Result<Vec<CommitMessage>> {
self.ensure_active()?;
self.partition_seq_cache.clear();
- if self.file_index_options.is_some() {
- return self.prepare_indexed_append_commit().await;
- }
let writers: Vec<(PartitionBucketKey, FileWriter)> =
self.partition_writers.drain().collect();
@@ -1009,38 +1003,6 @@ impl TableWrite {
Ok(messages)
}
- async fn prepare_indexed_append_commit(&mut self) ->
Result<Vec<CommitMessage>> {
- let closes =
- self.partition_writers
- .drain()
- .map(|((partition, bucket), writer)| async move {
- (partition, bucket, writer.prepare_commit().await)
- });
- // Do not cancel another partition's close when one fails: its
completed
- // files must remain reachable for abort cleanup.
- let results = futures::future::join_all(closes).await;
- let mut messages = Vec::new();
- let mut error = None;
- for (partition, bucket, result) in results {
- match result {
- Ok(files) if !files.data_files.is_empty() => {
- messages.push(CommitMessage::new(partition, bucket,
files.data_files));
- }
- Ok(_) => {}
- Err(err) => {
- error.get_or_insert(err);
- }
- }
- }
- if let Some(error) = error {
- self.failed = true;
- let commit = super::TableCommit::new(self.table.clone(),
self.commit_user.clone());
- let _ = commit.abort(&messages).await;
- return Err(error);
- }
- Ok(messages)
- }
-
async fn create_writer(&mut self, partition_bytes: Vec<u8>, bucket: i32)
-> Result<()> {
let partition_path = self.resolve_partition_path(&partition_bytes)?;
@@ -1142,6 +1104,7 @@ impl TableWrite {
write_buffer_size: self.write_buffer_size,
file_format: self.file_format.clone(),
data_file_prefix,
+ file_index_options: self.file_index_options.clone(),
},
)
.with_resources(self.resources.clone()),
@@ -1201,6 +1164,7 @@ impl TableWrite {
merge_engine: self.merge_engine,
deletion_vectors_enabled:
CoreOptions::new(self.table.schema().options())
.deletion_vectors_enabled(),
+ file_index_options: self.file_index_options.clone(),
},
next_seq,
)?