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 f9f279d7 feat: support deletion vector merge-on-read (#706)
f9f279d7 is described below

commit f9f279d7b94da3d24b4cd8ca886197e425e17b19
Author: Jingsong Lee <[email protected]>
AuthorDate: Wed Aug 12 19:33:00 2026 +0800

    feat: support deletion vector merge-on-read (#706)
---
 crates/integrations/datafusion/tests/pk_tables.rs |  69 ++++-
 crates/paimon/src/spec/schema.rs                  |  41 +++
 crates/paimon/src/table/kv_file_reader.rs         | 348 +++++++++++++++++++++-
 crates/paimon/src/table/table_read.rs             |  20 +-
 crates/paimon/src/table/table_scan.rs             |  98 +++++-
 docs/src/sql.md                                   |   9 +
 6 files changed, 546 insertions(+), 39 deletions(-)

diff --git a/crates/integrations/datafusion/tests/pk_tables.rs 
b/crates/integrations/datafusion/tests/pk_tables.rs
index 0f44c819..83a24d24 100644
--- a/crates/integrations/datafusion/tests/pk_tables.rs
+++ b/crates/integrations/datafusion/tests/pk_tables.rs
@@ -2110,11 +2110,8 @@ async fn 
test_pk_partitioned_fixed_bucket_predicate_query() {
 
 // ======================= DV + Deduplicate Regression =======================
 
-/// Regression: DV-enabled Deduplicate PK table must not error on read.
-/// Before the fix, removing the DV guard caused level-0 files to reach
-/// KeyValueFileReader which rejects deletion-vector files with a hard error.
-/// With the guard restored, level-0 files are skipped in scan (DV mode relies
-/// on compaction to produce higher-level files).
+/// DV-enabled Deduplicate PK tables keep the existing compacted-only behavior
+/// unless merge-on-read is explicitly enabled.
 #[tokio::test]
 async fn test_pk_dv_deduplicate_read_no_error() {
     let (_tmp, sql_context) = setup_sql_context().await;
@@ -2146,11 +2143,8 @@ async fn test_pk_dv_deduplicate_read_no_error() {
         .await
         .unwrap();
 
-    // Read must not error. DV mode skips level-0 files, so only compacted
-    // (level > 0) files are visible. Without compaction, all files are level-0
-    // and get skipped — count may be 0, but the read must succeed without 
error.
-    // Before the fix, this would hard-fail with "KeyValueFileReader does not
-    // support deletion vectors".
+    // The default remains compacted-only: level-0 files are skipped, so the
+    // read may be empty before compaction but must remain valid.
     let result = sql_context
         .sql("SELECT * FROM paimon.test_db.t_dv_dedup")
         .await
@@ -2164,6 +2158,61 @@ async fn test_pk_dv_deduplicate_read_no_error() {
     );
 }
 
+/// Merge-on-read makes level-0 files visible and resolves overlapping primary
+/// keys before applying the residual predicate.
+#[tokio::test]
+async fn test_pk_dv_merge_on_read_reads_level_zero() {
+    let (_tmp, sql_context) = setup_sql_context().await;
+
+    sql_context
+        .sql(
+            "CREATE TABLE paimon.test_db.t_dv_mor (
+                id INT NOT NULL, value INT,
+                PRIMARY KEY (id)
+            ) WITH (
+                'bucket' = '1',
+                'deletion-vectors.enabled' = 'true',
+                'deletion-vectors.merge-on-read' = 'true',
+                'source.split.target-size' = '1b',
+                'source.split.open-file-cost' = '1b'
+            )",
+        )
+        .await
+        .unwrap();
+
+    sql_context
+        .sql("INSERT INTO paimon.test_db.t_dv_mor VALUES (1, 10), (2, 20)")
+        .await
+        .unwrap()
+        .collect()
+        .await
+        .unwrap();
+    sql_context
+        .sql("INSERT INTO paimon.test_db.t_dv_mor VALUES (2, 200), (3, 30)")
+        .await
+        .unwrap()
+        .collect()
+        .await
+        .unwrap();
+
+    let rows = collect_id_value(
+        &sql_context,
+        "SELECT id, value FROM paimon.test_db.t_dv_mor ORDER BY id",
+    )
+    .await;
+    assert_eq!(rows, vec![(1, 10), (2, 200), (3, 30)]);
+
+    let stale_rows = collect_id_value(
+        &sql_context,
+        "SELECT id, value FROM paimon.test_db.t_dv_mor WHERE value = 20",
+    )
+    .await;
+    assert!(
+        stale_rows.is_empty(),
+        "the residual predicate must run after primary-key merging"
+    );
+}
+
 // ======================= Cross-Split Merge Correctness 
=======================
 
 /// Regression: a 1-byte split target forces every data file into its own
diff --git a/crates/paimon/src/spec/schema.rs b/crates/paimon/src/spec/schema.rs
index dd45712e..c763f0ba 100644
--- a/crates/paimon/src/spec/schema.rs
+++ b/crates/paimon/src/spec/schema.rs
@@ -3048,6 +3048,47 @@ mod tests {
         }
     }
 
+    #[test]
+    fn 
test_deletion_vector_merge_on_read_is_ignored_without_deletion_vectors() {
+        let schema = Schema::builder()
+            .column("id", DataType::Int(IntType::new()))
+            .column("value", DataType::Int(IntType::new()))
+            .primary_key(["id"])
+            .option("deletion-vectors.merge-on-read", "true")
+            .build()
+            .unwrap();
+        assert_eq!(
+            schema
+                .options()
+                .get("deletion-vectors.merge-on-read")
+                .map(String::as_str),
+            Some("true")
+        );
+
+        let table_schema = TableSchema::new(
+            0,
+            &Schema::builder()
+                .column("id", DataType::Int(IntType::new()))
+                .column("value", DataType::Int(IntType::new()))
+                .primary_key(["id"])
+                .build()
+                .unwrap(),
+        );
+        let changed = table_schema
+            .apply_changes(vec![crate::spec::SchemaChange::set_option(
+                "deletion-vectors.merge-on-read".to_string(),
+                "true".to_string(),
+            )])
+            .unwrap();
+        assert_eq!(
+            changed
+                .options()
+                .get("deletion-vectors.merge-on-read")
+                .map(String::as_str),
+            Some("true")
+        );
+    }
+
     #[test]
     fn 
test_deletion_vector_schema_validation_rejects_incompatible_changelog_producers()
 {
         for (producer, expected_message) in [
diff --git a/crates/paimon/src/table/kv_file_reader.rs 
b/crates/paimon/src/table/kv_file_reader.rs
index ee0c3ecc..d922c267 100644
--- a/crates/paimon/src/table/kv_file_reader.rs
+++ b/crates/paimon/src/table/kv_file_reader.rs
@@ -31,6 +31,7 @@ use super::sort_merge::{
     SortMergeReaderBuilder,
 };
 use crate::arrow::{build_target_arrow_schema, ParquetReadBudget};
+use crate::deletion_vector::DeletionVectorFactory;
 use crate::io::FileIO;
 use crate::spec::{
     BigIntType, DataField, DataFileMeta, DataType as PaimonDataType, 
MergeEngine,
@@ -39,7 +40,7 @@ use crate::spec::{
 };
 use crate::table::schema_manager::SchemaManager;
 use crate::table::ArrowRecordBatchStream;
-use crate::{DataSplit, Error};
+use crate::{DataSplit, DeletionFile, Error};
 use arrow_array::{RecordBatch, RecordBatchOptions};
 
 use async_stream::try_stream;
@@ -529,15 +530,29 @@ impl KeyValueFileReader {
 
         Ok(try_stream! {
             for split_group in &split_groups {
-                // DV mode should not reach KeyValueFileReader.
+                // A deletion-vector merge-on-read split can mix compacted
+                // sources carrying DVs with uncompacted level-0 files. Keep
+                // only the small per-file metadata here; load each bitmap when
+                // its sorted run reaches that physical file.
+                let mut deletion_files_by_split =
+                    HashMap::<usize, Arc<HashMap<String, 
DeletionFile>>>::new();
                 for split in split_group {
-                    if split
-                        .data_deletion_files()
-                        .is_some_and(|files| files.iter().any(Option::is_some))
-                    {
-                        Err(Error::Unsupported {
-                            message: "KeyValueFileReader does not support 
deletion vectors".to_string(),
-                        })?;
+                    let Some(deletion_files) = split.data_deletion_files() 
else {
+                        continue;
+                    };
+                    let by_name = split
+                        .data_files()
+                        .iter()
+                        .zip(deletion_files.iter())
+                        .filter_map(|(data_file, deletion_file)| {
+                            deletion_file
+                                .as_ref()
+                                .map(|file| (data_file.file_name.clone(), 
file.clone()))
+                        })
+                        .collect::<HashMap<_, _>>();
+                    if !by_name.is_empty() {
+                        deletion_files_by_split
+                            .insert(Arc::as_ptr(split) as usize, 
Arc::new(by_name));
                     }
                 }
                 for merge_group in plan_merge_groups(
@@ -570,6 +585,8 @@ impl KeyValueFileReader {
                         .with_batch_size(Some(read_batch_size))
                         
.with_parquet_read_budget(group_parquet_read_budget.clone());
                         let run_schema_manager = schema_manager.clone();
+                        let run_file_io = file_io.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 {
                                 let data_fields: Option<Vec<DataField>> =
@@ -577,14 +594,24 @@ impl KeyValueFileReader {
                                         let data_schema =
                                             
run_schema_manager.schema(file_meta.schema_id).await?;
                                         Some(data_schema.fields().to_vec())
-                                    } else {
-                                        None
-                                    };
+                                } else {
+                                    None
+                                };
+                                let deletion_file = deletion_files_by_split
+                                    .get(&(Arc::as_ptr(&split) as usize))
+                                    .and_then(|files| 
files.get(&file_meta.file_name))
+                                    .cloned();
+                                let deletion_vector = match deletion_file {
+                                    Some(file) => Some(Arc::new(
+                                        
DeletionVectorFactory::read(&run_file_io, &file).await?,
+                                    )),
+                                    None => None,
+                                };
                                 let mut file_stream = 
reader.read_single_file_stream(
                                     split.as_ref(),
                                     file_meta,
                                     data_fields,
-                                    None,
+                                    deletion_vector,
                                     split.row_ranges().map(|ranges| 
ranges.to_vec()),
                                 )?;
                                 while let Some(batch) = 
file_stream.next().await {
@@ -693,20 +720,24 @@ impl KeyValueFileReader {
 mod tests {
     use super::*;
     use crate::catalog::Identifier;
+    use crate::deletion_vector::DeletionVector;
     use crate::io::FileIOBuilder;
     use crate::spec::{
         stats::BinaryTableStats, BinaryRow, DataFileMeta, DataType, Datum, 
IntType,
         PredicateBuilder, Schema, TableSchema, VarCharType,
     };
-    use crate::table::source::DataSplitBuilder;
+    use crate::table::source::{DataSplitBuilder, DeletionFile};
     use crate::table::table_commit::TableCommit;
     use crate::table::{Table, TableWrite};
     use arrow_array::{Array, Int32Array, Int64Array, Int8Array, StringArray};
     use arrow_schema::{DataType as ArrowDataType, Field as ArrowField, Schema 
as ArrowSchema};
+    use bytes::Bytes;
     use futures::TryStreamExt;
     use parquet::arrow::AsyncArrowWriter;
     use parquet::file::metadata::ParquetMetaDataReader;
     use parquet::file::properties::WriterProperties;
+    use roaring::RoaringBitmap;
+    use std::collections::HashMap;
     use std::sync::Arc;
 
     #[tokio::test]
@@ -891,6 +922,35 @@ mod tests {
             .unwrap();
     }
 
+    async fn write_deletion_file(
+        file_io: &FileIO,
+        table_path: &str,
+        deleted_rows: &[u32],
+    ) -> DeletionFile {
+        let path = format!("{table_path}/index/dv");
+        file_io
+            .mkdirs(&format!("{table_path}/index/"))
+            .await
+            .unwrap();
+        let bitmap = deleted_rows.iter().copied().collect::<RoaringBitmap>();
+        let bytes = DeletionVector::from_bitmap(bitmap)
+            .serialize_to_bytes()
+            .unwrap();
+        let bitmap_length = 
i32::from_be_bytes(bytes[0..4].try_into().unwrap());
+        file_io
+            .new_output(&path)
+            .unwrap()
+            .write(Bytes::from(bytes))
+            .await
+            .unwrap();
+        DeletionFile::new(
+            path,
+            0,
+            i64::from(bitmap_length),
+            Some(deleted_rows.len() as i64),
+        )
+    }
+
     async fn read_rows(
         table: &Table,
         projection: Option<&[&str]>,
@@ -946,6 +1006,225 @@ mod tests {
         }
     }
 
+    /// Java-compatible DV merge-on-read is a batch visibility override: the
+    /// default still hides uncompacted level-0 files, while a dynamic override
+    /// includes them and merges overlapping key versions. A tiny split target
+    /// makes this also catch planners that incorrectly separate overlapping
+    /// files into independent raw splits.
+    #[tokio::test]
+    async fn dv_merge_on_read_exposes_and_merges_level_zero_files() {
+        let file_io = test_file_io();
+        let table_path = "memory:/dv_merge_on_read_level_zero";
+        setup_dirs(&file_io, table_path).await;
+        let table = pk_table(
+            &file_io,
+            table_path,
+            &[
+                ("deletion-vectors.enabled", "true"),
+                ("source.split.target-size", "1b"),
+                ("source.split.open-file-cost", "1b"),
+            ],
+        );
+
+        write_commit(&table, &int_batch(vec![1, 2], vec![Some(10), 
Some(20)])).await;
+        write_commit(&table, &int_batch(vec![1, 3], vec![Some(11), 
Some(30)])).await;
+
+        let hidden_plan = 
table.new_read_builder().new_scan().plan().await.unwrap();
+        assert!(
+            hidden_plan.splits().is_empty(),
+            "DV batch reads must keep hiding level-0 files by default"
+        );
+
+        let merge_on_read = table.copy_with_options(HashMap::from([(
+            "deletion-vectors.merge-on-read".to_string(),
+            "true".to_string(),
+        )]));
+        let read_builder = merge_on_read.new_read_builder();
+        let plan = read_builder.new_scan().plan().await.unwrap();
+        assert_eq!(
+            plan.splits().len(),
+            1,
+            "overlapping versions must share one split"
+        );
+        assert_eq!(plan.splits()[0].data_files().len(), 2);
+        assert!(!plan.splits()[0].raw_convertible());
+
+        let batches = read_builder
+            .new_read()
+            .unwrap()
+            .to_arrow(plan.splits())
+            .unwrap()
+            .try_collect::<Vec<_>>()
+            .await
+            .unwrap();
+        assert_eq!(int_column(&batches, "id"), vec![1, 2, 3]);
+        assert_eq!(int_column(&batches, "value"), vec![11, 20, 30]);
+
+        let stale_filter = 
PredicateBuilder::new(merge_on_read.schema().fields())
+            .equal("value", Datum::Int(10))
+            .unwrap();
+        let stale_batches = read_rows(&merge_on_read, None, 
Some(stale_filter)).await;
+        assert_eq!(
+            stale_batches
+                .iter()
+                .map(RecordBatch::num_rows)
+                .sum::<usize>(),
+            0,
+            "a predicate matching only the superseded L0 value must not 
resurrect it"
+        );
+    }
+
+    #[tokio::test]
+    async fn dynamic_dv_merge_on_read_is_ignored_without_deletion_vectors() {
+        let file_io = test_file_io();
+        let table_path = "memory:/ignored_dynamic_dv_merge_on_read";
+        setup_dirs(&file_io, table_path).await;
+        let table = pk_table(&file_io, table_path, &[]);
+        write_commit(&table, &int_batch(vec![1, 2], vec![Some(10), 
Some(20)])).await;
+        write_commit(&table, &int_batch(vec![1, 3], vec![Some(11), 
Some(30)])).await;
+
+        let merge_on_read = table.copy_with_options(HashMap::from([(
+            "deletion-vectors.merge-on-read".to_string(),
+            "true".to_string(),
+        )]));
+        let batches = read_rows(&merge_on_read, None, None).await;
+        assert_eq!(int_column(&batches, "id"), vec![1, 2, 3]);
+        assert_eq!(int_column(&batches, "value"), vec![11, 20, 30]);
+    }
+
+    /// A MOR split can contain both uncompacted records and compacted source
+    /// files with deletion vectors. DV filtering must happen per physical file
+    /// before the surviving records enter the key merge; otherwise a deleted
+    /// key with no replacement is resurrected.
+    #[tokio::test]
+    async fn dv_merge_on_read_applies_deletion_vectors_before_key_merge() {
+        let file_io = test_file_io();
+        let table_path = "memory:/dv_merge_on_read_with_dv";
+        setup_dirs(&file_io, table_path).await;
+        let table = pk_table(
+            &file_io,
+            table_path,
+            &[
+                ("deletion-vectors.enabled", "true"),
+                ("deletion-vectors.merge-on-read", "true"),
+            ],
+        );
+
+        write_commit(&table, &int_batch(vec![1, 2], vec![Some(10), 
Some(20)])).await;
+        write_commit(&table, &int_batch(vec![1, 3], vec![Some(11), 
Some(30)])).await;
+
+        let all_files = table
+            .new_read_builder()
+            .new_scan()
+            .with_scan_all_files()
+            .plan()
+            .await
+            .unwrap();
+        let mut files = all_files
+            .splits()
+            .iter()
+            .flat_map(|split| split.data_files().iter().cloned())
+            .collect::<Vec<_>>();
+        files.sort_by_key(|file| file.min_sequence_number);
+        assert_eq!(files.len(), 2);
+        files[0].level = 1;
+
+        let deletion_file = write_deletion_file(&file_io, table_path, 
&[1]).await;
+        let split = DataSplitBuilder::new()
+            .with_snapshot(2)
+            .with_partition(BinaryRow::new(0))
+            .with_bucket(0)
+            .with_bucket_path(format!("{table_path}/bucket-0"))
+            .with_total_buckets(1)
+            .with_data_files(files)
+            .with_data_deletion_files(vec![Some(deletion_file), None])
+            .with_raw_convertible(false)
+            .build()
+            .unwrap();
+
+        let batches = table
+            .new_read_builder()
+            .new_read()
+            .unwrap()
+            .to_arrow(&[split])
+            .unwrap()
+            .try_collect::<Vec<_>>()
+            .await
+            .unwrap();
+        assert_eq!(int_column(&batches, "id"), vec![1, 3]);
+        assert_eq!(int_column(&batches, "value"), vec![11, 30]);
+    }
+
+    /// Disjoint files in one split are consumed as sequential merge groups. A
+    /// DV attached to a late file must not be read before an earlier group can
+    /// emit its output.
+    #[tokio::test]
+    async fn dv_merge_on_read_loads_deletion_vectors_per_file() {
+        let file_io = test_file_io();
+        let table_path = "memory:/dv_merge_on_read_lazy_dv";
+        setup_dirs(&file_io, table_path).await;
+        let table = pk_table(
+            &file_io,
+            table_path,
+            &[
+                ("deletion-vectors.enabled", "true"),
+                ("deletion-vectors.merge-on-read", "true"),
+            ],
+        );
+
+        let mut files = Vec::new();
+        for i in 0..10 {
+            files.push(
+                write_multi_row_group_kv_file(
+                    &file_io,
+                    table_path,
+                    &format!("part-{i}.parquet"),
+                    i * 1_000,
+                    i64::from(i),
+                    i,
+                )
+                .await,
+            );
+        }
+        let mut deletion_files = vec![None; files.len()];
+        deletion_files[9] = Some(DeletionFile::new(
+            format!("{table_path}/index/not-yet-read.dv"),
+            0,
+            1,
+            Some(1),
+        ));
+        let split = DataSplitBuilder::new()
+            .with_snapshot(1)
+            .with_partition(BinaryRow::new(0))
+            .with_bucket(0)
+            .with_bucket_path(format!("{table_path}/bucket-0"))
+            .with_total_buckets(1)
+            .with_data_files(files)
+            .with_data_deletion_files(deletion_files)
+            .with_raw_convertible(false)
+            .build()
+            .unwrap();
+
+        let mut stream = table
+            .new_read_builder()
+            .new_read()
+            .unwrap()
+            .to_arrow(&[split])
+            .unwrap();
+        let first = stream
+            .next()
+            .await
+            .expect("the first output batch")
+            .expect("a late deletion vector must be loaded lazily");
+        assert_eq!(first.num_rows(), 128);
+
+        let err = stream.try_collect::<Vec<_>>().await.unwrap_err();
+        assert!(
+            err.to_string().contains("not-yet-read.dv"),
+            "the late file must still surface its missing DV when reached: 
{err:?}"
+        );
+    }
+
     fn int_key(value: i32) -> Vec<u8> {
         let mut builder = crate::spec::BinaryRowBuilder::new(1);
         builder.write_int(0, value);
@@ -1171,6 +1450,47 @@ mod tests {
         );
     }
 
+    #[tokio::test]
+    async fn dv_mor_table_read_bounds_merge_fan_in() {
+        let file_io = test_file_io();
+        let table_path = "memory:/dv_mor_merge_fan_in_limit";
+        let table = pk_table(
+            &file_io,
+            table_path,
+            &[
+                ("deletion-vectors.enabled", "true"),
+                ("deletion-vectors.merge-on-read", "true"),
+            ],
+        );
+        let split = DataSplitBuilder::new()
+            .with_snapshot(1)
+            .with_partition(BinaryRow::new(0))
+            .with_bucket(0)
+            .with_bucket_path(format!("{table_path}/bucket-0"))
+            .with_total_buckets(1)
+            .with_data_files(
+                (0..257)
+                    .map(|i| dummy_data_file(format!("file-{i}.parquet")))
+                    .collect(),
+            )
+            .build()
+            .unwrap();
+
+        let err = table
+            .new_read_builder()
+            .new_read()
+            .unwrap()
+            .to_arrow(&[split])
+            .unwrap()
+            .try_collect::<Vec<_>>()
+            .await
+            .unwrap_err();
+        assert!(
+            matches!(err, Error::Unsupported { message } if 
message.contains("sorted-run input streams")),
+            "DV merge-on-read must reject unbounded production merge fan-in"
+        );
+    }
+
     #[test]
     fn sorted_run_planning_limits_each_section_to_overlap_depth() {
         let file = |name: &str, min: i32, max: i32| {
diff --git a/crates/paimon/src/table/table_read.rs 
b/crates/paimon/src/table/table_read.rs
index c96ffab7..b36af841 100644
--- a/crates/paimon/src/table/table_read.rs
+++ b/crates/paimon/src/table/table_read.rs
@@ -41,6 +41,8 @@ use futures::{stream, StreamExt};
 use std::cmp::Ordering;
 use std::sync::Arc;
 
+const MAX_MERGE_INPUT_STREAMS: usize = 256;
+
 /// Table read: reads data from splits (e.g. produced by [TableScan::plan]).
 ///
 /// Reference: 
[pypaimon.read.table_read.TableRead](https://github.com/apache/paimon/blob/master/paimon-python/pypaimon/read/table_read.py)
@@ -684,7 +686,7 @@ impl<'a> PaimonTableRead<'a> {
                     .collect(),
                 read_batch_size: core_options.read_batch_size()?,
                 merge_splits: true,
-                max_merge_input_streams: Some(256),
+                max_merge_input_streams: Some(MAX_MERGE_INPUT_STREAMS),
                 // Diff primes the before and after streams in sequence. 
Keeping
                 // a row-group permit across yielded batches can otherwise let
                 // the first side block the second side indefinitely.
@@ -772,9 +774,10 @@ impl<'a> PaimonTableRead<'a> {
             return self.read_raw(data_splits);
         }
 
-        // Deletion-vector tables read raw by design: stale versions of a key
-        // are masked by DVs, not merged, and KeyValueFileReader does not
-        // support DVs. Keep the plain level-0 dispatch for them.
+        // Compacted deletion-vector splits read raw: their stale versions are
+        // masked directly by DVs. A split containing level-0 data goes through
+        // the key merge; KeyValueFileReader applies any attached per-file DVs
+        // before merging the uncompacted versions.
         let mut kv_splits = Vec::new();
         let mut raw_splits = Vec::new();
         for split in data_splits {
@@ -824,7 +827,9 @@ impl<'a> PaimonTableRead<'a> {
                     .collect(),
                 read_batch_size: core_options.read_batch_size()?,
                 merge_splits: false,
-                max_merge_input_streams: None,
+                max_merge_input_streams: 
(core_options.deletion_vectors_enabled()
+                    && core_options.deletion_vectors_merge_on_read())
+                .then_some(MAX_MERGE_INPUT_STREAMS),
                 parquet_read_budget: Some(self.parquet_read_budget()?),
             },
         );
@@ -1440,8 +1445,9 @@ fn scalar_compare(
 /// planning treats the missing stat as "no deletes" for compatibility, so the
 /// read side must fall back to the merge reader, which drops them.
 ///
-/// Deletion-vector tables keep the plain level-0 dispatch: stale versions are
-/// masked by DVs and KeyValueFileReader does not support DVs.
+/// Deletion-vector tables merge only splits containing level-0 files. Fully
+/// compacted splits stay on the raw path, while the merge reader applies any
+/// attached DVs before reconciling uncompacted key versions.
 fn pk_split_needs_merge(split: &DataSplit, dv_enabled: bool) -> bool {
     if dv_enabled {
         return split.data_files().iter().any(|f| f.level == 0);
diff --git a/crates/paimon/src/table/table_scan.rs 
b/crates/paimon/src/table/table_scan.rs
index a025fe60..e4cd60ef 100644
--- a/crates/paimon/src/table/table_scan.rs
+++ b/crates/paimon/src/table/table_scan.rs
@@ -646,6 +646,7 @@ fn should_skip_level_zero_for_scan(
     scan_all_files: bool,
     has_primary_keys: bool,
     deletion_vectors_enabled: bool,
+    deletion_vectors_merge_on_read: bool,
     merge_engine: crate::Result<crate::spec::MergeEngine>,
 ) -> bool {
     if scan_all_files {
@@ -655,7 +656,8 @@ fn should_skip_level_zero_for_scan(
         return false;
     }
 
-    deletion_vectors_enabled || merge_engine.is_ok_and(|e| e == 
crate::spec::MergeEngine::FirstRow)
+    (deletion_vectors_enabled && !deletion_vectors_merge_on_read)
+        || merge_engine.is_ok_and(|e| e == crate::spec::MergeEngine::FirstRow)
 }
 
 fn is_system_field_id(field_id: i32) -> bool {
@@ -1180,6 +1182,7 @@ impl<'a> PaimonTableScan<'a> {
             self.scan_all_files,
             has_primary_keys,
             deletion_vectors_enabled,
+            core_options.deletion_vectors_merge_on_read(),
             core_options.merge_engine(),
         );
 
@@ -1382,8 +1385,10 @@ impl<'a> PaimonTableScan<'a> {
     /// `KeyValueFileReader`.
     ///
     /// Exempt (full predicates kept):
-    /// - Deletion-vector tables: they read raw with per-row masks, stats are
-    ///   a superset of live rows, full pruning stays safe.
+    /// - Deletion-vector tables without merge-on-read: they read raw with
+    ///   per-row masks, stats are a superset of live rows, full pruning stays
+    ///   safe. With merge-on-read enabled, visible L0 versions require the
+    ///   same key-only pruning rule as an ordinary PK merge read.
     /// - `merge-engine=first-row`: planned with `skip_level_zero` and read
     ///   via `DataFileReader` (see `TableRead::to_arrow`), no merge on the
     ///   read path — pruning a file drops exactly the rows the raw path's
@@ -1393,13 +1398,17 @@ impl<'a> PaimonTableScan<'a> {
         let has_primary_keys = !self.table.schema().primary_keys().is_empty();
         let core_options = CoreOptions::new(self.table.schema().options());
         let deletion_vectors_enabled = core_options.deletion_vectors_enabled();
+        let deletion_vectors_merge_on_read = 
core_options.deletion_vectors_merge_on_read();
         // An unknown merge engine stays conservative (key-only pruning); the
         // read side fails on it anyway before returning rows.
         let first_row = matches!(
             core_options.merge_engine(),
             Ok(crate::spec::MergeEngine::FirstRow)
         );
-        if has_primary_keys && !deletion_vectors_enabled && !first_row {
+        if has_primary_keys
+            && (!deletion_vectors_enabled || deletion_vectors_merge_on_read)
+            && !first_row
+        {
             retain_primary_key_conjuncts(
                 &self.data_predicates,
                 self.table.schema().fields(),
@@ -1873,10 +1882,12 @@ impl<'a> PaimonTableScan<'a> {
         // sort-merge reader sees every version of a key. The comparator 
decodes
         // the trimmed-PK min/max keys written by the kv writer.
         //
-        // Deletion-vector and first-row tables read without merging (stale 
rows
-        // are masked by DVs / level-0 is skipped), so they keep plain 
size-based
-        // packing like Java's MergeTreeSplitGenerator fast path.
-        let read_merges_overlapping_keys = 
!core_options.deletion_vectors_enabled()
+        // Deletion-vector tables without merge-on-read and first-row tables 
read
+        // without merging (stale rows are masked by DVs / level-0 is skipped),
+        // so they keep plain size-based packing. DV merge-on-read includes L0
+        // files and must preserve overlapping key ranges just like ordinary 
MOR.
+        let read_merges_overlapping_keys = 
(!core_options.deletion_vectors_enabled()
+            || core_options.deletion_vectors_merge_on_read())
             && !matches!(
                 core_options.merge_engine(),
                 Ok(crate::spec::MergeEngine::FirstRow)
@@ -2799,6 +2810,7 @@ mod tests {
             false,
             true,
             false,
+            false,
             Ok(crate::spec::MergeEngine::FirstRow),
         ));
     }
@@ -2809,10 +2821,29 @@ mod tests {
             true,
             true,
             false,
+            false,
             Ok(crate::spec::MergeEngine::FirstRow),
         ));
     }
 
+    #[test]
+    fn test_dv_merge_on_read_controls_batch_level_zero_visibility() {
+        assert!(should_skip_level_zero_for_scan(
+            false,
+            true,
+            true,
+            false,
+            Ok(crate::spec::MergeEngine::Deduplicate),
+        ));
+        assert!(!should_skip_level_zero_for_scan(
+            false,
+            true,
+            true,
+            true,
+            Ok(crate::spec::MergeEngine::Deduplicate),
+        ));
+    }
+
     #[test]
     fn test_partition_filter_decode_failure_fails_open() {
         let fields = partition_string_field();
@@ -3592,6 +3623,57 @@ mod tests {
         );
     }
 
+    /// Enabling DV merge-on-read puts the table back on a key-merge path for
+    /// visible level-0 files. Non-key stats pruning is therefore unsafe for 
the
+    /// same reason as ordinary MOR: pruning the newest version can resurrect 
an
+    /// older matching value.
+    #[tokio::test]
+    async fn test_dv_merge_on_read_stats_pruning_ignores_non_key_conjuncts() {
+        let table_path = "memory:/test_dv_mor_stats_gate";
+        let table = 
pk_stats_gate_table(table_path).copy_with_options(HashMap::from([
+            ("deletion-vectors.enabled".to_string(), "true".to_string()),
+            (
+                "deletion-vectors.merge-on-read".to_string(),
+                "true".to_string(),
+            ),
+        ]));
+        setup_scan_trace_dirs(&table).await;
+
+        let mut old = pk_stats_file("old-version.parquet", (1, 5), (100, 200));
+        old.level = 1;
+        let mut new = pk_stats_file("new-version.parquet", (1, 5), (10, 60));
+        new.level = 1;
+        TableCommit::new(table.clone(), "dv-mor-gate-test".to_string())
+            .commit(vec![CommitMessage::new(
+                BinaryRowBuilder::new(0).build_serialized(),
+                0,
+                vec![old, new],
+            )])
+            .await
+            .unwrap();
+
+        let fields = vec![
+            DataField::new(0, "id".to_string(), DataType::Int(IntType::new())),
+            DataField::new(1, "value".to_string(), 
DataType::Int(IntType::new())),
+        ];
+        let value_filter = PredicateBuilder::new(&fields)
+            .greater_than("value", Datum::Int(90))
+            .unwrap();
+        let mut reader = table.new_read_builder();
+        reader.with_filter(value_filter);
+        let (plan, trace) = reader.new_scan().plan_with_trace().await.unwrap();
+
+        assert_eq!(trace.manifest_entries_pruned_by_data_stats, 0);
+        assert_eq!(
+            plan.splits()
+                .iter()
+                .map(|split| split.data_files().len())
+                .sum::<usize>(),
+            2,
+            "both key versions must reach the merge path"
+        );
+    }
+
     /// `merge-engine=first-row` PK tables read raw (no merge on the read
     /// path: planned with `skip_level_zero`, read via `DataFileReader`), so
     /// pruning a file by a non-key conjunct cannot resurrect anything — it
diff --git a/docs/src/sql.md b/docs/src/sql.md
index 34aa36a4..6eea6d22 100644
--- a/docs/src/sql.md
+++ b/docs/src/sql.md
@@ -1982,6 +1982,14 @@ Set via `WITH ('key' = 'value')` at table creation time, 
or dynamically via `SET
 | `'merge-engine' = 'partial-update'` | Basic partial-update engine for PK 
tables |
 | `'merge-engine' = 'aggregation'` | Basic aggregation engine for PK tables |
 
+For deletion-vector-enabled primary-key tables using the default `deduplicate`
+engine, batch scans hide uncompacted level-0 files by default. Set
+`'deletion-vectors.merge-on-read' = 'true'` to include those files and merge
+their key versions on read. Existing deletion vectors are applied before the
+key merge. This option affects batch snapshot reads only; it does not change
+streaming or changelog behavior. It takes effect only when
+`'deletion-vectors.enabled' = 'true'`; otherwise it is ignored.
+
 Rust supports the basic partial-update engine with latest-non-null semantics.
 Set either `'ignore-delete' = 'true'` or
 `'partial-update.ignore-delete' = 'true'` to ignore `DELETE` and
@@ -2086,6 +2094,7 @@ the normal physical format without wrapping the writer.
 | `'data-evolution.enabled' = 'true'` | Enable data evolution (partial-column 
writes, row-level UPDATE/MERGE/DELETE) |
 | `'global-index.enabled' = 'true'` | Enable global index metadata and reads |
 | `'deletion-vectors.enabled' = 'true'` | Enable deletion vectors |
+| `'deletion-vectors.merge-on-read' = 'true'` | Include and key-merge 
uncompacted level-0 files in DV-enabled deduplicate batch reads |
 | `'changelog-producer' = 'input'` | Changelog producer; primary-key tables 
support reads and writes in this mode |
 
 Cross-partition updates are not configured by an option: a primary-key table is

Reply via email to