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 23f60a1a feat(table): add primary-key full-text read path (#599)
23f60a1a is described below

commit 23f60a1acb94b8065b3d021c2a8b9c83baf182a8
Author: Junrui Lee <[email protected]>
AuthorDate: Thu Jul 23 20:57:19 2026 +0800

    feat(table): add primary-key full-text read path (#599)
---
 .../paimon/src/table/full_text_search_builder.rs   | 306 ++++++-
 crates/paimon/src/table/mod.rs                     |   6 +
 .../paimon/src/table/pk_full_text_bucket_search.rs | 642 ++++++++++++++
 crates/paimon/src/table/pk_full_text_read.rs       | 969 +++++++++++++++++++++
 crates/paimon/src/table/pk_full_text_scan.rs       | 711 +++++++++++++++
 crates/paimon/src/table/vector_search_builder.rs   |   6 +-
 6 files changed, 2634 insertions(+), 6 deletions(-)

diff --git a/crates/paimon/src/table/full_text_search_builder.rs 
b/crates/paimon/src/table/full_text_search_builder.rs
index 005d65e6..702c3295 100644
--- a/crates/paimon/src/table/full_text_search_builder.rs
+++ b/crates/paimon/src/table/full_text_search_builder.rs
@@ -25,14 +25,19 @@ use crate::spec::{
     CoreOptions, DataField, FileKind, GlobalIndexSearchMode, IndexFileMeta, 
IndexManifest,
     IndexManifestEntry, ROW_ID_FIELD_NAME,
 };
+use crate::table::data_file_reader::DataFileReader;
 use crate::table::full_text_index_adapter::{search_full_text_file, 
search_full_text_index};
 use crate::table::global_index_scanner::{
     deleted_row_ranges_for_data_evolution_dvs, search_limit_with_deleted_rows,
     unindexed_ranges_for_global_index_entries, RowRangeIndex,
 };
-use crate::table::{find_field_id_by_name, merge_row_ranges, RowRange, Table};
+use crate::table::pk_full_text_read::PrimaryKeyFullTextRead;
+use crate::table::pk_full_text_scan::PrimaryKeyFullTextScan;
+use crate::table::{
+    find_field_id_by_name, merge_row_ranges, ArrowRecordBatchStream, RowRange, 
Table,
+};
 use arrow_array::{Array, Int64Array, LargeStringArray, RecordBatch, 
StringArray};
-use futures::{StreamExt, TryStreamExt};
+use futures::{stream, StreamExt, TryStreamExt};
 use paimon_ftindex_core::io::PosWriter;
 use paimon_ftindex_core::{FullTextIndexConfig, FullTextIndexWriter};
 use roaring::RoaringTreemap;
@@ -111,7 +116,8 @@ impl<'a> FullTextSearchBuilder<'a> {
 
     pub async fn execute_scored(&self) -> crate::Result<SearchResult> {
         // Fail closed: returns data-derived row ranges outside 
`TableScan`/`TableRead`.
-        
CoreOptions::new(self.table.schema().options()).ensure_read_authorized()?;
+        let core = CoreOptions::new(self.table.schema().options());
+        core.ensure_read_authorized()?;
         let text_column =
             self.text_column
                 .as_deref()
@@ -128,6 +134,21 @@ impl<'a> FullTextSearchBuilder<'a> {
             message: "Limit must be set via with_limit()".to_string(),
         })?;
 
+        // Primary-key full-text search does not produce global row-ids: it 
maps
+        // hits to physical `(data file, row position)` pairs. A 
scored/row-range
+        // search is therefore unsupported on the PK path — callers must use
+        // `execute_read`. Fail loud rather than fall through to the append/DE
+        // global-index path (which would search the wrong index and could 
return
+        // an empty or wrong result). Mirrors the vector builder.
+        if resolves_to_pk_full_text_path(&core, text_column) {
+            return Err(crate::Error::DataInvalid {
+                message: "primary-key full-text search does not produce global 
row ids; use the \
+                          materialized read (execute_read) instead"
+                    .to_string(),
+                source: None,
+            });
+        }
+
         let mut search = FullTextSearch::new(
             normalize_query_text(query_text, text_column)?,
             limit,
@@ -166,6 +187,113 @@ impl<'a> FullTextSearchBuilder<'a> {
         )
         .await
     }
+
+    /// Run the full-text search and materialize the matching rows as Arrow 
batches,
+    /// ordered best-score-first with a `__paimon_search_score` column 
appended (the
+    /// internal `_PKEY_VECTOR_POSITION` column is never exposed).
+    ///
+    /// Only the primary-key full-text path can materialize rows: it produces
+    /// physical `(data file, row position)` hits that a subsequent read turns 
into
+    /// table rows. Dispatch mirrors Java `primaryKeyFullTextDefinition` — the 
PK
+    /// path is taken only when data-evolution is DISABLED and the queried 
column is
+    /// a configured `pk-full-text.index.columns` entry. The PK full-text read 
is
+    /// FAST-mode only; `FULL`/`DETAIL` fail loud (no silent degrade). A query 
that
+    /// does not resolve to the PK path also fails loud rather than silently
+    /// returning nothing, since the append/data-evolution materialized read 
is not
+    /// supported here.
+    pub async fn execute_read(&self) -> crate::Result<ArrowRecordBatchStream> {
+        // Fail closed: returns data outside `TableScan`/`TableRead`.
+        let core = CoreOptions::new(self.table.schema().options());
+        core.ensure_read_authorized()?;
+        let text_column =
+            self.text_column
+                .as_deref()
+                .ok_or_else(|| crate::Error::ConfigInvalid {
+                    message: "Text column must be set via 
with_text_column()".to_string(),
+                })?;
+        let query_text = self
+            .query_text
+            .as_deref()
+            .ok_or_else(|| crate::Error::ConfigInvalid {
+                message: "Query text must be set via 
with_query_text()".to_string(),
+            })?;
+        let limit = self.limit.ok_or_else(|| crate::Error::ConfigInvalid {
+            message: "Limit must be set via with_limit()".to_string(),
+        })?;
+
+        if !resolves_to_pk_full_text_path(&core, text_column) {
+            return Err(crate::Error::Unsupported {
+                message: "materialized full-text read (execute_read) is only 
supported on the \
+                          primary-key full-text path (data-evolution disabled 
and the column \
+                          configured in pk-full-text.index.columns)"
+                    .to_string(),
+            });
+        }
+
+        // FAST-only: reject FULL/DETAIL loud rather than silently degrading.
+        if core.global_index_search_mode()? != GlobalIndexSearchMode::Fast {
+            return Err(crate::Error::DataInvalid {
+                message: "primary-key full-text search supports only the FAST 
global-index search \
+                          mode"
+                    .to_string(),
+                source: None,
+            });
+        }
+
+        // Reject a non-positive limit at construction, before the empty-plan 
fast
+        // path — an empty table must not mask an invalid limit (mirrors Java
+        // `PrimaryKeyFullTextRead`, and matches `search_bucket`'s own guard).
+        if limit == 0 {
+            return Err(crate::Error::ConfigInvalid {
+                message: "Limit must be positive".to_string(),
+            });
+        }
+
+        // Resolve the queried column's schema field id for the scan/field-id 
guard.
+        let field_id = find_field_id_by_name(self.table.schema().fields(), 
text_column)
+            .ok_or_else(|| crate::Error::DataInvalid {
+                message: format!("full-text search column '{text_column}' does 
not exist"),
+                source: None,
+            })?;
+
+        let plan = PrimaryKeyFullTextScan::new(self.table, field_id, None)
+            .plan()
+            .await?;
+        if plan.splits.is_empty() {
+            return Ok(Box::pin(stream::empty()));
+        }
+
+        // A predicate-free materialization reader projecting the user table
+        // columns (mirrors `table_read.rs::new_data_file_reader` with an empty
+        // predicate list). The PK read appends the score column itself.
+        let materialize_reader = DataFileReader::new(
+            self.table.file_io().clone(),
+            self.table.schema_manager().clone(),
+            self.table.schema().id(),
+            self.table.schema().fields().to_vec(),
+            self.table.schema().fields().to_vec(),
+            Vec::new(),
+        );
+        let read = PrimaryKeyFullTextRead::new(
+            self.table.file_io().clone(),
+            materialize_reader,
+            self.table.location().trim_end_matches('/').to_string(),
+        );
+        read.read(&plan, query_text, limit).await
+    }
+}
+
+/// Whether a query on `text_column` resolves to the primary-key full-text read
+/// path. Mirrors Java `primaryKeyFullTextDefinition`: taken only when
+/// data-evolution is DISABLED and the column is a configured
+/// `pk-full-text.index.columns` entry (membership via the non-erroring 
accessor so
+/// a malformed config cannot abort an unrelated append/DE query).
+fn resolves_to_pk_full_text_path(core: &CoreOptions<'_>, text_column: &str) -> 
bool {
+    !core.data_evolution_enabled()
+        && core
+            .primary_key_full_text_index_columns()
+            .iter()
+            .any(|c| c == text_column)
 }
 
 /// Evaluate a full-text search query against full-text indexes found in the 
index manifest.
@@ -1203,4 +1331,176 @@ mod tests {
             "full-text search must fail closed for a query-auth table"
         );
     }
+
+    /// A primary-key full-text table: data-evolution off, `body` configured 
as a
+    /// `pk-full-text.index.columns` entry, so a `body` query resolves to the 
PK
+    /// full-text path. `extra` appends/overrides options (e.g. the search 
mode).
+    fn pk_full_text_table(name: &str, extra: &[(&str, &str)]) -> Table {
+        let mut builder = Schema::builder()
+            .column("id", DataType::Int(IntType::new()))
+            .column("body", DataType::VarChar(VarCharType::string_type()))
+            .primary_key(["id"])
+            .option("bucket", "1")
+            .option("pk-full-text.index.columns", "body");
+        for (k, v) in extra {
+            builder = builder.option(*k, *v);
+        }
+        let schema = builder.build().unwrap();
+        Table::new(
+            FileIOBuilder::new("memory").build().unwrap(),
+            Identifier::new("default", name),
+            format!("memory:/{name}"),
+            TableSchema::new(0, &schema),
+            None,
+        )
+    }
+
+    // (d) The PK full-text path produces physical positions, not global 
row-ids:
+    // `execute` / `execute_scored` must fail loud and point callers at 
execute_read.
+    #[tokio::test]
+    async fn pk_full_text_execute_fails_loud_use_execute_read() {
+        let table = pk_full_text_table("pk_ft_execute_loud", &[]);
+        let err = table
+            .new_full_text_search_builder()
+            .with_text_column("body")
+            .with_query_text(r#"{"match":{"query":"alpha"}}"#)
+            .with_limit(10)
+            .execute()
+            .await
+            .unwrap_err();
+        assert!(
+            matches!(&err, crate::Error::DataInvalid { message, .. } if 
message.contains("execute_read")),
+            "PK full-text execute must fail loud pointing at execute_read, 
got: {err}"
+        );
+    }
+
+    #[tokio::test]
+    async fn pk_full_text_execute_scored_fails_loud_use_execute_read() {
+        let table = pk_full_text_table("pk_ft_scored_loud", &[]);
+        let err = table
+            .new_full_text_search_builder()
+            .with_text_column("body")
+            .with_query_text(r#"{"match":{"query":"alpha"}}"#)
+            .with_limit(10)
+            .execute_scored()
+            .await
+            .unwrap_err();
+        assert!(
+            matches!(&err, crate::Error::DataInvalid { message, .. } if 
message.contains("global row ids")),
+            "PK full-text execute_scored must fail loud, got: {err}"
+        );
+    }
+
+    // (c) FAST-only: FULL / DETAIL global-index search modes must fail loud 
on the
+    // PK full-text read rather than silently degrade.
+    #[tokio::test]
+    async fn pk_full_text_execute_read_rejects_full_mode() {
+        let table = pk_full_text_table("pk_ft_full_mode", 
&[("global-index.search-mode", "full")]);
+        let err = match table
+            .new_full_text_search_builder()
+            .with_text_column("body")
+            .with_query_text(r#"{"match":{"query":"alpha"}}"#)
+            .with_limit(10)
+            .execute_read()
+            .await
+        {
+            Ok(_) => panic!("FULL mode PK full-text read must fail loud"),
+            Err(e) => e,
+        };
+        assert!(
+            matches!(&err, crate::Error::DataInvalid { message, .. } if 
message.contains("FAST")),
+            "FULL mode PK full-text read must fail loud, got: {err}"
+        );
+    }
+
+    #[tokio::test]
+    async fn pk_full_text_execute_read_rejects_detail_mode() {
+        let table = pk_full_text_table(
+            "pk_ft_detail_mode",
+            &[("global-index.search-mode", "detail")],
+        );
+        let err = match table
+            .new_full_text_search_builder()
+            .with_text_column("body")
+            .with_query_text(r#"{"match":{"query":"alpha"}}"#)
+            .with_limit(10)
+            .execute_read()
+            .await
+        {
+            Ok(_) => panic!("DETAIL mode PK full-text read must fail loud"),
+            Err(e) => e,
+        };
+        assert!(
+            matches!(&err, crate::Error::DataInvalid { message, .. } if 
message.contains("FAST")),
+            "DETAIL mode PK full-text read must fail loud, got: {err}"
+        );
+    }
+
+    // (e) A non-PK (append) table: execute_read is unsupported and must fail 
loud
+    // rather than silently return nothing; the append execute/execute_scored 
path
+    // stays unaffected (exercised by the raw-search tests above).
+    #[tokio::test]
+    async fn append_execute_read_fails_loud_unsupported() {
+        let file_io = FileIOBuilder::new("memory").build().unwrap();
+        let table = full_text_raw_table(&file_io, 
"memory:/append_ft_execute_read");
+        let err = match table
+            .new_full_text_search_builder()
+            .with_text_column("body")
+            .with_query_text(r#"{"match":{"query":"alpha"}}"#)
+            .with_limit(10)
+            .execute_read()
+            .await
+        {
+            Ok(_) => panic!("append full-text execute_read must fail loud"),
+            Err(e) => e,
+        };
+        assert!(
+            matches!(&err, crate::Error::Unsupported { message } if 
message.contains("primary-key full-text path")),
+            "append full-text execute_read must fail loud as unsupported, got: 
{err}"
+        );
+    }
+
+    // Dispatch: with data-evolution ENABLED, a configured pk-full-text column 
must
+    // NOT resolve to the PK path (mirrors Java 
`primaryKeyFullTextDefinition`), so
+    // execute_scored takes the append/DE path instead of the PK fail-loud 
guard.
+    #[tokio::test]
+    async fn data_evolution_enabled_does_not_take_pk_path() {
+        let de_on = HashMap::from([
+            ("pk-full-text.index.columns".to_string(), "body".to_string()),
+            ("data-evolution.enabled".to_string(), "true".to_string()),
+        ]);
+        let core = CoreOptions::new(&de_on);
+        assert!(
+            !resolves_to_pk_full_text_path(&core, "body"),
+            "data-evolution on must not resolve to the PK full-text path"
+        );
+        // Off + configured column -> PK path; a non-configured column -> not 
PK.
+        let de_off =
+            HashMap::from([("pk-full-text.index.columns".to_string(), 
"body".to_string())]);
+        let core_off = CoreOptions::new(&de_off);
+        assert!(resolves_to_pk_full_text_path(&core_off, "body"));
+        assert!(!resolves_to_pk_full_text_path(&core_off, "other"));
+    }
+
+    // A non-positive limit must fail loud at construction, even on an empty 
table
+    // (before the empty-plan fast path).
+    #[tokio::test]
+    async fn pk_full_text_execute_read_rejects_zero_limit() {
+        let table = pk_full_text_table("pk_ft_zero_limit", &[]);
+        let err = match table
+            .new_full_text_search_builder()
+            .with_text_column("body")
+            .with_query_text(r#"{"match":{"query":"alpha"}}"#)
+            .with_limit(0)
+            .execute_read()
+            .await
+        {
+            Ok(_) => panic!("zero limit PK full-text read must fail loud"),
+            Err(e) => e,
+        };
+        assert!(
+            matches!(&err, crate::Error::ConfigInvalid { message } if 
message.contains("positive")),
+            "zero limit must fail loud, got: {err}"
+        );
+    }
 }
diff --git a/crates/paimon/src/table/mod.rs b/crates/paimon/src/table/mod.rs
index 0b76d283..ca6db9c8 100644
--- a/crates/paimon/src/table/mod.rs
+++ b/crates/paimon/src/table/mod.rs
@@ -58,7 +58,13 @@ mod lumina_index_build_builder;
 pub(crate) mod merge_tree_split_generator;
 mod partition_filter;
 mod partition_stat;
+#[cfg(feature = "fulltext")]
+mod pk_full_text_bucket_search;
 mod pk_full_text_bucket_state;
+#[cfg(feature = "fulltext")]
+mod pk_full_text_read;
+#[cfg(feature = "fulltext")]
+mod pk_full_text_scan;
 mod pk_vector_data_file_reader;
 mod pk_vector_indexed_split_read;
 mod pk_vector_orchestrator;
diff --git a/crates/paimon/src/table/pk_full_text_bucket_search.rs 
b/crates/paimon/src/table/pk_full_text_bucket_search.rs
new file mode 100644
index 00000000..298733f0
--- /dev/null
+++ b/crates/paimon/src/table/pk_full_text_bucket_search.rs
@@ -0,0 +1,642 @@
+// 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.
+
+//! Bucket-local primary-key full-text search: the algorithmic core that turns 
a
+//! planned bucket split into scored candidates.
+//!
+//! Mirror of Java `PrimaryKeyFullTextBucketSearch.searchRankings`. For every
+//! current payload in the split it lays the payload's ordered source files 
out as
+//! a cumulative archive: file `i` owns archive rows `[offset_i, offset_i +
+//! row_count_i)`, where `offset_i` is the running sum of prior source row 
counts.
+//! A source is *active* iff its data file is present among the split's data 
files.
+//! When any source is inactive or any active source carries deletions, an 
include
+//! allow-list of live archive rows is built (active ranges minus 
deletion-vector
+//! positions, each deletion mapped to `offset + local_position`); otherwise 
the
+//! whole archive is read. Each returned archive row is mapped back to a 
physical
+//! `(data file, row position)` and emitted as a score-tagged candidate.
+
+use std::collections::HashMap;
+
+use roaring::RoaringTreemap;
+
+use crate::deletion_vector::DeletionVector;
+use crate::ftindex::reader::FullTextArchiveReader;
+use crate::io::FileIO;
+use crate::spec::PrimaryKeyIndexSourceMeta;
+use crate::table::pk_full_text_read::PrimaryKeyFullTextCandidate;
+use crate::table::pk_full_text_scan::PrimaryKeyFullTextSearchSplit;
+
+fn data_invalid(message: impl Into<String>) -> crate::Error {
+    crate::Error::DataInvalid {
+        message: message.into(),
+        source: None,
+    }
+}
+
+/// One source data file laid out at its cumulative archive `offset`. `active` 
is
+/// true iff the file is present among the split's data files; only active 
sources
+/// contribute live rows and may own a returned archive row. Mirrors the Java
+/// inner `SourceRange`.
+struct SourceRange {
+    file_name: String,
+    offset: i64,
+    row_count: i64,
+    active: bool,
+}
+
+impl SourceRange {
+    fn contains(&self, row_id: i64) -> bool {
+        row_id >= self.offset && row_id < self.offset + self.row_count
+    }
+}
+
+/// A payload's archive layout after validation: its ordered source ranges, the
+/// total archive row count, and the optional live-row include allow-list
+/// (`None` = read the whole archive; mirrors Java's `needsInclude ? … : 
null`).
+struct PreparedPayload {
+    source_ranges: Vec<SourceRange>,
+    total_row_count: i64,
+    include: Option<RoaringTreemap>,
+}
+
+/// Lay a payload's source files out as a cumulative archive and decide the
+/// include allow-list, validating each active source's row count against its
+/// data file. Mirrors the per-payload body of Java `searchRankings`.
+///
+/// `active_row_counts` maps an active data file name to its `DataFileMeta`
+/// row count; a source absent from the map is inactive (its data file is not 
in
+/// the split). `dvs` is keyed by data file name.
+fn prepare_payload(
+    source_meta: &PrimaryKeyIndexSourceMeta,
+    active_row_counts: &HashMap<String, i64>,
+    dvs: &HashMap<String, DeletionVector>,
+    payload_file_name: &str,
+) -> crate::Result<PreparedPayload> {
+    let mut source_ranges: Vec<SourceRange> = Vec::new();
+    let mut needs_include = false;
+    let mut total_row_count: i64 = 0;
+    for source in source_meta.source_files() {
+        let active_row_count = active_row_counts.get(source.file_name());
+        if let Some(&data_row_count) = active_row_count {
+            if source.row_count() != data_row_count {
+                return Err(data_invalid(format!(
+                    "full-text payload {payload_file_name} source row count 
does not match data file {}",
+                    source.file_name()
+                )));
+            }
+        }
+        let has_deletions = dvs.get(source.file_name()).is_some_and(|dv| 
!dv.is_empty());
+        let active = active_row_count.is_some();
+        needs_include |= !active || has_deletions;
+        source_ranges.push(SourceRange {
+            file_name: source.file_name().to_string(),
+            offset: total_row_count,
+            row_count: source.row_count(),
+            active,
+        });
+        total_row_count = total_row_count
+            .checked_add(source.row_count())
+            .ok_or_else(|| data_invalid("full-text source row counts overflow 
i64"))?;
+    }
+
+    let include = if needs_include {
+        Some(live_rows(&source_ranges, dvs)?)
+    } else {
+        None
+    };
+    Ok(PreparedPayload {
+        source_ranges,
+        total_row_count,
+        include,
+    })
+}
+
+/// Union of active source archive ranges `[offset, offset + row_count)` minus 
the
+/// deletion-vector positions mapped to archive positions (`offset + local`).
+/// Mirrors Java `liveRows`, including the fail-loud guard on out-of-range
+/// deletion positions.
+fn live_rows(
+    source_ranges: &[SourceRange],
+    dvs: &HashMap<String, DeletionVector>,
+) -> crate::Result<RoaringTreemap> {
+    let mut include = RoaringTreemap::new();
+    let mut deleted = RoaringTreemap::new();
+    for source in source_ranges {
+        if !source.active {
+            continue;
+        }
+        if source.row_count > 0 {
+            let start = source.offset as u64;
+            let end = (source.offset + source.row_count) as u64;
+            include.insert_range(start..end);
+        }
+        if let Some(dv) = dvs.get(&source.file_name) {
+            if !dv.is_empty() {
+                let row_count = source.row_count as u64;
+                for position in dv.iter() {
+                    if position >= row_count {
+                        return Err(data_invalid(format!(
+                            "deletion vector contains invalid row position 
{position} for source {}",
+                            source.file_name
+                        )));
+                    }
+                    deleted.insert(source.offset as u64 + position);
+                }
+            }
+        }
+    }
+    include -= &deleted;
+    Ok(include)
+}
+
+/// Resolve the active source owning `row_id`, failing loud if the id is 
outside
+/// `[0, total_row_count)` or lands in an inactive source (mirrors Java's two
+/// `checkArgument`s around `request.source(rowId)`).
+fn owning_active_source(
+    source_ranges: &[SourceRange],
+    total_row_count: i64,
+    row_id: i64,
+) -> crate::Result<&SourceRange> {
+    if row_id < 0 || row_id >= total_row_count {
+        return Err(data_invalid(format!(
+            "full-text index returned archive row position {row_id} outside 
row count {total_row_count}"
+        )));
+    }
+    match source_ranges.iter().find(|source| source.contains(row_id)) {
+        Some(source) if source.active => Ok(source),
+        _ => Err(data_invalid(format!(
+            "full-text index returned row position {row_id} from an inactive 
source"
+        ))),
+    }
+}
+
+/// Search one bucket's full-text payloads and return its scored candidates
+/// (unsorted; the read path fuses them cross-bucket via `top_k_by_score`).
+///
+/// `dvs` is keyed by data file name; `table_path` roots the index directory
+/// (`{table_path}/index/{payload}`). Mirrors Java
+/// `PrimaryKeyFullTextBucketSearch.searchRankings`, flattened to one candidate
+/// list per bucket.
+pub(crate) async fn search_bucket(
+    split: &PrimaryKeyFullTextSearchSplit,
+    query: &str,
+    limit: usize,
+    dvs: &HashMap<String, DeletionVector>,
+    file_io: &FileIO,
+    table_path: &str,
+    split_index: usize,
+) -> crate::Result<Vec<PrimaryKeyFullTextCandidate>> {
+    if limit == 0 {
+        return Err(data_invalid("full-text search limit must be positive"));
+    }
+
+    let data_split = &split.data_split;
+    let partition = data_split.partition().clone();
+    let bucket = data_split.bucket();
+
+    // Active data files by name, with a duplicate guard (mirror Java's `files`
+    // map). The split constructor already rejects duplicates, but re-checking
+    // keeps the invariant local to the search.
+    let mut active_row_counts: HashMap<String, i64> = HashMap::new();
+    for file in data_split.data_files() {
+        if active_row_counts
+            .insert(file.file_name.clone(), file.row_count)
+            .is_some()
+        {
+            return Err(data_invalid(format!(
+                "duplicate full-text source file {}",
+                file.file_name
+            )));
+        }
+    }
+
+    let mut candidates: Vec<PrimaryKeyFullTextCandidate> = Vec::new();
+    for payload in &split.current_payloads {
+        let gim = payload.global_index_meta.as_ref().ok_or_else(|| {
+            data_invalid(format!(
+                "full-text payload {} has no global index meta",
+                payload.file_name
+            ))
+        })?;
+        let source_meta = 
PrimaryKeyIndexSourceMeta::from_global_index_meta(gim)?;
+        let prepared = prepare_payload(&source_meta, &active_row_counts, dvs, 
&payload.file_name)?;
+
+        // An empty include means every row is deleted or inactive: nothing to
+        // search in this payload (mirror Java's `include.isEmpty()` skip).
+        if let Some(include) = &prepared.include {
+            if include.is_empty() {
+                continue;
+            }
+        }
+
+        let path = format!("{table_path}/index/{}", payload.file_name);
+        let input = file_io.new_input(&path)?;
+        let reader = FullTextArchiveReader::from_input_file(&input).await?;
+        let hits = match &prepared.include {
+            Some(include) => reader.search_with_include(query, limit, 
include)?,
+            None => reader.search(query, limit)?,
+        };
+
+        for (&row_id, &score) in hits.row_ids.iter().zip(hits.scores.iter()) {
+            if row_id < 0 || row_id >= prepared.total_row_count {
+                return Err(data_invalid(format!(
+                    "full-text index returned archive row position {row_id} 
outside row count {}",
+                    prepared.total_row_count
+                )));
+            }
+            // Defensive re-check of the engine's include filtering (mirror 
Java).
+            if let Some(include) = &prepared.include {
+                if !include.contains(row_id as u64) {
+                    continue;
+                }
+            }
+            let source =
+                owning_active_source(&prepared.source_ranges, 
prepared.total_row_count, row_id)?;
+            candidates.push(PrimaryKeyFullTextCandidate::new(
+                split_index,
+                partition.clone(),
+                bucket,
+                score,
+                source.file_name.clone(),
+                row_id - source.offset,
+            )?);
+        }
+    }
+
+    Ok(candidates)
+}
+
+#[cfg(test)]
+mod tests {
+    use super::*;
+    use crate::io::FileIOBuilder;
+    use crate::spec::stats::BinaryTableStats;
+    use crate::spec::{
+        BinaryRow, DataFileMeta, GlobalIndexMeta, IndexFileMeta, 
PrimaryKeyIndexSourceFile,
+    };
+    use crate::table::source::{DataSplit, DataSplitBuilder};
+    use bytes::Bytes;
+    use paimon_ftindex_core::io::PosWriter;
+    use paimon_ftindex_core::{FullTextIndexConfig, FullTextIndexWriter};
+    use roaring::RoaringBitmap;
+
+    const PK_FULL_TEXT_INDEX_TYPE: &str = "full-text";
+    const FILE_SOURCE_COMPACT: i32 = 1;
+
+    fn source_range(file_name: &str, offset: i64, row_count: i64, active: 
bool) -> SourceRange {
+        SourceRange {
+            file_name: file_name.to_string(),
+            offset,
+            row_count,
+            active,
+        }
+    }
+
+    // ---- helpers to build synthetic archives + source-meta payloads ----
+
+    /// Build a tiny full-text archive in memory via FT-PR1's core writer.
+    fn build_archive(docs: &[(i64, &str)]) -> Vec<u8> {
+        let mut writer = 
FullTextIndexWriter::new(FullTextIndexConfig::new()).unwrap();
+        for (row_id, text) in docs {
+            writer.add_document(*row_id, (*text).to_string()).unwrap();
+        }
+        let mut out = PosWriter::new(Vec::<u8>::new());
+        writer.write(&mut out).unwrap();
+        out.into_inner()
+    }
+
+    /// One Java `DataOutput#writeUTF` value (u16-BE length + modified UTF-8).
+    fn java_write_utf(s: &str) -> Vec<u8> {
+        let mut body = Vec::new();
+        for c in s.encode_utf16() {
+            if (0x0001..=0x007F).contains(&c) {
+                body.push(c as u8);
+            } else if c > 0x07FF {
+                body.push(0xE0 | (c >> 12) as u8);
+                body.push(0x80 | ((c >> 6) & 0x3F) as u8);
+                body.push(0x80 | (c & 0x3F) as u8);
+            } else {
+                body.push(0xC0 | (c >> 6) as u8);
+                body.push(0x80 | (c & 0x3F) as u8);
+            }
+        }
+        let mut out = (body.len() as u16).to_be_bytes().to_vec();
+        out.extend_from_slice(&body);
+        out
+    }
+
+    /// A `_SOURCE_META` frame (version 1) for a level and its ordered sources.
+    fn frame(data_level: i32, files: &[(&str, i64)]) -> Vec<u8> {
+        let mut out = Vec::new();
+        out.extend_from_slice(&1i32.to_be_bytes());
+        out.extend_from_slice(&data_level.to_be_bytes());
+        out.extend_from_slice(&(files.len() as i32).to_be_bytes());
+        for (name, rows) in files {
+            out.extend_from_slice(&java_write_utf(name));
+            out.extend_from_slice(&rows.to_be_bytes());
+        }
+        out
+    }
+
+    fn gim(field_id: i32, source_meta: Vec<u8>) -> GlobalIndexMeta {
+        GlobalIndexMeta {
+            row_range_start: 0,
+            row_range_end: 0,
+            index_field_id: field_id,
+            extra_field_ids: None,
+            index_meta: None,
+            source_meta: Some(source_meta),
+        }
+    }
+
+    fn ft_payload(file_name: &str, level: i32, files: &[(&str, i64)]) -> 
IndexFileMeta {
+        let total: i64 = files.iter().map(|(_, r)| *r).sum();
+        IndexFileMeta {
+            index_type: PK_FULL_TEXT_INDEX_TYPE.into(),
+            file_name: file_name.into(),
+            file_size: 1,
+            row_count: total as i32,
+            deletion_vectors_ranges: None,
+            global_index_meta: Some(gim(7, frame(level, files))),
+        }
+    }
+
+    fn dfm(name: &str, rows: i64) -> DataFileMeta {
+        DataFileMeta {
+            file_name: name.into(),
+            file_size: 1,
+            row_count: rows,
+            min_key: Vec::new(),
+            max_key: Vec::new(),
+            key_stats: BinaryTableStats::empty(),
+            value_stats: BinaryTableStats::empty(),
+            min_sequence_number: 0,
+            max_sequence_number: 0,
+            schema_id: 1,
+            level: 1,
+            extra_files: Vec::new(),
+            creation_time: None,
+            delete_row_count: None,
+            embedded_index: None,
+            file_source: Some(FILE_SOURCE_COMPACT),
+            value_stats_cols: None,
+            external_path: None,
+            first_row_id: Some(0),
+            write_cols: None,
+        }
+    }
+
+    fn data_split(files: Vec<DataFileMeta>) -> DataSplit {
+        DataSplitBuilder::new()
+            .with_snapshot(1)
+            .with_partition(BinaryRow::new(0))
+            .with_bucket(0)
+            .with_bucket_path("memory:/t/bucket-0".to_string())
+            .with_total_buckets(1)
+            .with_data_files(files)
+            .build()
+            .unwrap()
+    }
+
+    async fn write_archive(file_io: &FileIO, path: &str, bytes: Vec<u8>) {
+        let output = file_io.new_output(path).unwrap();
+        output.write(Bytes::from(bytes)).await.unwrap();
+    }
+
+    fn source_meta(files: &[(&str, i64)]) -> PrimaryKeyIndexSourceMeta {
+        let sources = files
+            .iter()
+            .map(|(n, r)| PrimaryKeyIndexSourceFile::new((*n).to_string(), 
*r).unwrap())
+            .collect();
+        PrimaryKeyIndexSourceMeta::new(1, sources).unwrap()
+    }
+
+    // ---- (a) include=None: all-active, no DV, all matches mapped ----
+    #[tokio::test]
+    async fn include_none_maps_all_matches_to_positions() {
+        // d0 owns archive rows 0..3, d1 owns 3..7; "alpha" at 0,2,4,6.
+        let file_io = FileIOBuilder::new("memory").build().unwrap();
+        let table_path = "memory:/ftbs_none";
+        let bytes = build_archive(&[
+            (0, "alpha"),
+            (1, "beta"),
+            (2, "alpha"),
+            (3, "beta"),
+            (4, "alpha"),
+            (5, "beta"),
+            (6, "alpha"),
+        ]);
+        write_archive(&file_io, &format!("{table_path}/index/ft-0"), 
bytes).await;
+
+        let split = PrimaryKeyFullTextSearchSplit::new(
+            data_split(vec![dfm("d0", 3), dfm("d1", 4)]),
+            vec![ft_payload("ft-0", 1, &[("d0", 3), ("d1", 4)])],
+            Vec::new(),
+        )
+        .unwrap();
+
+        let dvs: HashMap<String, DeletionVector> = HashMap::new();
+        let out = search_bucket(
+            &split,
+            r#"{"match":{"query":"alpha"}}"#,
+            10,
+            &dvs,
+            &file_io,
+            table_path,
+            0,
+        )
+        .await
+        .unwrap();
+
+        let mut got: Vec<(String, i64)> = out
+            .iter()
+            .map(|c| (c.data_file_name.clone(), c.row_position))
+            .collect();
+        got.sort();
+        assert_eq!(
+            got,
+            vec![
+                ("d0".to_string(), 0),
+                ("d0".to_string(), 2),
+                ("d1".to_string(), 1),
+                ("d1".to_string(), 3),
+            ]
+        );
+    }
+
+    // ---- (b) a DV-deleted archive position is excluded from results ----
+    #[tokio::test]
+    async fn deletion_vector_excludes_matched_position() {
+        // "beta" at archive rows 1,3,5; DV deletes d0 local pos 1 (archive 1).
+        let file_io = FileIOBuilder::new("memory").build().unwrap();
+        let table_path = "memory:/ftbs_dv";
+        let bytes = build_archive(&[
+            (0, "alpha"),
+            (1, "beta"),
+            (2, "alpha"),
+            (3, "beta"),
+            (4, "alpha"),
+            (5, "beta"),
+            (6, "alpha"),
+        ]);
+        write_archive(&file_io, &format!("{table_path}/index/ft-0"), 
bytes).await;
+
+        let split = PrimaryKeyFullTextSearchSplit::new(
+            data_split(vec![dfm("d0", 3), dfm("d1", 4)]),
+            vec![ft_payload("ft-0", 1, &[("d0", 3), ("d1", 4)])],
+            Vec::new(),
+        )
+        .unwrap();
+
+        let mut bitmap = RoaringBitmap::new();
+        bitmap.insert(1); // d0 local position 1
+        let mut dvs: HashMap<String, DeletionVector> = HashMap::new();
+        dvs.insert("d0".to_string(), DeletionVector::from_bitmap(bitmap));
+
+        let out = search_bucket(
+            &split,
+            r#"{"match":{"query":"beta"}}"#,
+            10,
+            &dvs,
+            &file_io,
+            table_path,
+            0,
+        )
+        .await
+        .unwrap();
+
+        let mut got: Vec<(String, i64)> = out
+            .iter()
+            .map(|c| (c.data_file_name.clone(), c.row_position))
+            .collect();
+        got.sort();
+        // Archive row 1 (d0 pos 1) is deleted; only d1 rows 3,5 survive.
+        assert_eq!(got, vec![("d1".to_string(), 0), ("d1".to_string(), 2)]);
+    }
+
+    // ---- (c) a row_id landing in an inactive source → Err ----
+    #[test]
+    fn inactive_source_row_id_errors() {
+        // d0 active 0..3, d1 inactive 3..7.
+        let ranges = vec![
+            source_range("d0", 0, 3, true),
+            source_range("d1", 3, 4, false),
+        ];
+        assert!(owning_active_source(&ranges, 7, 1).is_ok());
+        assert!(owning_active_source(&ranges, 7, 4).is_err());
+    }
+
+    // ---- (e) a row_id out of the archive's total range → Err ----
+    #[test]
+    fn out_of_range_row_id_errors() {
+        let ranges = vec![source_range("d0", 0, 3, true)];
+        assert!(owning_active_source(&ranges, 3, 2).is_ok());
+        assert!(owning_active_source(&ranges, 3, 3).is_err()); // == total
+        assert!(owning_active_source(&ranges, 3, -1).is_err());
+    }
+
+    // ---- (d) source row_count != active DataFileMeta.row_count → Err ----
+    #[test]
+    fn source_row_count_mismatch_errors() {
+        let meta = source_meta(&[("d0", 3), ("d1", 4)]);
+        let dvs: HashMap<String, DeletionVector> = HashMap::new();
+        // Active d1 has 5 rows on disk but the payload recorded 4 → mismatch.
+        let mut active: HashMap<String, i64> = HashMap::new();
+        active.insert("d0".to_string(), 3);
+        active.insert("d1".to_string(), 5);
+        assert!(prepare_payload(&meta, &active, &dvs, "ft-0").is_err());
+
+        // Matching row counts prepare cleanly.
+        let mut ok_active: HashMap<String, i64> = HashMap::new();
+        ok_active.insert("d0".to_string(), 3);
+        ok_active.insert("d1".to_string(), 4);
+        assert!(prepare_payload(&meta, &ok_active, &dvs, "ft-0").is_ok());
+    }
+
+    // ---- include decision: None when clean, Some when a source is inactive 
----
+    #[test]
+    fn include_is_none_only_when_clean() {
+        let meta = source_meta(&[("d0", 3), ("d1", 4)]);
+        let dvs: HashMap<String, DeletionVector> = HashMap::new();
+
+        // Both active, no DV → include None (read everything).
+        let mut all_active: HashMap<String, i64> = HashMap::new();
+        all_active.insert("d0".to_string(), 3);
+        all_active.insert("d1".to_string(), 4);
+        let clean = prepare_payload(&meta, &all_active, &dvs, "ft-0").unwrap();
+        assert!(clean.include.is_none());
+        assert_eq!(clean.total_row_count, 7);
+
+        // d1 inactive → include Some, containing only d0's live rows 0..3.
+        let mut only_d0: HashMap<String, i64> = HashMap::new();
+        only_d0.insert("d0".to_string(), 3);
+        let partial = prepare_payload(&meta, &only_d0, &dvs, "ft-0").unwrap();
+        let include = partial.include.expect("inactive source forces include");
+        assert_eq!(include.iter().collect::<Vec<u64>>(), vec![0, 1, 2]);
+    }
+
+    // ---- live_rows subtracts DV positions mapped by offset ----
+    #[test]
+    fn live_rows_subtracts_mapped_deletions() {
+        let ranges = vec![
+            source_range("d0", 0, 3, true),
+            source_range("d1", 3, 4, true),
+        ];
+        let mut bitmap = RoaringBitmap::new();
+        bitmap.insert(2); // d1 local pos 2 → archive 5
+        let mut dvs: HashMap<String, DeletionVector> = HashMap::new();
+        dvs.insert("d1".to_string(), DeletionVector::from_bitmap(bitmap));
+
+        let live = live_rows(&ranges, &dvs).unwrap();
+        assert_eq!(live.iter().collect::<Vec<u64>>(), vec![0, 1, 2, 3, 4, 6]);
+    }
+
+    // ---- live_rows fails loud on an out-of-range deletion position ----
+    #[test]
+    fn live_rows_rejects_out_of_range_deletion() {
+        let ranges = vec![source_range("d0", 0, 3, true)];
+        let mut bitmap = RoaringBitmap::new();
+        bitmap.insert(3); // >= row_count 3
+        let mut dvs: HashMap<String, DeletionVector> = HashMap::new();
+        dvs.insert("d0".to_string(), DeletionVector::from_bitmap(bitmap));
+        assert!(live_rows(&ranges, &dvs).is_err());
+    }
+
+    // ---- limit must be positive ----
+    #[tokio::test]
+    async fn zero_limit_errors() {
+        let file_io = FileIOBuilder::new("memory").build().unwrap();
+        let split = PrimaryKeyFullTextSearchSplit::new(
+            data_split(vec![dfm("d0", 3)]),
+            vec![ft_payload("ft-0", 1, &[("d0", 3)])],
+            Vec::new(),
+        )
+        .unwrap();
+        let dvs: HashMap<String, DeletionVector> = HashMap::new();
+        assert!(search_bucket(
+            &split,
+            r#"{"match":{"query":"x"}}"#,
+            0,
+            &dvs,
+            &file_io,
+            "memory:/x",
+            0
+        )
+        .await
+        .is_err());
+    }
+}
diff --git a/crates/paimon/src/table/pk_full_text_read.rs 
b/crates/paimon/src/table/pk_full_text_read.rs
new file mode 100644
index 00000000..a8dd70a8
--- /dev/null
+++ b/crates/paimon/src/table/pk_full_text_read.rs
@@ -0,0 +1,969 @@
+// 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.
+
+//! Primary-key full-text search read foundation: the score-based candidate 
type,
+//! its cross-bucket score-descending Top-K, and the FAST-only materialized 
read
+//! orchestration (plan the buckets, search each, fuse by score, materialize 
the
+//! winning physical rows best-score-first with a `__paimon_search_score` 
column).
+
+use std::cmp::Ordering;
+use std::collections::{BTreeMap, HashMap};
+
+use arrow_array::RecordBatch;
+use futures::{stream, TryStreamExt};
+
+use crate::deletion_vector::DeletionVector;
+use crate::io::FileIO;
+use crate::spec::BinaryRow;
+use crate::table::data_file_reader::DataFileReader;
+use crate::table::pk_full_text_bucket_search::search_bucket;
+use crate::table::pk_full_text_scan::{PrimaryKeyFullTextScanPlan, 
PrimaryKeyFullTextSearchSplit};
+use crate::table::pk_vector_indexed_split_read::{PkVectorIndexedSplit, 
PkVectorIndexedSplitRead};
+use crate::table::source::DataSplitBuilder;
+use crate::table::vector_search_builder::{
+    collect_ranked_rows, reorder_and_strip_position, RankedRow,
+};
+use crate::table::{ArrowRecordBatchStream, RowRange};
+
+fn data_invalid(message: impl Into<String>) -> crate::Error {
+    crate::Error::DataInvalid {
+        message: message.into(),
+        source: None,
+    }
+}
+
+/// A full-text search hit tagged with its source bucket. `partition`/`bucket` 
are
+/// the cross-bucket merge dimensions carried explicitly (mirroring Java
+/// `PrimaryKeySearchPosition`) so the tie-break does not have to re-derive 
them
+/// from `split_index`; `split_index` is the re-association handle back to the
+/// planned split. `score` is the final full-text relevance score (higher is
+/// better) and is validated finite at construction.
+#[derive(Clone)]
+pub(crate) struct PrimaryKeyFullTextCandidate {
+    pub(crate) split_index: usize,
+    pub(crate) partition: BinaryRow,
+    pub(crate) bucket: i32,
+    pub(crate) score: f32,
+    pub(crate) data_file_name: String,
+    pub(crate) row_position: i64,
+}
+
+impl PrimaryKeyFullTextCandidate {
+    /// Build a candidate, rejecting a non-finite `score` (NaN / ±Infinity).
+    /// A non-finite score means a corrupt ranking (e.g. a malformed archive 
or a
+    /// bogus scorer output) and would poison the score-descending order, so 
fail
+    /// loud rather than emit it. Mirrors Java `PrimaryKeySearchPosition`.
+    pub(crate) fn new(
+        split_index: usize,
+        partition: BinaryRow,
+        bucket: i32,
+        score: f32,
+        data_file_name: String,
+        row_position: i64,
+    ) -> crate::Result<Self> {
+        if !score.is_finite() {
+            return Err(data_invalid(format!(
+                "full-text search score must be finite, got {score} for 
{data_file_name} @ {row_position}"
+            )));
+        }
+        Ok(Self {
+            split_index,
+            partition,
+            bucket,
+            score,
+            data_file_name,
+            row_position,
+        })
+    }
+}
+
+/// 5-level best-first (largest score = best) key. Level 1 orders `score`
+/// DESCENDING via `f32::total_cmp` (comparing `other` against `self`) so 
higher
+/// scores sort first; scores are already validated finite at construction. 
Level
+/// 2 uses the partition's serialized bytes; Rust `Vec<u8>::cmp` is unsigned
+/// lexicographic then shorter-is-less, exactly the spec's contract
+/// (`[0x7f] < [0x80] < [0xff]`). The remaining levels break ties 
deterministically
+/// by bucket, data file name, then physical row position.
+fn candidate_cmp(a: &PrimaryKeyFullTextCandidate, b: 
&PrimaryKeyFullTextCandidate) -> Ordering {
+    b.score
+        .total_cmp(&a.score)
+        .then_with(|| {
+            a.partition
+                .to_serialized_bytes()
+                .cmp(&b.partition.to_serialized_bytes())
+        })
+        .then_with(|| a.bucket.cmp(&b.bucket))
+        .then_with(|| a.data_file_name.cmp(&b.data_file_name))
+        .then_with(|| a.row_position.cmp(&b.row_position))
+}
+
+/// Collect all candidates, order score-descending (best-first) with the
+/// deterministic tie-break, and keep the best `limit`.
+pub(crate) fn top_k_by_score(
+    mut candidates: Vec<PrimaryKeyFullTextCandidate>,
+    limit: usize,
+) -> Vec<PrimaryKeyFullTextCandidate> {
+    candidates.sort_by(candidate_cmp);
+    candidates.truncate(limit);
+    candidates
+}
+
+/// Group best-score survivors into one single-file `PkVectorIndexedSplit` per
+/// `(partition, bucket, data_file)`, re-associating each file to its
+/// `DataFileMeta` + aligned deletion file in the source bucket split. Mirrors 
the
+/// vector `build_indexed_splits`, but carries the raw full-text scores 
verbatim
+/// (`scores: Some(..)`) rather than applying `distance_to_score`: the 
full-text
+/// scores are already final relevance scores, so any transform would corrupt 
them
+/// (spec D1). Groups are emitted in ascending group-key order for a 
deterministic
+/// file/position materialization order; the caller reorders back to 
best-first.
+fn build_full_text_indexed_splits(
+    survivors: Vec<PrimaryKeyFullTextCandidate>,
+    splits: &[PrimaryKeyFullTextSearchSplit],
+) -> crate::Result<Vec<PkVectorIndexedSplit>> {
+    // Group key: (partition bytes, bucket, file name). BTreeMap keeps 
ascending
+    // group order deterministically. Value: (split_index, Vec<(position, 
score)>).
+    type GroupKey = (Vec<u8>, i32, String);
+    let mut groups: BTreeMap<GroupKey, (usize, Vec<(i64, f32)>)> = 
BTreeMap::new();
+    for candidate in survivors {
+        let key = (
+            candidate.partition.to_serialized_bytes(),
+            candidate.bucket,
+            candidate.data_file_name.clone(),
+        );
+        let entry = groups
+            .entry(key)
+            .or_insert_with(|| (candidate.split_index, Vec::new()));
+        // A (partition, bucket, file) group must map to one source split.
+        // Divergent split_index means malformed input (e.g. duplicate 
buckets);
+        // materializing against the wrong split would be a silent wrong-read.
+        if entry.0 != candidate.split_index {
+            return Err(data_invalid(format!(
+                "full-text search hits for {} map to different splits ({} and 
{})",
+                candidate.data_file_name, entry.0, candidate.split_index
+            )));
+        }
+        entry.1.push((candidate.row_position, candidate.score));
+    }
+
+    let mut out = Vec::with_capacity(groups.len());
+    for ((_partition, _bucket, file_name), (split_index, mut hits)) in groups {
+        // Sort positions ascending; reject a duplicate (file, position).
+        hits.sort_by_key(|(pos, _)| *pos);
+        for pair in hits.windows(2) {
+            if pair[0].0 == pair[1].0 {
+                return Err(data_invalid(format!(
+                    "duplicate (file, position) in full-text search result: {} 
@ {}",
+                    file_name, pair[0].0
+                )));
+            }
+        }
+
+        // Re-associate the file to its DataFileMeta + aligned deletion file.
+        let source = &splits
+            .get(split_index)
+            .ok_or_else(|| {
+                data_invalid(format!(
+                    "full-text search hit references split index {split_index} 
out of range (splits: {})",
+                    splits.len()
+                ))
+            })?
+            .data_split;
+        let file_idx = source
+            .data_files()
+            .iter()
+            .position(|f| f.file_name == file_name)
+            .ok_or_else(|| {
+                data_invalid(format!(
+                    "full-text search hit references data file {file_name} not 
present in its bucket split"
+                ))
+            })?;
+        let file_meta = source.data_files()[file_idx].clone();
+        let deletion_file = source
+            .data_deletion_files()
+            .and_then(|dfs| dfs.get(file_idx).cloned().flatten());
+
+        // Every hit's physical position must be in range for its data file.
+        for &(pos, _) in &hits {
+            if pos < 0 || pos >= file_meta.row_count {
+                return Err(data_invalid(format!(
+                    "full-text search position {pos} is outside data file 
{file_name} row count {}",
+                    file_meta.row_count
+                )));
+            }
+        }
+
+        // Coalesce ascending positions into inclusive ranges; scores are the 
raw
+        // full-text scores aligned to ascending-position order.
+        let mut row_ranges: Vec<RowRange> = Vec::new();
+        let mut scores: Vec<f32> = Vec::with_capacity(hits.len());
+        let mut start = hits[0].0;
+        let mut end = hits[0].0;
+        scores.push(hits[0].1);
+        for &(pos, score) in &hits[1..] {
+            if pos == end + 1 {
+                end = pos;
+            } else {
+                row_ranges.push(RowRange::new(start, end));
+                start = pos;
+                end = pos;
+            }
+            scores.push(score);
+        }
+        row_ranges.push(RowRange::new(start, end));
+
+        let mut builder = DataSplitBuilder::new()
+            .with_snapshot(source.snapshot_id())
+            .with_partition(source.partition().clone())
+            .with_bucket(source.bucket())
+            .with_bucket_path(source.bucket_path().to_string())
+            .with_total_buckets(source.total_buckets())
+            .with_data_files(vec![file_meta]);
+        if let Some(df) = deletion_file {
+            builder = builder.with_data_deletion_files(vec![Some(df)]);
+        }
+        let split = builder.build()?;
+
+        out.push(PkVectorIndexedSplit {
+            split,
+            row_ranges,
+            scores: Some(scores),
+        });
+    }
+    Ok(out)
+}
+
+/// FAST-only primary-key full-text materialized read. Given a planned set of
+/// per-bucket search splits, it searches each bucket through the full-text 
archive
+/// reader, fuses the hits cross-bucket by score, materializes the winning 
physical
+/// rows, and re-orders them best-score-first with a `__paimon_search_score` 
column
+/// (the internal `_PKEY_VECTOR_POSITION` column is stripped). Mirrors Java
+/// `PrimaryKeyFullTextRead` feeding its result splits into an ordinary table 
read:
+/// the search decides which rows, the reader decides which columns.
+///
+/// The planning + FAST-mode guard live in the caller 
(`FullTextSearchBuilder`);
+/// this type is the Table-free search+materialize core (it holds only the
+/// predicate-free materialization reader, the `FileIO`, and the table path 
for the
+/// archive directory), so it can be driven end-to-end from a hand-built plan.
+pub(crate) struct PrimaryKeyFullTextRead {
+    file_io: FileIO,
+    materialize_reader: DataFileReader,
+    table_path: String,
+}
+
+impl PrimaryKeyFullTextRead {
+    pub(crate) fn new(
+        file_io: FileIO,
+        materialize_reader: DataFileReader,
+        table_path: String,
+    ) -> Self {
+        Self {
+            file_io,
+            materialize_reader,
+            table_path,
+        }
+    }
+
+    /// Search every planned bucket, fuse the hits by score into a global 
Top-`limit`,
+    /// materialize the winning rows, and emit them best-score-first with the
+    /// unified score column. An empty plan or an empty result yields an empty
+    /// stream. `query` is passed verbatim to the archive reader (spec D2).
+    pub(crate) async fn read(
+        &self,
+        plan: &PrimaryKeyFullTextScanPlan,
+        query: &str,
+        limit: usize,
+    ) -> crate::Result<ArrowRecordBatchStream> {
+        // Every planned split must belong to the plan's resolved snapshot; 
mixing
+        // snapshots would search/materialize physical rows against the wrong
+        // version. The scan already pins one snapshot, so a mismatch is 
malformed
+        // input — fail loud rather than read inconsistently (mirrors Java
+        // `PrimaryKeyFullTextRead`).
+        for split in &plan.splits {
+            if split.data_split.snapshot_id() != plan.snapshot_id {
+                return Err(data_invalid(format!(
+                    "full-text search split snapshot id {} does not match plan 
snapshot {}",
+                    split.data_split.snapshot_id(),
+                    plan.snapshot_id
+                )));
+            }
+        }
+
+        // Per-bucket search -> collected candidates. The bucket DVs are 
loaded from
+        // the same split the materialization later reads (an accepted 
redundancy
+        // between the search and materialization phases, mirroring the vector 
path).
+        let mut candidates: Vec<PrimaryKeyFullTextCandidate> = Vec::new();
+        for (split_index, split) in plan.splits.iter().enumerate() {
+            let dv_factory = self
+                .materialize_reader
+                .build_split_dv_factory(&split.data_split)
+                .await?;
+            let mut dvs: HashMap<String, DeletionVector> = HashMap::new();
+            for file in split.data_split.data_files() {
+                if let Some(dv) =
+                    
DataFileReader::deletion_vector_for_file(dv_factory.as_ref(), &file.file_name)
+                {
+                    dvs.insert(file.file_name.clone(), (*dv).clone());
+                }
+            }
+            let mut hits = search_bucket(
+                split,
+                query,
+                limit,
+                &dvs,
+                &self.file_io,
+                &self.table_path,
+                split_index,
+            )
+            .await?;
+            candidates.append(&mut hits);
+        }
+
+        // Global cross-bucket fusion: score-descending Top-`limit`.
+        let survivors = top_k_by_score(candidates, limit);
+        if survivors.is_empty() {
+            return Ok(Box::pin(stream::empty()));
+        }
+
+        // Rank each survivor by its best-first position, keyed by its FULL 
physical
+        // position `(partition bytes, bucket, file, row position)` so the 
physical
+        // materialization order can be reduced back to best-first (gap 8). 
Keying on
+        // just file+position would collide across partitions/buckets.
+        let mut rank_of: HashMap<(Vec<u8>, i32, String, i64), usize> = 
HashMap::new();
+        for (rank, c) in survivors.iter().enumerate() {
+            rank_of.insert(
+                (
+                    c.partition.to_serialized_bytes(),
+                    c.bucket,
+                    c.data_file_name.clone(),
+                    c.row_position,
+                ),
+                rank,
+            );
+        }
+
+        let indexed_splits = build_full_text_indexed_splits(survivors, 
&plan.splits)?;
+
+        // Materialize every indexed split, retaining each batch and, per row, 
the
+        // (rank, batch_index, row_index) tuple so we can reorder to 
best-first.
+        // Top-K is small, so full in-memory collection is acceptable.
+        let mut batches: Vec<RecordBatch> = Vec::new();
+        let mut ranked: Vec<RankedRow> = Vec::new();
+        for indexed in indexed_splits {
+            let partition_bytes = 
indexed.split.partition().to_serialized_bytes();
+            let bucket = indexed.split.bucket();
+            let file_name = indexed.split.data_files()[0].file_name.clone();
+            let mut stream =
+                
PkVectorIndexedSplitRead::new(self.materialize_reader.clone()).read(&indexed)?;
+            while let Some(batch) = stream.try_next().await? {
+                let batch_index = batches.len();
+                collect_ranked_rows(
+                    &batch,
+                    batch_index,
+                    &partition_bytes,
+                    bucket,
+                    &file_name,
+                    &rank_of,
+                    &mut ranked,
+                )?;
+                batches.push(batch);
+            }
+        }
+
+        // Reorder to best-first and drop the internal position column.
+        let output = reorder_and_strip_position(&batches, ranked)?;
+        Ok(Box::pin(stream::iter(output.into_iter().map(Ok))))
+    }
+}
+
+#[cfg(test)]
+mod tests {
+    use super::*;
+    use crate::spec::BinaryRow;
+
+    fn candidate(
+        split_index: usize,
+        partition_bytes: Vec<u8>,
+        bucket: i32,
+        file: &str,
+        pos: i64,
+        score: f32,
+    ) -> PrimaryKeyFullTextCandidate {
+        PrimaryKeyFullTextCandidate::new(
+            split_index,
+            BinaryRow::from_bytes(1, partition_bytes),
+            bucket,
+            score,
+            file.to_string(),
+            pos,
+        )
+        .unwrap()
+    }
+
+    fn ids(c: &[PrimaryKeyFullTextCandidate]) -> Vec<(i32, String, i64)> {
+        c.iter()
+            .map(|c| (c.bucket, c.data_file_name.clone(), c.row_position))
+            .collect()
+    }
+
+    #[test]
+    fn orders_score_descending() {
+        let part = vec![0x00];
+        let out = top_k_by_score(
+            vec![
+                candidate(0, part.clone(), 0, "f", 0, 1.0),
+                candidate(0, part.clone(), 0, "f", 1, 3.0),
+                candidate(0, part.clone(), 0, "f", 2, 2.0),
+            ],
+            3,
+        );
+        assert_eq!(
+            out.iter().map(|c| c.score).collect::<Vec<f32>>(),
+            vec![3.0, 2.0, 1.0]
+        );
+    }
+
+    #[test]
+    fn tie_break_orders_partition_bytes_as_unsigned() {
+        // Equal score so partition bytes decide: 0x7f < 0x80 < 0xff 
(unsigned).
+        let out = top_k_by_score(
+            vec![
+                candidate(2, vec![0xff], 0, "f", 0, 1.0),
+                candidate(0, vec![0x7f], 0, "f", 0, 1.0),
+                candidate(1, vec![0x80], 0, "f", 0, 1.0),
+            ],
+            3,
+        );
+        assert_eq!(
+            out.iter()
+                .map(|c| c.partition.to_serialized_bytes().pop().unwrap())
+                .collect::<Vec<u8>>(),
+            vec![0x7f, 0x80, 0xff]
+        );
+    }
+
+    #[test]
+    fn tie_break_orders_bucket_then_file_then_position() {
+        // Equal score and equal partition; lower levels decide in order:
+        // bucket asc, then data_file_name asc, then row_position asc.
+        let part = vec![0x00];
+        let out = top_k_by_score(
+            vec![
+                candidate(0, part.clone(), 1, "b", 9, 5.0),
+                candidate(0, part.clone(), 0, "b", 0, 5.0),
+                candidate(0, part.clone(), 0, "a", 7, 5.0),
+                candidate(0, part.clone(), 0, "a", 3, 5.0),
+            ],
+            4,
+        );
+        assert_eq!(
+            ids(&out),
+            vec![
+                (0, "a".to_string(), 3),
+                (0, "a".to_string(), 7),
+                (0, "b".to_string(), 0),
+                (1, "b".to_string(), 9),
+            ]
+        );
+    }
+
+    #[test]
+    fn truncates_to_limit() {
+        let part = vec![0x00];
+        let out = top_k_by_score(
+            vec![
+                candidate(0, part.clone(), 0, "f", 0, 1.0),
+                candidate(0, part.clone(), 0, "f", 1, 3.0),
+                candidate(0, part.clone(), 0, "f", 2, 2.0),
+            ],
+            1,
+        );
+        // Only the single best score survives.
+        assert_eq!(ids(&out), vec![(0, "f".to_string(), 1)]);
+        assert_eq!(out[0].score, 3.0);
+    }
+
+    #[test]
+    fn empty_candidates_yield_empty() {
+        assert!(top_k_by_score(Vec::new(), 5).is_empty());
+    }
+
+    #[test]
+    fn new_rejects_non_finite_scores() {
+        let mk = |score: f32| {
+            PrimaryKeyFullTextCandidate::new(0, BinaryRow::new(0), 0, score, 
"f".to_string(), 0)
+        };
+        assert!(mk(f32::NAN).is_err());
+        assert!(mk(f32::INFINITY).is_err());
+        assert!(mk(f32::NEG_INFINITY).is_err());
+        // A finite score is accepted.
+        assert!(mk(1.5).is_ok());
+    }
+}
+
+#[cfg(test)]
+mod read_tests {
+    use super::*;
+    use crate::arrow::build_target_arrow_schema;
+    use crate::io::{FileIO, FileIOBuilder};
+    use crate::spec::stats::BinaryTableStats;
+    use crate::spec::{DataField, DataFileMeta, DataType, GlobalIndexMeta, 
IndexFileMeta, IntType};
+    use crate::table::pk_full_text_bucket_state::PK_FULL_TEXT_INDEX_TYPE;
+    use crate::table::pk_full_text_scan::PrimaryKeyFullTextSearchSplit;
+    use crate::table::pk_vector_position_read::{PKEY_VECTOR_POSITION_COLUMN, 
SEARCH_SCORE_COLUMN};
+    use crate::table::schema_manager::SchemaManager;
+    use crate::table::source::{DataSplit, DataSplitBuilder, DeletionFile};
+    use arrow_array::{Array, Float32Array, Int32Array, RecordBatch};
+    use bytes::Bytes;
+    use paimon_ftindex_core::io::PosWriter;
+    use paimon_ftindex_core::{FullTextIndexConfig, FullTextIndexWriter};
+    use paimon_mosaic_core::spec::COMPRESSION_NONE;
+    use paimon_mosaic_core::writer::{MosaicWriter, OutputFile, WriterOptions};
+    use std::io;
+    use std::sync::Arc;
+
+    const FILE_SOURCE_COMPACT: i32 = 1;
+
+    // ---- helpers ----
+
+    struct MemOutputFile {
+        data: Vec<u8>,
+    }
+    impl MemOutputFile {
+        fn new() -> Self {
+            Self { data: Vec::new() }
+        }
+    }
+    impl OutputFile for MemOutputFile {
+        fn write(&mut self, data: &[u8]) -> io::Result<()> {
+            self.data.extend_from_slice(data);
+            Ok(())
+        }
+        fn flush(&mut self) -> io::Result<()> {
+            Ok(())
+        }
+        fn pos(&self) -> u64 {
+            self.data.len() as u64
+        }
+    }
+
+    fn id_field() -> DataField {
+        DataField::new(0, "id".to_string(), DataType::Int(IntType::new()))
+    }
+    fn id_fields() -> Vec<DataField> {
+        vec![id_field()]
+    }
+    fn id_batch(ids: Vec<i32>) -> RecordBatch {
+        let schema = build_target_arrow_schema(&id_fields()).unwrap();
+        RecordBatch::try_new(schema, 
vec![Arc::new(Int32Array::from(ids))]).unwrap()
+    }
+
+    fn write_mosaic(batch: &RecordBatch) -> Bytes {
+        let out = MemOutputFile::new();
+        let mut writer = MosaicWriter::new(
+            out,
+            batch.schema().as_ref(),
+            WriterOptions {
+                compression: COMPRESSION_NONE,
+                num_buckets: 2,
+                row_group_max_size: u64::MAX,
+                ..Default::default()
+            },
+        )
+        .unwrap();
+        writer.write_batch(batch).unwrap();
+        writer.close().unwrap();
+        Bytes::from(writer.output().data.to_vec())
+    }
+
+    /// Build a tiny full-text archive in memory via the core writer.
+    fn build_archive(docs: &[(i64, &str)]) -> Vec<u8> {
+        let mut writer = 
FullTextIndexWriter::new(FullTextIndexConfig::new()).unwrap();
+        for (row_id, text) in docs {
+            writer.add_document(*row_id, (*text).to_string()).unwrap();
+        }
+        let mut out = PosWriter::new(Vec::<u8>::new());
+        writer.write(&mut out).unwrap();
+        out.into_inner()
+    }
+
+    async fn write_bytes(file_io: &FileIO, path: &str, bytes: Vec<u8>) {
+        file_io
+            .new_output(path)
+            .unwrap()
+            .write(Bytes::from(bytes))
+            .await
+            .unwrap();
+    }
+
+    /// One Java `DataOutput#writeUTF` value (u16-BE length + modified UTF-8).
+    fn java_write_utf(s: &str) -> Vec<u8> {
+        let mut body = Vec::new();
+        for c in s.encode_utf16() {
+            if (0x0001..=0x007F).contains(&c) {
+                body.push(c as u8);
+            } else if c > 0x07FF {
+                body.push(0xE0 | (c >> 12) as u8);
+                body.push(0x80 | ((c >> 6) & 0x3F) as u8);
+                body.push(0x80 | (c & 0x3F) as u8);
+            } else {
+                body.push(0xC0 | (c >> 6) as u8);
+                body.push(0x80 | (c & 0x3F) as u8);
+            }
+        }
+        let mut out = (body.len() as u16).to_be_bytes().to_vec();
+        out.extend_from_slice(&body);
+        out
+    }
+
+    /// A `_SOURCE_META` frame (version 1) for a level and its ordered sources.
+    fn frame(data_level: i32, files: &[(&str, i64)]) -> Vec<u8> {
+        let mut out = Vec::new();
+        out.extend_from_slice(&1i32.to_be_bytes());
+        out.extend_from_slice(&data_level.to_be_bytes());
+        out.extend_from_slice(&(files.len() as i32).to_be_bytes());
+        for (name, rows) in files {
+            out.extend_from_slice(&java_write_utf(name));
+            out.extend_from_slice(&rows.to_be_bytes());
+        }
+        out
+    }
+
+    fn ft_payload(file_name: &str, files: &[(&str, i64)]) -> IndexFileMeta {
+        let total: i64 = files.iter().map(|(_, r)| *r).sum();
+        IndexFileMeta {
+            index_type: PK_FULL_TEXT_INDEX_TYPE.into(),
+            file_name: file_name.into(),
+            file_size: 1,
+            row_count: total as i32,
+            deletion_vectors_ranges: None,
+            global_index_meta: Some(GlobalIndexMeta {
+                row_range_start: 0,
+                row_range_end: 0,
+                index_field_id: 1,
+                extra_field_ids: None,
+                index_meta: None,
+                source_meta: Some(frame(1, files)),
+            }),
+        }
+    }
+
+    fn data_file(file_name: &str, file_size: i64, row_count: i64) -> 
DataFileMeta {
+        DataFileMeta {
+            file_name: file_name.into(),
+            file_size,
+            row_count,
+            min_key: Vec::new(),
+            max_key: Vec::new(),
+            key_stats: BinaryTableStats::empty(),
+            value_stats: BinaryTableStats::empty(),
+            min_sequence_number: 0,
+            max_sequence_number: 0,
+            schema_id: 1,
+            level: 1,
+            extra_files: Vec::new(),
+            creation_time: None,
+            delete_row_count: None,
+            embedded_index: None,
+            file_source: Some(FILE_SOURCE_COMPACT),
+            value_stats_cols: None,
+            external_path: None,
+            first_row_id: Some(0),
+            write_cols: None,
+        }
+    }
+
+    async fn write_deletion_file(
+        file_io: &FileIO,
+        path: &str,
+        deleted_rows: &[u32],
+    ) -> DeletionFile {
+        const MAGIC_NUMBER: i32 = 1581511376;
+        let mut bitmap = roaring::RoaringBitmap::new();
+        for row in deleted_rows {
+            bitmap.insert(*row);
+        }
+        let mut bitmap_bytes = Vec::new();
+        bitmap.serialize_into(&mut bitmap_bytes).unwrap();
+        let bitmap_length = 4 + bitmap_bytes.len() as i32;
+        let mut blob = Vec::new();
+        blob.extend_from_slice(&bitmap_length.to_be_bytes());
+        blob.extend_from_slice(&MAGIC_NUMBER.to_be_bytes());
+        blob.extend_from_slice(&bitmap_bytes);
+        blob.extend_from_slice(&0i32.to_be_bytes());
+        file_io
+            .new_output(path)
+            .unwrap()
+            .write(Bytes::from(blob))
+            .await
+            .unwrap();
+        DeletionFile::new(
+            path.to_string(),
+            0,
+            bitmap_length as i64,
+            Some(deleted_rows.len() as i64),
+        )
+    }
+
+    /// Build the physical mosaic + a single-file `DataSplit`, plus a 
predicate-free
+    /// `DataFileReader` projecting `id`. `deleted_rows`, when non-empty, 
attaches a DV.
+    async fn build_data(
+        file_io: &FileIO,
+        table_path: &str,
+        ids: Vec<i32>,
+        deleted_rows: &[u32],
+    ) -> (DataFileReader, DataSplit) {
+        let bucket_path = format!("{table_path}/bucket-0");
+        let file_name = "d0.mosaic";
+        let row_count = ids.len() as i64;
+        let data = write_mosaic(&id_batch(ids));
+        write_bytes(
+            file_io,
+            &format!("{bucket_path}/{file_name}"),
+            data.to_vec(),
+        )
+        .await;
+
+        let mut split_builder = DataSplitBuilder::new()
+            .with_snapshot(1)
+            .with_partition(BinaryRow::new(0))
+            .with_bucket(0)
+            .with_bucket_path(bucket_path)
+            .with_total_buckets(1)
+            .with_data_files(vec![data_file(file_name, data.len() as i64, 
row_count)]);
+        if !deleted_rows.is_empty() {
+            let df =
+                write_deletion_file(file_io, 
&format!("{table_path}/index/dv-0"), deleted_rows)
+                    .await;
+            split_builder = 
split_builder.with_data_deletion_files(vec![Some(df)]);
+        }
+        let split = split_builder.build().unwrap();
+
+        let schema_manager = SchemaManager::new(file_io.clone(), 
table_path.to_string());
+        let reader = DataFileReader::new(
+            file_io.clone(),
+            schema_manager,
+            1,
+            id_fields(),
+            id_fields(),
+            Vec::new(),
+        );
+        (reader, split)
+    }
+
+    fn column_i32(batch: &RecordBatch, name: &str) -> Vec<i32> {
+        let idx = batch.schema().index_of(name).unwrap();
+        batch
+            .column(idx)
+            .as_any()
+            .downcast_ref::<Int32Array>()
+            .unwrap()
+            .values()
+            .to_vec()
+    }
+
+    async fn collect(stream: ArrowRecordBatchStream) -> Vec<RecordBatch> {
+        stream.try_collect::<Vec<_>>().await.unwrap()
+    }
+
+    // ---- build_full_text_indexed_splits: raw scores, no distance transform 
----
+    #[test]
+    fn build_splits_carries_raw_scores_aligned_by_position() {
+        let split = PrimaryKeyFullTextSearchSplit::new(
+            DataSplitBuilder::new()
+                .with_snapshot(1)
+                .with_partition(BinaryRow::new(0))
+                .with_bucket(0)
+                .with_bucket_path("memory:/t/bucket-0".to_string())
+                .with_total_buckets(1)
+                .with_data_files(vec![data_file("d0.mosaic", 1, 4)])
+                .build()
+                .unwrap(),
+            vec![ft_payload("ft-0", &[("d0.mosaic", 4)])],
+            Vec::new(),
+        )
+        .unwrap();
+
+        // Out-of-order candidates; scores must survive verbatim, aligned to
+        // ascending position order (no distance_to_score transform).
+        let cands = vec![
+            PrimaryKeyFullTextCandidate::new(0, BinaryRow::new(0), 0, 0.1, 
"d0.mosaic".into(), 3)
+                .unwrap(),
+            PrimaryKeyFullTextCandidate::new(0, BinaryRow::new(0), 0, 0.9, 
"d0.mosaic".into(), 0)
+                .unwrap(),
+            PrimaryKeyFullTextCandidate::new(0, BinaryRow::new(0), 0, 0.5, 
"d0.mosaic".into(), 2)
+                .unwrap(),
+        ];
+        let out = build_full_text_indexed_splits(cands, 
std::slice::from_ref(&split)).unwrap();
+        assert_eq!(out.len(), 1);
+        // positions 0,2,3 -> ranges [0,0],[2,3]; scores aligned ascending: 
0.9,0.5,0.1.
+        assert_eq!(
+            out[0].row_ranges,
+            vec![RowRange::new(0, 0), RowRange::new(2, 3)]
+        );
+        assert_eq!(out[0].scores, Some(vec![0.9, 0.5, 0.1]));
+    }
+
+    #[test]
+    fn build_splits_rejects_duplicate_position() {
+        let split = PrimaryKeyFullTextSearchSplit::new(
+            DataSplitBuilder::new()
+                .with_snapshot(1)
+                .with_partition(BinaryRow::new(0))
+                .with_bucket(0)
+                .with_bucket_path("memory:/t/bucket-0".to_string())
+                .with_total_buckets(1)
+                .with_data_files(vec![data_file("d0.mosaic", 1, 4)])
+                .build()
+                .unwrap(),
+            vec![ft_payload("ft-0", &[("d0.mosaic", 4)])],
+            Vec::new(),
+        )
+        .unwrap();
+        let cands = vec![
+            PrimaryKeyFullTextCandidate::new(0, BinaryRow::new(0), 0, 0.9, 
"d0.mosaic".into(), 1)
+                .unwrap(),
+            PrimaryKeyFullTextCandidate::new(0, BinaryRow::new(0), 0, 0.5, 
"d0.mosaic".into(), 1)
+                .unwrap(),
+        ];
+        assert!(build_full_text_indexed_splits(cands, 
std::slice::from_ref(&split)).is_err());
+    }
+
+    // ---- (a) round-trip: archive -> search -> materialize -> 
best-score-first ----
+    #[tokio::test]
+    async fn execute_read_materializes_rows_best_score_first() {
+        let file_io = FileIOBuilder::new("memory").build().unwrap();
+        let table_path = "memory:/pk_ft_roundtrip";
+        // pos2 ("alpha alpha alpha") scores higher for "alpha" than pos0 
("alpha").
+        let archive = build_archive(&[
+            (0, "alpha"),
+            (1, "beta"),
+            (2, "alpha alpha alpha"),
+            (3, "gamma"),
+        ]);
+        write_bytes(&file_io, &format!("{table_path}/index/ft-0"), 
archive).await;
+
+        let (reader, split) = build_data(&file_io, table_path, vec![100, 101, 
102, 103], &[]).await;
+        let plan = PrimaryKeyFullTextScanPlan {
+            snapshot_id: 1,
+            splits: vec![PrimaryKeyFullTextSearchSplit::new(
+                split,
+                vec![ft_payload("ft-0", &[("d0.mosaic", 4)])],
+                Vec::new(),
+            )
+            .unwrap()],
+        };
+
+        let read = PrimaryKeyFullTextRead::new(file_io.clone(), reader, 
table_path.to_string());
+        let batches = collect(
+            read.read(&plan, r#"{"match":{"query":"alpha"}}"#, 10)
+                .await
+                .unwrap(),
+        )
+        .await;
+
+        // Best-score-first: pos2 (id 102, higher tf) then pos0 (id 100).
+        assert_eq!(column_i32(&batches[0], "id"), vec![102, 100]);
+        // The unified score column is present; the internal position column 
is not.
+        assert!(batches[0].schema().index_of(SEARCH_SCORE_COLUMN).is_ok());
+        assert!(batches[0]
+            .schema()
+            .index_of(PKEY_VECTOR_POSITION_COLUMN)
+            .is_err());
+        // Scores are descending (best first) and finite.
+        let scores: Vec<f32> = batches
+            .iter()
+            .flat_map(|b| {
+                let idx = b.schema().index_of(SEARCH_SCORE_COLUMN).unwrap();
+                b.column(idx)
+                    .as_any()
+                    .downcast_ref::<Float32Array>()
+                    .unwrap()
+                    .values()
+                    .to_vec()
+            })
+            .collect();
+        assert!(
+            scores[0] >= scores[1],
+            "scores must be best-first: {scores:?}"
+        );
+    }
+
+    // ---- (b) a DV-deleted row is absent from the materialized result ----
+    #[tokio::test]
+    async fn execute_read_excludes_dv_deleted_row() {
+        let file_io = FileIOBuilder::new("memory").build().unwrap();
+        let table_path = "memory:/pk_ft_dv";
+        let archive = build_archive(&[
+            (0, "alpha"),
+            (1, "beta"),
+            (2, "alpha alpha alpha"),
+            (3, "gamma"),
+        ]);
+        write_bytes(&file_io, &format!("{table_path}/index/ft-0"), 
archive).await;
+
+        // Delete physical position 2 (the strong "alpha" hit) -> only pos0 
survives.
+        let (reader, split) =
+            build_data(&file_io, table_path, vec![100, 101, 102, 103], 
&[2]).await;
+        let plan = PrimaryKeyFullTextScanPlan {
+            snapshot_id: 1,
+            splits: vec![PrimaryKeyFullTextSearchSplit::new(
+                split,
+                vec![ft_payload("ft-0", &[("d0.mosaic", 4)])],
+                Vec::new(),
+            )
+            .unwrap()],
+        };
+
+        let read = PrimaryKeyFullTextRead::new(file_io.clone(), reader, 
table_path.to_string());
+        let batches = collect(
+            read.read(&plan, r#"{"match":{"query":"alpha"}}"#, 10)
+                .await
+                .unwrap(),
+        )
+        .await;
+
+        let ids: Vec<i32> = batches.iter().flat_map(|b| column_i32(b, 
"id")).collect();
+        assert_eq!(ids, vec![100], "DV-deleted row (id 102) must be absent");
+    }
+
+    // ---- empty result -> empty stream (no candidates) ----
+    #[tokio::test]
+    async fn execute_read_no_match_yields_empty_stream() {
+        let file_io = FileIOBuilder::new("memory").build().unwrap();
+        let table_path = "memory:/pk_ft_empty";
+        let archive = build_archive(&[(0, "alpha"), (1, "beta")]);
+        write_bytes(&file_io, &format!("{table_path}/index/ft-0"), 
archive).await;
+        let (reader, split) = build_data(&file_io, table_path, vec![100, 101], 
&[]).await;
+        let plan = PrimaryKeyFullTextScanPlan {
+            snapshot_id: 1,
+            splits: vec![PrimaryKeyFullTextSearchSplit::new(
+                split,
+                vec![ft_payload("ft-0", &[("d0.mosaic", 2)])],
+                Vec::new(),
+            )
+            .unwrap()],
+        };
+        let read = PrimaryKeyFullTextRead::new(file_io.clone(), reader, 
table_path.to_string());
+        let batches = collect(
+            read.read(&plan, r#"{"match":{"query":"zeta"}}"#, 10)
+                .await
+                .unwrap(),
+        )
+        .await;
+        assert!(batches.is_empty(), "no match must yield an empty stream");
+    }
+}
diff --git a/crates/paimon/src/table/pk_full_text_scan.rs 
b/crates/paimon/src/table/pk_full_text_scan.rs
new file mode 100644
index 00000000..081ece10
--- /dev/null
+++ b/crates/paimon/src/table/pk_full_text_scan.rs
@@ -0,0 +1,711 @@
+// 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.
+
+//! Primary-key full-text search planning: resolve a snapshot, plan its data
+//! splits, scan the index manifest for this column's full-text payloads, and
+//! accumulate one search split per bucket. Mirror of Java
+//! `PrimaryKeyFullTextScan` and `PrimaryKeyFullTextSearchSplit`.
+
+use std::collections::{BTreeMap, HashSet};
+
+use crate::spec::{
+    should_read_pk_index_source, BinaryRow, DataFileMeta, FileKind, 
IndexFileMeta, IndexManifest,
+    Predicate, PrimaryKeyIndexSourceMeta,
+};
+use crate::table::pk_full_text_bucket_state::{PkFullTextBucketState, 
PK_FULL_TEXT_INDEX_TYPE};
+use crate::table::source::{DataSplit, DataSplitBuilder, DeletionFile};
+use crate::table::Table;
+
+fn data_invalid(message: impl Into<String>) -> crate::Error {
+    crate::Error::DataInvalid {
+        message: message.into(),
+        source: None,
+    }
+}
+
+/// Compaction-visible data files and full-text payloads for one snapshot 
bucket.
+///
+/// Mirror of Java `PrimaryKeyFullTextSearchSplit`. The constructor enforces 
the
+/// bucket-split invariants (unique data files; each payload covers ≥1 active
+/// source; no active file double-covered; uncovered files 
active/non-duplicate/
+/// not-covered; and covered ∪ uncovered == every active data file in the 
split).
+pub(crate) struct PrimaryKeyFullTextSearchSplit {
+    pub data_split: DataSplit,
+    pub current_payloads: Vec<IndexFileMeta>,
+    // Carried for parity with Java `PrimaryKeyFullTextSearchSplit` (which 
preserves
+    // the uncovered files) and validated by the split constructor's 
completeness
+    // invariant. FAST mode never searches these files, so nothing reads the 
field
+    // after construction.
+    #[allow(dead_code)]
+    pub uncovered_data_files: Vec<String>,
+}
+
+impl PrimaryKeyFullTextSearchSplit {
+    pub(crate) fn new(
+        data_split: DataSplit,
+        current_payloads: Vec<IndexFileMeta>,
+        uncovered_data_files: Vec<String>,
+    ) -> crate::Result<Self> {
+        // Unique data files in the split.
+        let mut source_files: HashSet<String> = HashSet::new();
+        for data_file in data_split.data_files() {
+            if !source_files.insert(data_file.file_name.clone()) {
+                return Err(data_invalid(format!(
+                    "data file {} appears more than once in a full-text bucket 
split",
+                    data_file.file_name
+                )));
+            }
+        }
+
+        // Each current payload must cover ≥1 active source, and no active 
source
+        // may be covered by more than one payload.
+        let mut covered: HashSet<String> = HashSet::new();
+        for payload in &current_payloads {
+            let gim = payload.global_index_meta.as_ref().ok_or_else(|| {
+                data_invalid(format!(
+                    "full-text payload {} has no global index meta",
+                    payload.file_name
+                ))
+            })?;
+            let source_meta = 
PrimaryKeyIndexSourceMeta::from_global_index_meta(gim)?;
+            let mut covers_active_source = false;
+            for source in source_meta.source_files() {
+                if !source_files.contains(source.file_name()) {
+                    continue;
+                }
+                covers_active_source = true;
+                if !covered.insert(source.file_name().to_string()) {
+                    return Err(data_invalid(format!(
+                        "data file {} is covered by more than one full-text 
payload",
+                        source.file_name()
+                    )));
+                }
+            }
+            if !covers_active_source {
+                return Err(data_invalid(format!(
+                    "full-text payload {} does not cover an active data file 
in its bucket split",
+                    payload.file_name
+                )));
+            }
+        }
+
+        // Uncovered files must be active, non-duplicate, and not covered.
+        let mut uncovered: HashSet<String> = HashSet::new();
+        for source in &uncovered_data_files {
+            if !source_files.contains(source) {
+                return Err(data_invalid(format!(
+                    "uncovered full-text data file {source} is outside its 
bucket split"
+                )));
+            }
+            if !uncovered.insert(source.clone()) {
+                return Err(data_invalid(format!(
+                    "uncovered full-text data file {source} appears more than 
once"
+                )));
+            }
+            if covered.contains(source) {
+                return Err(data_invalid(format!(
+                    "data file {source} cannot be both indexed and uncovered"
+                )));
+            }
+        }
+
+        // Completeness: covered ∪ uncovered == every active data file in the 
split.
+        if covered.len() + uncovered.len() != source_files.len() {
+            return Err(data_invalid(
+                "every full-text source file must be indexed or explicitly 
uncovered",
+            ));
+        }
+
+        Ok(Self {
+            data_split,
+            current_payloads,
+            uncovered_data_files,
+        })
+    }
+}
+
+/// One bucket's should-read data files combined into a single split, keeping 
data
+/// files and deletion files in strict parallel order and rejecting duplicate 
file
+/// names. Mirrors Java `PrimaryKeyFullTextScan.BucketAccumulator`: only files 
that
+/// pass `should_read_pk_index_source` are retained.
+struct BucketAccumulator {
+    snapshot_id: i64,
+    partition: BinaryRow,
+    bucket: i32,
+    bucket_path: Option<String>,
+    total_buckets: Option<i32>,
+    data_files: Vec<DataFileMeta>,
+    deletion_files: Vec<Option<DeletionFile>>,
+    seen: HashSet<String>,
+    any_deletion: bool,
+}
+
+impl BucketAccumulator {
+    fn new(snapshot_id: i64, partition: BinaryRow, bucket: i32) -> Self {
+        Self {
+            snapshot_id,
+            partition,
+            bucket,
+            bucket_path: None,
+            total_buckets: None,
+            data_files: Vec::new(),
+            deletion_files: Vec::new(),
+            seen: HashSet::new(),
+            any_deletion: false,
+        }
+    }
+
+    fn add(&mut self, split: &DataSplit) -> crate::Result<()> {
+        if split.snapshot_id() != self.snapshot_id {
+            return Err(data_invalid(
+                "data split snapshot id does not match plan snapshot",
+            ));
+        }
+        if split.partition().to_serialized_bytes() != 
self.partition.to_serialized_bytes() {
+            return Err(data_invalid(
+                "data split partition does not match bucket group",
+            ));
+        }
+        if split.bucket() != self.bucket {
+            return Err(data_invalid(
+                "data split bucket does not match bucket group",
+            ));
+        }
+        match &self.bucket_path {
+            Some(p) if p != split.bucket_path() => {
+                return Err(data_invalid("inconsistent bucket path within 
bucket group"))
+            }
+            None => self.bucket_path = Some(split.bucket_path().to_string()),
+            _ => {}
+        }
+        match self.total_buckets {
+            Some(tb) if tb != split.total_buckets() => {
+                return Err(data_invalid(
+                    "inconsistent total buckets within bucket group",
+                ))
+            }
+            None => self.total_buckets = Some(split.total_buckets()),
+            _ => {}
+        }
+        let dvs = split.data_deletion_files();
+        for (i, file) in split.data_files().iter().enumerate() {
+            // Only compaction-visible sources take part in full-text planning,
+            // mirroring Java's `PrimaryKeyIndexSourcePolicy.shouldRead` 
filter.
+            if !should_read_pk_index_source(file) {
+                continue;
+            }
+            if !self.seen.insert(file.file_name.clone()) {
+                return Err(data_invalid(format!(
+                    "duplicate data file in bucket group: {}",
+                    file.file_name
+                )));
+            }
+            self.data_files.push(file.clone());
+            let df = dvs.and_then(|d| d.get(i).cloned().flatten());
+            if df.is_some() {
+                self.any_deletion = true;
+            }
+            self.deletion_files.push(df);
+        }
+        Ok(())
+    }
+
+    fn is_empty(&self) -> bool {
+        self.data_files.is_empty()
+    }
+
+    fn build(self) -> crate::Result<DataSplit> {
+        let mut builder = DataSplitBuilder::new()
+            .with_snapshot(self.snapshot_id)
+            .with_partition(self.partition)
+            .with_bucket(self.bucket)
+            .with_bucket_path(
+                self.bucket_path
+                    .ok_or_else(|| data_invalid("bucket group has no bucket 
path"))?,
+            )
+            .with_total_buckets(self.total_buckets.unwrap_or(1))
+            .with_data_files(self.data_files)
+            .with_raw_convertible(false);
+        if self.any_deletion {
+            builder = builder.with_data_deletion_files(self.deletion_files);
+        }
+        builder.build()
+    }
+}
+
+/// The per-bucket search splits produced by planning.
+pub(crate) struct PrimaryKeyFullTextScanPlan {
+    // The snapshot the plan was resolved against; the read guards every split
+    // against it before searching.
+    pub snapshot_id: i64,
+    pub splits: Vec<PrimaryKeyFullTextSearchSplit>,
+}
+
+pub(crate) struct PrimaryKeyFullTextScan<'a> {
+    table: &'a Table,
+    text_field_id: i32,
+    filter: Option<Predicate>,
+}
+
+impl<'a> PrimaryKeyFullTextScan<'a> {
+    pub(crate) fn new(table: &'a Table, text_field_id: i32, filter: 
Option<Predicate>) -> Self {
+        Self {
+            table,
+            text_field_id,
+            filter,
+        }
+    }
+
+    pub(crate) async fn plan(&self) -> 
crate::Result<PrimaryKeyFullTextScanPlan> {
+        let snapshot_manager = self.table.snapshot_manager();
+
+        // Data splits first, via the table's own scan resolution (time travel 
/
+        // scan.snapshot-id aware), then derive the snapshot from the scan 
output so
+        // the index manifest and the data splits stay on ONE snapshot (mirror 
Java
+        // `PrimaryKeyFullTextScan`, which resolves a single snapshot up 
front). The
+        // residual filter is pushed into the read builder so scan planning 
drops
+        // files whose stats cannot match the predicate.
+        let mut read_builder = self.table.new_read_builder();
+        if let Some(filter) = &self.filter {
+            read_builder.with_filter(filter.clone());
+        }
+        let data_splits = read_builder
+            .new_scan()
+            .with_scan_all_files()
+            .plan()
+            .await?
+            .splits()
+            .to_vec();
+
+        let Some(first_split) = data_splits.first() else {
+            return Ok(PrimaryKeyFullTextScanPlan {
+                snapshot_id: 0,
+                splits: Vec::new(),
+            });
+        };
+        let snapshot_id = first_split.snapshot_id();
+        let snapshot = snapshot_manager.get_snapshot(snapshot_id).await?;
+
+        let mut entries: Vec<(BinaryRow, i32, IndexFileMeta)> = Vec::new();
+        if let Some(name) = snapshot.index_manifest() {
+            let path = snapshot_manager.manifest_path(name);
+            for entry in IndexManifest::read(self.table.file_io(), 
&path).await? {
+                // The on-disk index manifest is combined to live ADD entries 
only.
+                if entry.kind != FileKind::Add {
+                    return Err(data_invalid(format!(
+                        "index manifest entry {} is not active (kind {:?})",
+                        entry.index_file.file_name, entry.kind
+                    )));
+                }
+                // Planner guards (mirror Java `matchesDefinition`): full-text 
index
+                // type, global index meta + source meta present, matching 
field id.
+                if entry.index_file.index_type != PK_FULL_TEXT_INDEX_TYPE {
+                    continue;
+                }
+                let Some(gim) = entry.index_file.global_index_meta.as_ref() 
else {
+                    continue;
+                };
+                if gim.source_meta.is_none() {
+                    continue;
+                }
+                if gim.index_field_id != self.text_field_id {
+                    continue;
+                }
+                let partition = 
BinaryRow::from_serialized_bytes(&entry.partition)?;
+                entries.push((partition, entry.bucket, 
entry.index_file.clone()));
+            }
+        }
+
+        let splits = plan_from_inputs(snapshot_id, data_splits, entries, 
self.text_field_id)?;
+        Ok(PrimaryKeyFullTextScanPlan {
+            snapshot_id,
+            splits,
+        })
+    }
+}
+
+/// Pure planning core, drivable without a live snapshot: group full-text 
payloads
+/// and data splits by `(partition, bucket)`, then assemble one search split 
per
+/// bucket that has active data. Index-only buckets are dropped, not errored.
+fn plan_from_inputs(
+    snapshot_id: i64,
+    data_splits: Vec<DataSplit>,
+    index_entries: Vec<(BinaryRow, i32, IndexFileMeta)>,
+    text_field_id: i32,
+) -> crate::Result<Vec<PrimaryKeyFullTextSearchSplit>> {
+    type Key = (Vec<u8>, i32);
+
+    // Phase A: group full-text payloads by (partition, bucket).
+    let mut payloads_by_bucket: BTreeMap<Key, Vec<IndexFileMeta>> = 
BTreeMap::new();
+    for (partition, bucket, payload) in index_entries {
+        let key = (partition.to_serialized_bytes(), bucket);
+        payloads_by_bucket.entry(key).or_default().push(payload);
+    }
+
+    // Phase B: group data splits by (partition, bucket).
+    let mut accum_by_bucket: BTreeMap<Key, BucketAccumulator> = 
BTreeMap::new();
+    for split in &data_splits {
+        // Skip negative-bucket data splits before grouping, mirroring Java
+        // `PrimaryKeyFullTextScan.plan`: only real buckets take part in
+        // full-text search planning.
+        if split.bucket() < 0 {
+            continue;
+        }
+        let key = (split.partition().to_serialized_bytes(), split.bucket());
+        let acc = accum_by_bucket.entry(key).or_insert_with(|| {
+            BucketAccumulator::new(snapshot_id, split.partition().clone(), 
split.bucket())
+        });
+        acc.add(split)?;
+    }
+
+    // Phase C: assemble one split per bucket that has active data.
+    let mut out = Vec::new();
+    for (key, acc) in accum_by_bucket {
+        if acc.is_empty() {
+            continue;
+        }
+        let payloads = payloads_by_bucket.remove(&key).unwrap_or_default();
+        let data_split = acc.build()?;
+        let active_data_files = data_split.data_files().to_vec();
+        let state = PkFullTextBucketState::from_active_data_files(
+            text_field_id,
+            &active_data_files,
+            payloads,
+        )?;
+
+        // Covered = active source files a current payload maps. Since the 
state is
+        // fed exactly the split's active files, this is the full current 
mapping;
+        // the active-name guard mirrors Java's defensive 
`activeSources.contains`.
+        let active_names: HashSet<&str> = active_data_files
+            .iter()
+            .map(|f| f.file_name.as_str())
+            .collect();
+        let covered: HashSet<String> = state
+            .payload_by_source_file()
+            .keys()
+            .filter(|source| active_names.contains(source.as_str()))
+            .cloned()
+            .collect();
+        let uncovered: Vec<String> = active_data_files
+            .iter()
+            .filter(|f| !covered.contains(&f.file_name))
+            .map(|f| f.file_name.clone())
+            .collect();
+
+        out.push(PrimaryKeyFullTextSearchSplit::new(
+            data_split,
+            state.current_payloads().to_vec(),
+            uncovered,
+        )?);
+    }
+    // Index-only buckets left in payloads_by_bucket are intentionally dropped.
+    Ok(out)
+}
+
+#[cfg(test)]
+mod tests {
+    use super::*;
+    use crate::spec::stats::BinaryTableStats;
+    use crate::spec::GlobalIndexMeta;
+
+    /// COMPACT file source discriminant (matches Java `FileSource.COMPACT`).
+    const FILE_SOURCE_COMPACT: i32 = 1;
+
+    /// A COMPACT data file at `level` (>0 to be a live source by default).
+    fn dfm(name: &str, rows: i64, level: i32) -> DataFileMeta {
+        DataFileMeta {
+            file_name: name.into(),
+            file_size: 1,
+            row_count: rows,
+            min_key: Vec::new(),
+            max_key: Vec::new(),
+            key_stats: BinaryTableStats::empty(),
+            value_stats: BinaryTableStats::empty(),
+            min_sequence_number: 0,
+            max_sequence_number: 0,
+            schema_id: 1,
+            level,
+            extra_files: Vec::new(),
+            creation_time: None,
+            delete_row_count: None,
+            embedded_index: None,
+            file_source: Some(FILE_SOURCE_COMPACT),
+            value_stats_cols: None,
+            external_path: None,
+            first_row_id: Some(0),
+            write_cols: None,
+        }
+    }
+
+    /// One Java `DataOutput#writeUTF` value (u16-BE length + modified UTF-8).
+    fn java_write_utf(s: &str) -> Vec<u8> {
+        let mut body = Vec::new();
+        for c in s.encode_utf16() {
+            if (0x0001..=0x007F).contains(&c) {
+                body.push(c as u8);
+            } else if c > 0x07FF {
+                body.push(0xE0 | (c >> 12) as u8);
+                body.push(0x80 | ((c >> 6) & 0x3F) as u8);
+                body.push(0x80 | (c & 0x3F) as u8);
+            } else {
+                body.push(0xC0 | (c >> 6) as u8);
+                body.push(0x80 | (c & 0x3F) as u8);
+            }
+        }
+        let mut out = (body.len() as u16).to_be_bytes().to_vec();
+        out.extend_from_slice(&body);
+        out
+    }
+
+    /// A valid `_SOURCE_META` frame (version 1) for the given level and files.
+    fn frame(data_level: i32, files: &[(&str, i64)]) -> Vec<u8> {
+        let mut out = Vec::new();
+        out.extend_from_slice(&1i32.to_be_bytes()); // version
+        out.extend_from_slice(&data_level.to_be_bytes());
+        out.extend_from_slice(&(files.len() as i32).to_be_bytes());
+        for (name, rows) in files {
+            out.extend_from_slice(&java_write_utf(name));
+            out.extend_from_slice(&rows.to_be_bytes());
+        }
+        out
+    }
+
+    fn gim(field_id: i32, start: i64, end: i64, source_meta: Option<Vec<u8>>) 
-> GlobalIndexMeta {
+        GlobalIndexMeta {
+            row_range_start: start,
+            row_range_end: end,
+            index_field_id: field_id,
+            extra_field_ids: None,
+            index_meta: None,
+            source_meta,
+        }
+    }
+
+    fn payload(
+        file_name: &str,
+        index_type: &str,
+        row_count: i32,
+        global_index_meta: Option<GlobalIndexMeta>,
+    ) -> IndexFileMeta {
+        IndexFileMeta {
+            index_type: index_type.into(),
+            file_name: file_name.into(),
+            file_size: 1,
+            row_count,
+            deletion_vectors_ranges: None,
+            global_index_meta,
+        }
+    }
+
+    fn ft_payload(file_name: &str, level: i32, files: &[(&str, i64)]) -> 
IndexFileMeta {
+        let total: i64 = files.iter().map(|(_, r)| *r).sum();
+        payload(
+            file_name,
+            PK_FULL_TEXT_INDEX_TYPE,
+            total as i32,
+            Some(gim(7, 0, total - 1, Some(frame(level, files)))),
+        )
+    }
+
+    fn data_split(bucket: i32, files: Vec<DataFileMeta>) -> DataSplit {
+        DataSplitBuilder::new()
+            .with_snapshot(1)
+            .with_partition(BinaryRow::new(0))
+            .with_bucket(bucket)
+            .with_bucket_path(format!("memory:/t/bucket-{bucket}"))
+            .with_total_buckets(1)
+            .with_data_files(files)
+            .build()
+            .unwrap()
+    }
+
+    // (a) planning groups payloads/data by (partition,bucket) and builds a 
split.
+    #[test]
+    fn builds_one_split_per_bucket_with_data() {
+        let entries = vec![
+            (BinaryRow::new(0), 0, ft_payload("ft-0", 1, &[("d0", 3)])),
+            (BinaryRow::new(0), 1, ft_payload("ft-1", 1, &[("d1", 4)])),
+        ];
+        let splits = plan_from_inputs(
+            1,
+            vec![
+                data_split(0, vec![dfm("d0", 3, 1)]),
+                data_split(1, vec![dfm("d1", 4, 1)]),
+            ],
+            entries,
+            7,
+        )
+        .unwrap();
+        assert_eq!(splits.len(), 2);
+        let b0 = splits.iter().find(|s| s.data_split.bucket() == 0).unwrap();
+        assert_eq!(b0.current_payloads.len(), 1);
+        assert_eq!(b0.current_payloads[0].file_name, "ft-0");
+        assert!(b0.uncovered_data_files.is_empty());
+        let b1 = splits.iter().find(|s| s.data_split.bucket() == 1).unwrap();
+        assert_eq!(b1.current_payloads[0].file_name, "ft-1");
+    }
+
+    // (b) split carries current_payloads + uncovered_data_files: one file 
covered
+    // by a current payload, another active file with no matching payload -> 
uncovered.
+    #[test]
+    fn split_carries_current_and_uncovered() {
+        // d0 (level 1) is covered by ft-0; d1 (level 2) has no payload -> 
uncovered.
+        let entries = vec![(BinaryRow::new(0), 0, ft_payload("ft-0", 1, 
&[("d0", 3)]))];
+        let splits = plan_from_inputs(
+            1,
+            vec![data_split(0, vec![dfm("d0", 3, 1), dfm("d1", 5, 2)])],
+            entries,
+            7,
+        )
+        .unwrap();
+        assert_eq!(splits.len(), 1);
+        let s = &splits[0];
+        assert_eq!(s.current_payloads.len(), 1);
+        assert_eq!(s.current_payloads[0].file_name, "ft-0");
+        assert_eq!(s.uncovered_data_files, vec!["d1".to_string()]);
+        // Completeness: covered (d0) ∪ uncovered (d1) == active files (d0, 
d1).
+        assert_eq!(s.data_split.data_files().len(), 2);
+    }
+
+    // (c) completeness invariant holds through the plan, and a constructed
+    // violation is rejected by the split constructor.
+    #[test]
+    fn split_new_rejects_incomplete_cover() {
+        // Data split has d0 + d1 active, ft-0 covers only d0, and uncovered is
+        // empty -> d1 is neither covered nor uncovered -> Err.
+        let ds = data_split(0, vec![dfm("d0", 3, 1), dfm("d1", 5, 2)]);
+        let err = PrimaryKeyFullTextSearchSplit::new(
+            ds,
+            vec![ft_payload("ft-0", 1, &[("d0", 3)])],
+            Vec::new(),
+        );
+        assert!(err.is_err(), "incomplete covered∪uncovered must fail loud");
+    }
+
+    #[test]
+    fn split_new_accepts_complete_cover() {
+        let ds = data_split(0, vec![dfm("d0", 3, 1), dfm("d1", 5, 2)]);
+        let ok = PrimaryKeyFullTextSearchSplit::new(
+            ds,
+            vec![ft_payload("ft-0", 1, &[("d0", 3)])],
+            vec!["d1".to_string()],
+        );
+        assert!(ok.is_ok());
+    }
+
+    #[test]
+    fn split_new_rejects_double_covered_source() {
+        // Two payloads both claim d0 as an active source -> double cover -> 
Err.
+        let ds = data_split(0, vec![dfm("d0", 3, 1)]);
+        let err = PrimaryKeyFullTextSearchSplit::new(
+            ds,
+            vec![
+                ft_payload("ft-a", 1, &[("d0", 3)]),
+                ft_payload("ft-b", 2, &[("d0", 3)]),
+            ],
+            Vec::new(),
+        );
+        assert!(err.is_err());
+    }
+
+    #[test]
+    fn split_new_rejects_payload_covering_no_active_source() {
+        // ft-x's only source (ghost) is not an active data file -> Err.
+        let ds = data_split(0, vec![dfm("d0", 3, 1)]);
+        let err = PrimaryKeyFullTextSearchSplit::new(
+            ds,
+            vec![
+                ft_payload("ft-0", 1, &[("d0", 3)]),
+                ft_payload("ft-x", 1, &[("ghost", 9)]),
+            ],
+            Vec::new(),
+        );
+        assert!(err.is_err());
+    }
+
+    // (d) a payload with the wrong index_type is excluded from the current 
set.
+    #[test]
+    fn wrong_index_type_payload_excluded() {
+        let bad = payload(
+            "gi-0",
+            "GLOBAL_INDEX",
+            3,
+            Some(gim(7, 0, 2, Some(frame(1, &[("d0", 3)])))),
+        );
+        let splits = plan_from_inputs(
+            1,
+            vec![data_split(0, vec![dfm("d0", 3, 1)])],
+            vec![(BinaryRow::new(0), 0, bad)],
+            7,
+        )
+        .unwrap();
+        assert_eq!(splits.len(), 1);
+        // No current payload -> the sole active file is uncovered.
+        assert!(splits[0].current_payloads.is_empty());
+        assert_eq!(splits[0].uncovered_data_files, vec!["d0".to_string()]);
+    }
+
+    #[test]
+    fn drops_index_only_bucket_without_error() {
+        // Payload for bucket 0 but NO data split -> no split, no error.
+        let entries = vec![(BinaryRow::new(0), 0, ft_payload("ft-0", 1, 
&[("d0", 3)]))];
+        let splits = plan_from_inputs(1, Vec::new(), entries, 7).unwrap();
+        assert!(splits.is_empty());
+    }
+
+    #[test]
+    fn accumulator_skips_non_should_read_files() {
+        // d0 level 0 (APPEND-like: level 0 COMPACT is not should_read) is 
skipped;
+        // the bucket has only d1 as an active source.
+        let entries = vec![(BinaryRow::new(0), 0, ft_payload("ft-1", 1, 
&[("d1", 4)]))];
+        let splits = plan_from_inputs(
+            1,
+            vec![data_split(0, vec![dfm("d0", 3, 0), dfm("d1", 4, 1)])],
+            entries,
+            7,
+        )
+        .unwrap();
+        assert_eq!(splits.len(), 1);
+        assert_eq!(splits[0].data_split.data_files().len(), 1);
+        assert_eq!(splits[0].data_split.data_files()[0].file_name, "d1");
+        assert_eq!(splits[0].current_payloads[0].file_name, "ft-1");
+    }
+
+    #[test]
+    fn skips_negative_bucket_data_split() {
+        // A negative-bucket data split is excluded from planning (mirror 
Java's
+        // `if (dataSplit.bucket() < 0) continue;`); only the valid bucket ≥ 0
+        // produces a split. (The builder reserves the -1 sentinel, so -2 
stands
+        // in for a negative bucket to exercise the `bucket < 0` skip.)
+        let entries = vec![(BinaryRow::new(0), 0, ft_payload("ft-0", 1, 
&[("d0", 3)]))];
+        let splits = plan_from_inputs(
+            1,
+            vec![
+                data_split(-2, vec![dfm("d-neg", 9, 1)]),
+                data_split(0, vec![dfm("d0", 3, 1)]),
+            ],
+            entries,
+            7,
+        )
+        .unwrap();
+        assert_eq!(splits.len(), 1);
+        assert_eq!(splits[0].data_split.bucket(), 0);
+        assert_eq!(splits[0].data_split.data_files().len(), 1);
+        assert_eq!(splits[0].data_split.data_files()[0].file_name, "d0");
+        assert_eq!(splits[0].current_payloads[0].file_name, "ft-0");
+    }
+}
diff --git a/crates/paimon/src/table/vector_search_builder.rs 
b/crates/paimon/src/table/vector_search_builder.rs
index 273dced9..9bf00de0 100644
--- a/crates/paimon/src/table/vector_search_builder.rs
+++ b/crates/paimon/src/table/vector_search_builder.rs
@@ -1845,7 +1845,7 @@ async fn rerank_indexed_positional(
 
 /// One materialized row tagged with its best-first `rank` and its 
`(batch_index,
 /// row_index)` location in the retained materialization batches.
-struct RankedRow {
+pub(crate) struct RankedRow {
     rank: usize,
     batch_index: usize,
     row_index: usize,
@@ -1857,7 +1857,7 @@ struct RankedRow {
 /// map to a candidate rank (the batch came from that candidate's file), so a 
miss
 /// fails loud rather than silently dropping a row.
 #[allow(clippy::too_many_arguments)]
-fn collect_ranked_rows(
+pub(crate) fn collect_ranked_rows(
     batch: &RecordBatch,
     batch_index: usize,
     partition_bytes: &[u8],
@@ -1908,7 +1908,7 @@ fn collect_ranked_rows(
 /// `_PKEY_VECTOR_POSITION` column, yielding a single output batch (empty input
 /// yields no batches). The projected user columns and `__paimon_search_score` 
are
 /// retained.
-fn reorder_and_strip_position(
+pub(crate) fn reorder_and_strip_position(
     batches: &[RecordBatch],
     mut ranked: Vec<RankedRow>,
 ) -> crate::Result<Vec<RecordBatch>> {

Reply via email to