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, &timestamp_type, PredicateOperator::IsNull, 
&[]),
+        selection([1])
+    );
+    assert_eq!(
+        reader.evaluate(
+            "ts",
+            0,
+            &timestamp_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,
             )?

Reply via email to