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 9746004b Support managed BLOBs in native primary-key Parquet tables 
(#954)
9746004b is described below

commit 9746004b1080780dee78dd6e0f9feb50cc054a4c
Author: Jingsong Lee <[email protected]>
AuthorDate: Fri Sep 25 15:35:03 2026 +0800

    Support managed BLOBs in native primary-key Parquet tables (#954)
---
 crates/paimon/src/arrow/format/blob.rs             |   52 +
 crates/paimon/src/spec/schema.rs                   |  280 ++++-
 crates/paimon/src/table/kv_file_writer.rs          |   32 +-
 crates/paimon/src/table/managed_blob_reader.rs     |  413 ++++++++
 crates/paimon/src/table/managed_blob_reference.rs  |  287 +++++
 .../paimon/src/table/managed_blob_table_tests.rs   | 1121 ++++++++++++++++++++
 crates/paimon/src/table/managed_blob_writer.rs     |  467 ++++++++
 crates/paimon/src/table/mod.rs                     |    5 +
 crates/paimon/src/table/read_builder.rs            |    4 +-
 crates/paimon/src/table/table_read.rs              |   37 +-
 10 files changed, 2687 insertions(+), 11 deletions(-)

diff --git a/crates/paimon/src/arrow/format/blob.rs 
b/crates/paimon/src/arrow/format/blob.rs
index 8c91cbf1..e15ef57b 100644
--- a/crates/paimon/src/arrow/format/blob.rs
+++ b/crates/paimon/src/arrow/format/blob.rs
@@ -1998,6 +1998,58 @@ impl BlobFormatWriter {
             lengths: Vec::new(),
         })
     }
+
+    /// Append one managed BLOB and return the payload range in this pack.
+    /// The range excludes the four-byte entry magic and twelve-byte trailer,
+    /// matching Java `BlobFormatWriter`'s descriptor callback.
+    pub(crate) async fn write_managed_value(&mut self, value: &[u8]) -> 
crate::Result<(i64, i64)> {
+        let start = self.bytes_written;
+        let schema = 
Arc::new(arrow_schema::Schema::new(vec![arrow_schema::Field::new(
+            "blob",
+            ArrowDataType::LargeBinary,
+            false,
+        )]));
+        let batch = RecordBatch::try_new(
+            schema,
+            vec![Arc::new(LargeBinaryArray::from(vec![Some(value)]))],
+        )
+        .map_err(|error| Error::DataInvalid {
+            message: format!("Failed to build managed BLOB input: {error}"),
+            source: Some(Box::new(error)),
+        })?;
+        self.write(&batch).await?;
+        let entry_len =
+            self.bytes_written
+                .checked_sub(start)
+                .ok_or_else(|| Error::DataInvalid {
+                    message: "Managed BLOB writer position moved 
backwards".to_string(),
+                    source: None,
+                })?;
+        let payload_len =
+            entry_len
+                .checked_sub(BLOB_ENTRY_OVERHEAD)
+                .ok_or_else(|| Error::DataInvalid {
+                    message: "Managed BLOB entry is shorter than its 
framing".to_string(),
+                    source: None,
+                })?;
+        let offset =
+            start
+                .checked_add(BLOB_INLINE_HEADER_SIZE)
+                .ok_or_else(|| Error::DataInvalid {
+                    message: "Managed BLOB payload offset overflows 
u64".to_string(),
+                    source: None,
+                })?;
+        Ok((
+            i64::try_from(offset).map_err(|error| Error::DataInvalid {
+                message: "Managed BLOB payload offset exceeds i64".to_string(),
+                source: Some(Box::new(error)),
+            })?,
+            i64::try_from(payload_len).map_err(|error| Error::DataInvalid {
+                message: "Managed BLOB payload length exceeds i64".to_string(),
+                source: Some(Box::new(error)),
+            })?,
+        ))
+    }
 }
 
 const BLOB_WRITE_BUFFER_SIZE: u64 = 8 * 1024 * 1024; // 8 MB
diff --git a/crates/paimon/src/spec/schema.rs b/crates/paimon/src/spec/schema.rs
index 56651d9f..9da4de9a 100644
--- a/crates/paimon/src/spec/schema.rs
+++ b/crates/paimon/src/spec/schema.rs
@@ -1197,7 +1197,8 @@ impl Schema {
         validate_no_reserved_field_names(fields)?;
         Self::validate_key_field_types(fields, primary_keys, options)?;
         Self::validate_row_tracking(primary_keys, options)?;
-        Self::validate_blob_fields(fields, partition_keys, options)?;
+        Self::validate_blob_fields(fields, partition_keys, primary_keys, 
options)?;
+        Self::validate_primary_key_blob_configuration(fields, primary_keys, 
options)?;
         Self::validate_vector_store_fields(fields, partition_keys, options)?;
         
PartialUpdateConfig::new(options).validate_create_mode(!primary_keys.is_empty())?;
         validate_no_aggregation_on_sequence_field(options)?;
@@ -1480,9 +1481,23 @@ impl Schema {
     fn validate_blob_fields(
         fields: &[DataField],
         partition_keys: &[String],
+        primary_keys: &[String],
         options: &HashMap<String, String>,
     ) -> crate::Result<()> {
         let blob_field_names = Self::top_level_blob_field_names(fields);
+        for field in fields {
+            if !Self::is_top_level_blob_file_type(field.data_type())
+                && Self::contains_blob_type(field.data_type())
+            {
+                return Err(crate::Error::ConfigInvalid {
+                    message: format!(
+                        "Field '{}' has unsupported nested BLOB type {:?}. 
BLOB is only supported as a top-level BLOB, ARRAY<BLOB>, or MAP<X, BLOB> 
field.",
+                        field.name(),
+                        field.data_type()
+                    ),
+                });
+            }
+        }
         if blob_field_names.is_empty() {
             return Ok(());
         }
@@ -1504,10 +1519,36 @@ impl Schema {
             });
         }
 
-        if !core_options.data_evolution_enabled() {
+        for name in core_options.blob_fields() {
+            if !blob_field_names.contains(&name.as_str()) {
+                return Err(crate::Error::ConfigInvalid {
+                    message: format!(
+                        "Field '{name}' in '{BLOB_FIELD_OPTION}' must be a 
BLOB, ARRAY<BLOB> or MAP<X, BLOB> field in table schema."
+                    ),
+                });
+            }
+        }
+        for (option_name, names) in [
+            (BLOB_DESCRIPTOR_FIELD_OPTION, &blob_descriptor_fields),
+            (BLOB_VIEW_FIELD_OPTION, &blob_view_fields),
+        ] {
+            for name in names {
+                let scalar_blob = fields.iter().any(|field| {
+                    field.name() == name && matches!(field.data_type(), 
DataType::Blob(_))
+                });
+                if !scalar_blob {
+                    return Err(crate::Error::ConfigInvalid {
+                        message: format!(
+                            "Field '{name}' in '{option_name}' must be a 
scalar BLOB field in table schema. ARRAY<BLOB> and MAP<X, BLOB> are only 
supported by '{BLOB_FIELD_OPTION}'."
+                        ),
+                    });
+                }
+            }
+        }
+
+        if primary_keys.is_empty() && !core_options.data_evolution_enabled() {
             return Err(crate::Error::ConfigInvalid {
-                message: "Data evolution config must enabled for table with 
BLOB type column."
-                    .to_string(),
+                message: "Data evolution config must enabled for table with 
BLOB, ARRAY<BLOB> or MAP<X, BLOB> type column.".to_string(),
             });
         }
 
@@ -1530,6 +1571,80 @@ impl Schema {
         Ok(())
     }
 
+    fn validate_primary_key_blob_configuration(
+        fields: &[DataField],
+        primary_keys: &[String],
+        options: &HashMap<String, String>,
+    ) -> crate::Result<()> {
+        if primary_keys.is_empty() {
+            return Ok(());
+        }
+        let core = CoreOptions::new(options);
+        let inline = core.blob_inline_fields();
+        let managed = fields
+            .iter()
+            .filter(|field| {
+                Self::is_top_level_blob_file_type(field.data_type())
+                    && !inline.contains(field.name())
+            })
+            .map(|field| field.name())
+            .collect::<Vec<_>>();
+        if managed.is_empty() {
+            return Ok(());
+        }
+        if !matches!(
+            core.merge_engine()?,
+            MergeEngine::Deduplicate | MergeEngine::PartialUpdate | 
MergeEngine::FirstRow
+        ) {
+            return Err(crate::Error::ConfigInvalid {
+                message: "Primary-key managed BLOB tables only support 
deduplicate, partial-update or first-row merge engine.".to_string(),
+            });
+        }
+        if core.try_changelog_producer()? != ChangelogProducer::None {
+            return Err(crate::Error::ConfigInvalid {
+                message: "Primary-key managed BLOB tables only support 
changelog-producer 'none'."
+                    .to_string(),
+            });
+        }
+        if options.contains_key("data-file.external-paths") {
+            return Err(crate::Error::ConfigInvalid {
+                message:
+                    "Primary-key managed BLOB tables do not support 
'data-file.external-paths'."
+                        .to_string(),
+            });
+        }
+        if options
+            .get("pk-clustering-override")
+            .is_some_and(|v| v.eq_ignore_ascii_case("true"))
+        {
+            return Err(crate::Error::ConfigInvalid {
+                message: "Primary-key managed BLOB tables do not support 
'pk-clustering-override'."
+                    .to_string(),
+            });
+        }
+        for (name, keys) in [
+            ("primary keys", primary_keys.to_vec()),
+            ("bucket keys", core.bucket_key().unwrap_or_default()),
+            (
+                "sequence fields",
+                core.sequence_fields()
+                    .iter()
+                    .map(|s| s.to_string())
+                    .collect(),
+            ),
+        ] {
+            if let Some(field) = managed
+                .iter()
+                .find(|field| keys.iter().any(|key| key == **field))
+            {
+                return Err(crate::Error::ConfigInvalid {
+                    message: format!("Managed BLOB field '{field}' cannot be 
used as {name}."),
+                });
+            }
+        }
+        Ok(())
+    }
+
     fn validate_vector_store_fields(
         fields: &[DataField],
         partition_keys: &[String],
@@ -2018,13 +2133,37 @@ impl Schema {
     fn top_level_blob_field_names(fields: &[DataField]) -> Vec<&str> {
         fields
             .iter()
-            .filter_map(|field| match field.data_type() {
-                DataType::Blob(_) => Some(field.name()),
-                _ => None,
-            })
+            .filter(|field| 
Self::is_top_level_blob_file_type(field.data_type()))
+            .map(|field| field.name())
             .collect()
     }
 
+    fn is_top_level_blob_file_type(data_type: &DataType) -> bool {
+        match data_type {
+            DataType::Blob(_) => true,
+            DataType::Array(array) => matches!(array.element_type(), 
DataType::Blob(_)),
+            DataType::Map(map) => matches!(map.value_type(), 
DataType::Blob(_)),
+            _ => false,
+        }
+    }
+
+    fn contains_blob_type(data_type: &DataType) -> bool {
+        match data_type {
+            DataType::Blob(_) => true,
+            DataType::Array(array) => 
Self::contains_blob_type(array.element_type()),
+            DataType::Map(map) => {
+                Self::contains_blob_type(map.key_type())
+                    || Self::contains_blob_type(map.value_type())
+            }
+            DataType::Multiset(multiset) => 
Self::contains_blob_type(multiset.element_type()),
+            DataType::Row(row) => row
+                .fields()
+                .iter()
+                .any(|field| Self::contains_blob_type(field.data_type())),
+            _ => false,
+        }
+    }
+
     /// Returns top-level Vector field names for dedicated vector-store checks.
     fn top_level_vector_field_names(fields: &[DataField]) -> Vec<&str> {
         fields
@@ -2817,6 +2956,131 @@ mod tests {
         assert_eq!(schema.fields().len(), 2);
     }
 
+    #[test]
+    fn test_primary_key_managed_blob_does_not_require_data_evolution() {
+        let schema = Schema::builder()
+            .column("id", DataType::Int(IntType::new()))
+            .column("payload", DataType::Blob(BlobType::new()))
+            .column(
+                "payloads",
+                
DataType::Array(ArrayType::new(DataType::Blob(BlobType::new()))),
+            )
+            .primary_key(["id"])
+            .option("bucket", "1")
+            .build()
+            .unwrap();
+        assert_eq!(schema.primary_keys(), &["id"]);
+        assert!(!CoreOptions::new(schema.options()).data_evolution_enabled());
+    }
+
+    #[test]
+    fn test_primary_key_managed_blob_rejects_incompatible_write_modes() {
+        for (option, value, expected) in [
+            ("merge-engine", "aggregation", "merge engine"),
+            ("changelog-producer", "input", "changelog-producer"),
+            (
+                "data-file.external-paths",
+                "file:///tmp/external",
+                "external-paths",
+            ),
+            ("pk-clustering-override", "true", "pk-clustering-override"),
+        ] {
+            let result = Schema::builder()
+                .column("id", DataType::Int(IntType::new()))
+                .column("payload", DataType::Blob(BlobType::new()))
+                .primary_key(["id"])
+                .option("bucket", "1")
+                .option(option, value)
+                .build();
+            assert!(
+                matches!(result, Err(crate::Error::ConfigInvalid { ref message 
})
+                    if message.contains(expected)),
+                "expected '{option}={value}' to be rejected for {expected}, 
got {result:?}"
+            );
+        }
+    }
+
+    #[test]
+    fn test_primary_key_managed_blob_rejects_key_and_ordering_fields() {
+        let result = Schema::builder()
+            .column("id", DataType::Int(IntType::new()))
+            .column("payload", DataType::Blob(BlobType::new()))
+            .primary_key(["payload"])
+            .option("bucket", "1")
+            .build();
+        assert!(
+            matches!(result, Err(crate::Error::ConfigInvalid { ref message })
+            if message.contains("Managed BLOB") && message.contains("primary 
keys"))
+        );
+
+        for (option, value, expected) in [
+            ("bucket-key", "payload", "bucket keys"),
+            ("sequence.field", "payload", "sequence fields"),
+        ] {
+            let result = Schema::builder()
+                .column("id", DataType::Int(IntType::new()))
+                .column("payload", DataType::Blob(BlobType::new()))
+                .primary_key(["id"])
+                .option("bucket", "2")
+                .option(option, value)
+                .build();
+            assert!(
+                matches!(result, Err(crate::Error::ConfigInvalid { ref message 
})
+                if message.contains("Managed BLOB") && 
message.contains(expected)),
+                "expected {option} to reject managed BLOB values, got 
{result:?}"
+            );
+        }
+    }
+
+    #[test]
+    fn 
test_blob_schema_validation_rejects_unsupported_nested_shapes_and_inline_options()
 {
+        let nested = DataType::Row(RowType::new(vec![DataField::new(
+            0,
+            "inner".to_string(),
+            DataType::Blob(BlobType::new()),
+        )]));
+        for field_type in [
+            nested,
+            DataType::Map(MapType::new(
+                DataType::Blob(BlobType::new()),
+                DataType::Int(IntType::new()),
+            )),
+            DataType::Array(ArrayType::new(DataType::Array(ArrayType::new(
+                DataType::Blob(BlobType::new()),
+            )))),
+        ] {
+            let result = Schema::builder()
+                .column("id", DataType::Int(IntType::new()))
+                .column("payload", field_type)
+                .primary_key(["id"])
+                .option("bucket", "1")
+                .build();
+            assert!(
+                matches!(result, Err(crate::Error::ConfigInvalid { ref message 
})
+                if message.contains("unsupported nested BLOB")),
+                "unsupported nesting should fail at schema validation: 
{result:?}"
+            );
+        }
+
+        for option in ["blob-descriptor-field", "blob-view-field"] {
+            let result = Schema::builder()
+                .column("id", DataType::Int(IntType::new()))
+                .column(
+                    "payloads",
+                    
DataType::Array(ArrayType::new(DataType::Blob(BlobType::new()))),
+                )
+                .primary_key(["id"])
+                .option("bucket", "1")
+                .option(option, "payloads")
+                .build();
+            assert!(
+                matches!(result, Err(crate::Error::ConfigInvalid { ref message 
})
+                if message.contains("scalar BLOB")),
+                "{option} is scalar-only in Java: {result:?}"
+            );
+        }
+    }
+
     #[test]
     fn test_blob_field_option_promotes_binary_column() {
         let schema = Schema::builder()
diff --git a/crates/paimon/src/table/kv_file_writer.rs 
b/crates/paimon/src/table/kv_file_writer.rs
index 05ae9aad..a8dd6842 100644
--- a/crates/paimon/src/table/kv_file_writer.rs
+++ b/crates/paimon/src/table/kv_file_writer.rs
@@ -40,6 +40,8 @@ use crate::spec::{
     SEQUENCE_NUMBER_FIELD_NAME, VALUE_KIND_FIELD_ID, VALUE_KIND_FIELD_NAME,
 };
 use crate::table::data_file_index_writer::FileIndexOptions;
+use crate::table::managed_blob_reference::ManagedBlobReferences;
+use crate::table::managed_blob_writer::ManagedBlobWriteState;
 use crate::table::prepared_files::PreparedFiles;
 use crate::table::sort_merge::{AggregateMergeFunction, BufferedBatch, 
MergeRow};
 use crate::Result;
@@ -69,6 +71,7 @@ pub(crate) struct KeyValueFileWriter {
     written_files: Vec<DataFileMeta>,
     /// Completed changelog file metadata.
     written_changelog_files: Vec<DataFileMeta>,
+    managed_blob_writer: ManagedBlobWriteState,
 }
 
 /// Configuration for [`KeyValueFileWriter`], grouping file-location, schema,
@@ -154,6 +157,8 @@ impl KeyValueFileWriter {
             }
         }
 
+        let managed_blob_writer = ManagedBlobWriteState::new(&file_io, 
&config)?;
+
         Ok(Self {
             file_io,
             config,
@@ -165,6 +170,7 @@ impl KeyValueFileWriter {
             buffer_reservation: None,
             written_files: Vec::new(),
             written_changelog_files: Vec::new(),
+            managed_blob_writer,
         })
     }
 
@@ -187,6 +193,7 @@ impl KeyValueFileWriter {
         if batch.num_rows() == 0 {
             return Ok(());
         }
+        let batch = self.managed_blob_writer.externalize(batch).await?;
         let batch_bytes: usize = batch
             .columns()
             .iter()
@@ -437,8 +444,15 @@ impl KeyValueFileWriter {
                 .map(|options| options.create_writer())
                 .transpose()?
         };
-
         let physical_schema = build_physical_schema(&user_schema);
+        let mut blob_references = ManagedBlobReferences::new(
+            &self.config.value_fields,
+            &CoreOptions::new(&self.config.table_options),
+            &physical_schema,
+            self.managed_blob_writer.enabled(),
+            write.is_changelog,
+        )?;
+
         let file_name = format!(
             "{}{}-{}.{}",
             write.file_prefix,
@@ -548,6 +562,11 @@ impl KeyValueFileWriter {
                     message: format!("Failed to create physical batch: {e}"),
                     source: None,
                 })?;
+            if let Err(error) = ManagedBlobReferences::collect(&mut 
blob_references, &chunk_batch) {
+                let _ = writer.close().await;
+                let _ = self.file_io.delete_file(&file_path).await;
+                return Err(error);
+            }
             if let Err(error) = writer.write(&chunk_batch).await {
                 let _ = writer.close().await;
                 let _ = self.file_io.delete_file(&file_path).await;
@@ -674,6 +693,15 @@ impl KeyValueFileWriter {
             }
         }
 
+        ManagedBlobReferences::finish(
+            blob_references,
+            &self.file_io,
+            &file_path,
+            &bucket_dir,
+            &mut meta,
+        )
+        .await?;
+
         Ok(meta)
     }
 
@@ -1035,11 +1063,13 @@ impl KeyValueFileWriter {
                 let _ = self.file_io.delete_file(&path).await;
             }
         }
+        self.managed_blob_writer.abort().await;
     }
 
     /// Flush remaining buffer and return all written file metadata.
     pub(crate) async fn prepare_commit(&mut self) -> Result<PreparedFiles> {
         self.flush().await?;
+        self.managed_blob_writer.prepare_commit().await?;
         Ok(PreparedFiles {
             data_files: std::mem::take(&mut self.written_files),
             changelog_files: std::mem::take(&mut self.written_changelog_files),
diff --git a/crates/paimon/src/table/managed_blob_reader.rs 
b/crates/paimon/src/table/managed_blob_reader.rs
new file mode 100644
index 00000000..6cdaf808
--- /dev/null
+++ b/crates/paimon/src/table/managed_blob_reader.rs
@@ -0,0 +1,413 @@
+// 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.
+
+//! Resolve BLOB descriptors after primary-key merge and row selection.
+
+use super::blob_resolver::{resolve_blob_column, BlobReadLimiter};
+use super::managed_blob_writer::{managed_blob_kind, ManagedBlobKind};
+use super::{ArrowRecordBatchStream, Table};
+use crate::arrow::format::FilePredicates;
+use crate::io::FileIO;
+use crate::spec::{CoreOptions, DataField, Predicate};
+use crate::Result;
+use arrow_array::builder::LargeBinaryBuilder;
+use arrow_array::{
+    Array, ArrayRef, LargeBinaryArray, ListArray, MapArray, RecordBatch, 
StructArray,
+};
+use arrow_schema::DataType as ArrowDataType;
+use futures::StreamExt;
+use std::collections::HashSet;
+use std::sync::Arc;
+
+pub(crate) fn resolve_primary_key_blob_stream(
+    stream: ArrowRecordBatchStream,
+    fields: &[DataField],
+    options: &CoreOptions<'_>,
+    file_io: FileIO,
+    parallelism: usize,
+) -> ArrowRecordBatchStream {
+    let selected = resolved_primary_key_blob_fields(fields, options);
+    if selected.is_empty() {
+        return stream;
+    }
+    let limiter = BlobReadLimiter::with_parallelism(parallelism);
+    Box::pin(async_stream::try_stream! {
+        let mut stream = stream;
+        while let Some(batch) = stream.next().await {
+            let batch = batch?;
+            yield resolve_batch(batch, &selected, &file_io, &limiter).await?;
+        }
+    })
+}
+
+pub(crate) fn resolved_primary_key_blob_fields(
+    fields: &[DataField],
+    options: &CoreOptions<'_>,
+) -> Vec<(usize, ManagedBlobKind)> {
+    if options.blob_as_descriptor() {
+        return Vec::new();
+    }
+    let descriptor_fields = options.blob_descriptor_fields();
+    let inline_fields = options.blob_inline_fields();
+    fields
+        .iter()
+        .enumerate()
+        .filter_map(|(index, field)| {
+            let kind = managed_blob_kind(field.data_type())?;
+            (!inline_fields.contains(field.name()) || 
descriptor_fields.contains(field.name()))
+                .then_some((index, kind))
+        })
+        .collect()
+}
+
+fn resolved_blob_indices(fields: &[DataField], options: &CoreOptions<'_>) -> 
HashSet<usize> {
+    resolved_primary_key_blob_fields(fields, options)
+        .into_iter()
+        .map(|(index, _)| index)
+        .collect()
+}
+
+fn predicate_uses_resolved_blob(predicate: &Predicate, resolved: 
&HashSet<usize>) -> bool {
+    let mut referenced = HashSet::new();
+    predicate.collect_leaf_field_indices(&mut referenced);
+    referenced.iter().any(|index| resolved.contains(index))
+}
+
+/// Drop payload predicates from file and stats pruning while retaining safe
+/// predicates on ordinary columns. The full filter remains on TableRead.
+pub(crate) fn scan_predicates(table: &Table, predicates: &[Predicate]) -> 
Vec<Predicate> {
+    if table.schema().primary_keys().is_empty() {
+        return predicates.to_vec();
+    }
+    let options = table.schema().core_options();
+    let resolved = resolved_blob_indices(table.schema().fields(), &options);
+    if resolved.is_empty() {
+        return predicates.to_vec();
+    }
+    predicates
+        .iter()
+        .filter(|predicate| !predicate_uses_resolved_blob(predicate, 
&resolved))
+        .cloned()
+        .collect()
+}
+
+/// Holds the extra columns and exact residual filter needed when a primary-key
+/// predicate compares BLOB payloads rather than their Parquet descriptors.
+pub(crate) struct ManagedBlobReadPlan {
+    scan_fields: Vec<DataField>,
+    output_fields: Vec<DataField>,
+    predicates: FilePredicates,
+}
+
+impl ManagedBlobReadPlan {
+    pub(crate) fn new(
+        read_type: &[DataField],
+        predicates: &[Predicate],
+        table_fields: &[DataField],
+        options: &CoreOptions<'_>,
+    ) -> Option<Self> {
+        let resolved = resolved_blob_indices(table_fields, options);
+        if resolved.is_empty() {
+            return None;
+        }
+        predicates
+            .iter()
+            .any(|predicate| predicate_uses_resolved_blob(predicate, 
&resolved))
+            .then(|| {
+                let predicates = FilePredicates {
+                    predicates: predicates.to_vec(),
+                    row_filter_factory: None,
+                    file_fields: table_fields.to_vec(),
+                };
+                let scan_fields =
+                    crate::arrow::residual::widen_scan_fields(read_type, 
Some(&predicates));
+                Self {
+                    scan_fields,
+                    output_fields: read_type.to_vec(),
+                    predicates,
+                }
+            })
+    }
+
+    pub(crate) fn scan_fields(&self) -> &[DataField] {
+        &self.scan_fields
+    }
+
+    pub(crate) fn finish(
+        self,
+        stream: ArrowRecordBatchStream,
+        options: &CoreOptions<'_>,
+        file_io: FileIO,
+        parallelism: usize,
+    ) -> ArrowRecordBatchStream {
+        let stream = resolve_primary_key_blob_stream(
+            stream,
+            &self.scan_fields,
+            options,
+            file_io,
+            parallelism,
+        );
+        Box::pin(async_stream::try_stream! {
+            let mut stream = stream;
+            while let Some(batch) = stream.next().await {
+                let batch = 
crate::arrow::residual::filter_record_batch_by_predicates(
+                    batch?, &self.predicates, &self.scan_fields,
+                )?;
+                let indices = self.output_fields
+                    .iter()
+                    .map(|field| batch.schema().index_of(field.name()))
+                    .collect::<std::result::Result<Vec<_>, _>>()
+                    .map_err(|error| crate::Error::DataInvalid {
+                        message: format!("Managed BLOB output projection is 
missing a column: {error}"),
+                        source: Some(Box::new(error)),
+                    })?;
+                yield batch.project(&indices).map_err(|error| 
crate::Error::DataInvalid {
+                    message: format!("Failed to project managed BLOB read 
output: {error}"),
+                    source: Some(Box::new(error)),
+                })?;
+            }
+        })
+    }
+}
+
+async fn resolve_batch(
+    batch: RecordBatch,
+    fields: &[(usize, ManagedBlobKind)],
+    file_io: &FileIO,
+    limiter: &BlobReadLimiter,
+) -> Result<RecordBatch> {
+    let mut columns = batch.columns().to_vec();
+    for &(index, kind) in fields {
+        let column = match kind {
+            ManagedBlobKind::Scalar => {
+                let values = blob_values(columns[index].as_ref())?;
+                Arc::new(resolve_blob_column(values, file_io, 
limiter.clone()).await?) as ArrayRef
+            }
+            ManagedBlobKind::Array => {
+                let array = columns[index]
+                    .as_any()
+                    .downcast_ref::<ListArray>()
+                    .ok_or_else(|| invalid("ARRAY<BLOB> requires ListArray"))?;
+                let values = blob_values(array.values().as_ref())?;
+                let visible = visible_child_values(values, 
array.value_offsets(), array)?;
+                let resolved =
+                    Arc::new(resolve_blob_column(&visible, file_io, 
limiter.clone()).await?);
+                let ArrowDataType::List(element) = array.data_type() else {
+                    unreachable!()
+                };
+                Arc::new(
+                    ListArray::try_new(
+                        element.clone(),
+                        array.offsets().clone(),
+                        resolved,
+                        array.nulls().cloned(),
+                    )
+                    .map_err(|error| invalid(&error.to_string()))?,
+                )
+            }
+            ManagedBlobKind::Map => {
+                let map = columns[index]
+                    .as_any()
+                    .downcast_ref::<MapArray>()
+                    .ok_or_else(|| invalid("MAP<X, BLOB> requires MapArray"))?;
+                let values = blob_values(map.entries().column(1).as_ref())?;
+                let visible = visible_child_values(values, 
map.value_offsets(), map)?;
+                let resolved =
+                    Arc::new(resolve_blob_column(&visible, file_io, 
limiter.clone()).await?);
+                let ArrowDataType::Map(entries_field, ordered) = 
map.data_type() else {
+                    unreachable!()
+                };
+                let ArrowDataType::Struct(entry_fields) = 
entries_field.data_type() else {
+                    unreachable!()
+                };
+                let entries = StructArray::try_new(
+                    entry_fields.clone(),
+                    vec![map.entries().column(0).clone(), resolved],
+                    None,
+                )
+                .map_err(|error| invalid(&error.to_string()))?;
+                Arc::new(
+                    MapArray::try_new(
+                        entries_field.clone(),
+                        map.offsets().clone(),
+                        entries,
+                        map.nulls().cloned(),
+                        *ordered,
+                    )
+                    .map_err(|error| invalid(&error.to_string()))?,
+                )
+            }
+        };
+        columns[index] = column;
+    }
+    RecordBatch::try_new(batch.schema(), columns).map_err(|error| 
invalid(&error.to_string()))
+}
+
+/// Arrow may retain child values beneath a null collection parent. They are
+/// invisible to the row and must not trigger descriptor I/O or fail the read
+/// when a stale URI happens to remain in an unused child slot.
+fn visible_child_values(
+    values: &LargeBinaryArray,
+    offsets: &[i32],
+    parents: &dyn Array,
+) -> Result<LargeBinaryArray> {
+    if offsets.len() != parents.len() + 1 {
+        return Err(invalid("BLOB collection offsets do not match parent 
rows"));
+    }
+    let mut visible = vec![false; values.len()];
+    for row in 0..parents.len() {
+        if !parents.is_valid(row) {
+            continue;
+        }
+        let start = usize::try_from(offsets[row])
+            .map_err(|_| invalid("Negative BLOB collection offset"))?;
+        let end = usize::try_from(offsets[row + 1])
+            .map_err(|_| invalid("Negative BLOB collection offset"))?;
+        let range = visible
+            .get_mut(start..end)
+            .ok_or_else(|| invalid("BLOB collection offset exceeds child 
values"))?;
+        range.fill(true);
+    }
+    let mut builder = LargeBinaryBuilder::new();
+    for (index, visible) in visible.into_iter().enumerate() {
+        if visible && values.is_valid(index) {
+            builder.append_value(values.value(index));
+        } else {
+            builder.append_null();
+        }
+    }
+    Ok(builder.finish())
+}
+
+fn blob_values(array: &dyn Array) -> Result<&LargeBinaryArray> {
+    array
+        .as_any()
+        .downcast_ref::<LargeBinaryArray>()
+        .ok_or_else(|| invalid("BLOB values require LargeBinaryArray"))
+}
+
+fn invalid(message: &str) -> crate::Error {
+    crate::Error::DataInvalid {
+        message: message.to_string(),
+        source: None,
+    }
+}
+
+#[cfg(test)]
+mod tests {
+    use super::*;
+    use crate::io::FileIOBuilder;
+    use crate::spec::BlobDescriptor;
+    use arrow_array::StringArray;
+    use arrow_buffer::{NullBuffer, OffsetBuffer, ScalarBuffer};
+    use arrow_schema::{Field, Schema};
+
+    #[tokio::test]
+    async fn null_array_parent_does_not_fetch_hidden_descriptor() {
+        let missing =
+            BlobDescriptor::new("memory:/missing-managed.blob".to_string(), 0, 
4).serialize();
+        let values = 
Arc::new(LargeBinaryArray::from(vec![Some(missing.as_slice())]));
+        let element = Arc::new(Field::new("element", 
ArrowDataType::LargeBinary, true));
+        let array = Arc::new(
+            ListArray::try_new(
+                element.clone(),
+                OffsetBuffer::new(ScalarBuffer::from(vec![0, 1])),
+                values,
+                Some(NullBuffer::from(vec![false])),
+            )
+            .unwrap(),
+        );
+        let schema = Arc::new(Schema::new(vec![Field::new(
+            "items",
+            ArrowDataType::List(element),
+            true,
+        )]));
+        let batch = RecordBatch::try_new(schema, vec![array]).unwrap();
+        let io = FileIOBuilder::new("memory").build().unwrap();
+        let resolved = resolve_batch(
+            batch,
+            &[(0, ManagedBlobKind::Array)],
+            &io,
+            &BlobReadLimiter::new(),
+        )
+        .await
+        .unwrap();
+        let array = resolved
+            .column(0)
+            .as_any()
+            .downcast_ref::<ListArray>()
+            .unwrap();
+        assert!(array.is_null(0));
+        assert!(array.values().is_null(0));
+    }
+
+    #[tokio::test]
+    async fn null_map_parent_does_not_fetch_hidden_descriptor() {
+        let missing =
+            BlobDescriptor::new("memory:/missing-managed.blob".to_string(), 0, 
4).serialize();
+        let entry_fields = vec![
+            Arc::new(Field::new("key", ArrowDataType::Utf8, false)),
+            Arc::new(Field::new("value", ArrowDataType::LargeBinary, true)),
+        ];
+        let entries = StructArray::try_new(
+            entry_fields.clone().into(),
+            vec![
+                Arc::new(StringArray::from(vec!["hidden"])),
+                
Arc::new(LargeBinaryArray::from(vec![Some(missing.as_slice())])),
+            ],
+            None,
+        )
+        .unwrap();
+        let entry_field = Arc::new(Field::new(
+            "entries",
+            ArrowDataType::Struct(entry_fields.into()),
+            false,
+        ));
+        let map = Arc::new(
+            MapArray::try_new(
+                entry_field.clone(),
+                OffsetBuffer::new(ScalarBuffer::from(vec![0, 1])),
+                entries,
+                Some(NullBuffer::from(vec![false])),
+                false,
+            )
+            .unwrap(),
+        );
+        let schema = Arc::new(Schema::new(vec![Field::new(
+            "named",
+            ArrowDataType::Map(entry_field, false),
+            true,
+        )]));
+        let batch = RecordBatch::try_new(schema, vec![map]).unwrap();
+        let io = FileIOBuilder::new("memory").build().unwrap();
+        let resolved = resolve_batch(
+            batch,
+            &[(0, ManagedBlobKind::Map)],
+            &io,
+            &BlobReadLimiter::new(),
+        )
+        .await
+        .unwrap();
+        let map = resolved
+            .column(0)
+            .as_any()
+            .downcast_ref::<MapArray>()
+            .unwrap();
+        assert!(map.is_null(0));
+        assert!(map.entries().column(1).is_null(0));
+    }
+}
diff --git a/crates/paimon/src/table/managed_blob_reference.rs 
b/crates/paimon/src/table/managed_blob_reference.rs
new file mode 100644
index 00000000..7f107d64
--- /dev/null
+++ b/crates/paimon/src/table/managed_blob_reference.rs
@@ -0,0 +1,287 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Java-compatible `.blobref` sidecars for managed primary-key BLOB packs.
+
+use super::managed_blob_writer::{managed_blob_fields, ManagedBlobKind};
+use crate::io::FileIO;
+use crate::spec::{
+    BlobDescriptor, CoreOptions, DataField, DataFileMeta, RowKind, 
VALUE_KIND_FIELD_NAME,
+};
+use crate::Result;
+use arrow_array::{Array, Int8Array, LargeBinaryArray, ListArray, MapArray, 
RecordBatch};
+use arrow_schema::Schema as ArrowSchema;
+use bytes::Bytes;
+use std::collections::BTreeSet;
+
+const MAGIC: i32 = 0x50424C52;
+const VERSION: u8 = 1;
+pub(crate) const REFERENCE_FILE_SUFFIX: &str = ".blobref";
+const MANAGED_BLOB_SUFFIX: &str = ".managed.blob";
+
+/// One reference is a storage root and a file name, matching Java
+/// `ManagedBlobReferenceFile.Reference`.
+#[derive(Clone, Debug, Eq, Ord, PartialEq, PartialOrd)]
+pub(crate) struct ManagedBlobReference {
+    pub storage_root: String,
+    pub file_name: String,
+}
+
+#[derive(Default)]
+pub(crate) struct ManagedBlobReferenceCollector {
+    references: BTreeSet<ManagedBlobReference>,
+}
+
+/// Per-file reference tracking, absent for ordinary and changelog files.
+pub(crate) struct ManagedBlobReferences {
+    fields: Vec<(usize, ManagedBlobKind)>,
+    collector: ManagedBlobReferenceCollector,
+}
+
+impl ManagedBlobReferences {
+    pub(crate) fn new(
+        value_fields: &[DataField],
+        table_options: &CoreOptions<'_>,
+        physical_schema: &ArrowSchema,
+        managed_enabled: bool,
+        is_changelog: bool,
+    ) -> Result<Option<Self>> {
+        if !managed_enabled || is_changelog {
+            return Ok(None);
+        }
+        let fields = managed_blob_fields(value_fields, table_options)
+            .into_iter()
+            .map(|(index, kind)| {
+                let name = value_fields[index].name();
+                physical_schema
+                    .index_of(name)
+                    .map(|physical_index| (physical_index, kind))
+                    .map_err(|error| {
+                        invalid(&format!(
+                            "Managed BLOB field '{name}' is missing from 
physical schema: {error}"
+                        ))
+                    })
+            })
+            .collect::<Result<Vec<_>>>()?;
+        Ok((!fields.is_empty()).then(|| Self {
+            fields,
+            collector: ManagedBlobReferenceCollector::default(),
+        }))
+    }
+
+    pub(crate) fn collect(references: &mut Option<Self>, batch: &RecordBatch) 
-> Result<()> {
+        if let Some(references) = references {
+            references
+                .collector
+                .collect_batch(batch, &references.fields)?;
+        }
+        Ok(())
+    }
+
+    pub(crate) async fn finish(
+        references: Option<Self>,
+        file_io: &FileIO,
+        data_path: &str,
+        bucket_dir: &str,
+        meta: &mut DataFileMeta,
+    ) -> Result<()> {
+        if let Some(references) = references {
+            match references.collector.write(file_io, data_path).await {
+                Ok(name) => meta.extra_files.push(name),
+                Err(error) => {
+                    for path in meta.collect_files(bucket_dir) {
+                        let _ = file_io.delete_file(&path).await;
+                    }
+                    return Err(error);
+                }
+            }
+        }
+        Ok(())
+    }
+}
+
+impl ManagedBlobReferenceCollector {
+    pub(crate) fn collect_batch(
+        &mut self,
+        batch: &RecordBatch,
+        fields: &[(usize, ManagedBlobKind)],
+    ) -> Result<()> {
+        let kind_index = batch
+            .schema()
+            .fields()
+            .iter()
+            .position(|field| field.name() == VALUE_KIND_FIELD_NAME);
+        let kinds = kind_index
+            .map(|index| {
+                batch
+                    .column(index)
+                    .as_any()
+                    .downcast_ref::<Int8Array>()
+                    .ok_or_else(|| invalid("_VALUE_KIND column must be Int8"))
+            })
+            .transpose()?;
+        for row in 0..batch.num_rows() {
+            let kind = kinds
+                .filter(|array| array.is_valid(row))
+                .map_or(RowKind::Insert.to_value(), |array| array.value(row));
+            if RowKind::from_value(kind)?.is_retract() {
+                continue;
+            }
+            for &(index, field_kind) in fields {
+                let column = batch.column(index);
+                if column.is_null(row) {
+                    continue;
+                }
+                match field_kind {
+                    ManagedBlobKind::Scalar => {
+                        let values = binary_column(column.as_ref())?;
+                        self.collect_value(values.value(row))?;
+                    }
+                    ManagedBlobKind::Array => {
+                        let array = column
+                            .as_any()
+                            .downcast_ref::<ListArray>()
+                            .ok_or_else(|| invalid("Managed ARRAY<BLOB> 
requires ListArray"))?;
+                        let values = binary_column(array.values().as_ref())?;
+                        for value in 
array.value_offsets()[row]..array.value_offsets()[row + 1] {
+                            let value = value as usize;
+                            if values.is_valid(value) {
+                                self.collect_value(values.value(value))?;
+                            }
+                        }
+                    }
+                    ManagedBlobKind::Map => {
+                        let map = column
+                            .as_any()
+                            .downcast_ref::<MapArray>()
+                            .ok_or_else(|| invalid("Managed MAP<X, BLOB> 
requires MapArray"))?;
+                        let values = 
binary_column(map.entries().column(1).as_ref())?;
+                        for value in 
map.value_offsets()[row]..map.value_offsets()[row + 1] {
+                            let value = value as usize;
+                            if values.is_valid(value) {
+                                self.collect_value(values.value(value))?;
+                            }
+                        }
+                    }
+                }
+            }
+        }
+        Ok(())
+    }
+
+    fn collect_value(&mut self, value: &[u8]) -> Result<()> {
+        if !BlobDescriptor::is_blob_descriptor(value) {
+            return Ok(());
+        }
+        let descriptor = BlobDescriptor::deserialize(value)?;
+        let uri = descriptor.uri();
+        if !uri.ends_with(MANAGED_BLOB_SUFFIX) {
+            return Ok(());
+        }
+        let slash = uri
+            .rfind('/')
+            .ok_or_else(|| invalid("Managed BLOB descriptor URI has no parent 
directory"))?;
+        let file_name = &uri[slash + 1..];
+        if file_name.is_empty() || file_name == "." || file_name == ".." {
+            return Err(invalid("Managed BLOB descriptor has an invalid file 
name"));
+        }
+        self.references.insert(ManagedBlobReference {
+            storage_root: uri[..slash].to_string(),
+            file_name: file_name.to_string(),
+        });
+        Ok(())
+    }
+
+    pub(crate) async fn write(&self, file_io: &FileIO, data_path: &str) -> 
Result<String> {
+        let name = data_path
+            .rsplit('/')
+            .next()
+            .ok_or_else(|| invalid("Managed BLOB data file has no file 
name"))?;
+        let sidecar_name = format!("{name}{REFERENCE_FILE_SUFFIX}");
+        let path = format!("{data_path}{REFERENCE_FILE_SUFFIX}");
+        let bytes = self.serialize()?;
+        let output = file_io.new_output(&path)?;
+        if let Err(error) = output.write(Bytes::from(bytes)).await {
+            let _ = file_io.delete_file(&path).await;
+            return Err(error);
+        }
+        Ok(sidecar_name)
+    }
+
+    pub(crate) fn serialize(&self) -> Result<Vec<u8>> {
+        let count = i32::try_from(self.references.len())
+            .map_err(|_| invalid("Managed BLOB reference count exceeds i32"))?;
+        let mut payload = Vec::new();
+        payload.push(VERSION);
+        payload.extend_from_slice(&count.to_be_bytes());
+        for reference in &self.references {
+            write_java_utf(&reference.storage_root, &mut payload)?;
+            write_java_utf(&reference.file_name, &mut payload)?;
+        }
+        let checksum = crc32fast::hash(&payload);
+        let mut bytes = Vec::with_capacity(4 + payload.len() + 4);
+        bytes.extend_from_slice(&MAGIC.to_be_bytes());
+        bytes.extend_from_slice(&payload);
+        bytes.extend_from_slice(&checksum.to_be_bytes());
+        Ok(bytes)
+    }
+
+    #[cfg(test)]
+    pub(crate) fn references_for_test(
+        &mut self,
+        references: impl IntoIterator<Item = ManagedBlobReference>,
+    ) {
+        self.references.extend(references);
+    }
+}
+
+fn binary_column(array: &dyn Array) -> Result<&LargeBinaryArray> {
+    array
+        .as_any()
+        .downcast_ref::<LargeBinaryArray>()
+        .ok_or_else(|| invalid("Managed BLOB descriptor column requires 
LargeBinaryArray"))
+}
+
+fn invalid(message: &str) -> crate::Error {
+    crate::Error::DataInvalid {
+        message: message.to_string(),
+        source: None,
+    }
+}
+
+/// `DataOutputStream.writeUTF` uses modified UTF-8 over UTF-16 code units.
+/// A plain Rust UTF-8 string differs for NUL and supplementary characters.
+fn write_java_utf(value: &str, output: &mut Vec<u8>) -> Result<()> {
+    let mut encoded = Vec::new();
+    for unit in value.encode_utf16() {
+        if (1..=0x7f).contains(&unit) {
+            encoded.push(unit as u8);
+        } else if unit <= 0x7ff {
+            encoded.push(0xc0 | (unit >> 6) as u8);
+            encoded.push(0x80 | (unit & 0x3f) as u8);
+        } else {
+            encoded.push(0xe0 | (unit >> 12) as u8);
+            encoded.push(0x80 | ((unit >> 6) & 0x3f) as u8);
+            encoded.push(0x80 | (unit & 0x3f) as u8);
+        }
+    }
+    let length = u16::try_from(encoded.len())
+        .map_err(|_| invalid("Managed BLOB reference path exceeds Java UTF 
length"))?;
+    output.extend_from_slice(&length.to_be_bytes());
+    output.extend_from_slice(&encoded);
+    Ok(())
+}
diff --git a/crates/paimon/src/table/managed_blob_table_tests.rs 
b/crates/paimon/src/table/managed_blob_table_tests.rs
new file mode 100644
index 00000000..7a9cd190
--- /dev/null
+++ b/crates/paimon/src/table/managed_blob_table_tests.rs
@@ -0,0 +1,1121 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+use super::managed_blob_reference::{ManagedBlobReference, 
ManagedBlobReferenceCollector};
+use super::table_write::tests::{setup_dirs, test_file_io};
+use super::{Table, TableCommit, TableWrite};
+use crate::arrow::format::create_format_reader;
+use crate::catalog::Identifier;
+use crate::io::FileIO;
+use crate::spec::{
+    bucket_path_under, ArrayType, BigIntType, BlobDescriptor, BlobType, 
DataField, DataType,
+    IntType, MapType, Schema, TableSchema, TinyIntType, VarCharType, 
SEQUENCE_NUMBER_FIELD_ID,
+    SEQUENCE_NUMBER_FIELD_NAME, VALUE_KIND_FIELD_ID, VALUE_KIND_FIELD_NAME,
+};
+use crate::spec::{Datum, PredicateBuilder};
+use arrow_array::{
+    Array, Int32Array, LargeBinaryArray, ListArray, MapArray, RecordBatch, 
StringArray, StructArray,
+};
+use arrow_buffer::{NullBuffer, OffsetBuffer, ScalarBuffer};
+use arrow_schema::{DataType as ArrowDataType, Field as ArrowField, Schema as 
ArrowSchema};
+use futures::TryStreamExt;
+use std::sync::Arc;
+
+fn scalar_batch_with_kinds(rows: &[(i32, Option<&[u8]>, i8)]) -> RecordBatch {
+    let schema = Arc::new(ArrowSchema::new(vec![
+        ArrowField::new("id", ArrowDataType::Int32, false),
+        ArrowField::new("payload", ArrowDataType::LargeBinary, true),
+        ArrowField::new(VALUE_KIND_FIELD_NAME, ArrowDataType::Int8, false),
+    ]));
+    RecordBatch::try_new(
+        schema,
+        vec![
+            Arc::new(Int32Array::from(
+                rows.iter().map(|(id, _, _)| *id).collect::<Vec<_>>(),
+            )),
+            Arc::new(LargeBinaryArray::from(
+                rows.iter().map(|(_, value, _)| *value).collect::<Vec<_>>(),
+            )),
+            Arc::new(arrow_array::Int8Array::from(
+                rows.iter().map(|(_, _, kind)| *kind).collect::<Vec<_>>(),
+            )),
+        ],
+    )
+    .unwrap()
+}
+
+fn scalar_table(file_io: &FileIO, path: &str, extra_options: &[(&str, &str)]) 
-> Table {
+    let mut builder = Schema::builder()
+        .column("id", DataType::Int(IntType::new()))
+        .column("payload", DataType::Blob(BlobType::new()))
+        .primary_key(["id"])
+        .option("bucket", "1")
+        .option("file.format", "parquet");
+    for &(key, value) in extra_options {
+        builder = builder.option(key, value);
+    }
+    let schema = builder.build().unwrap();
+    Table::new(
+        file_io.clone(),
+        Identifier::new("default", "managed_blob_test"),
+        path.to_string(),
+        TableSchema::new(0, &schema),
+        None,
+    )
+}
+
+fn nested_table(file_io: &FileIO, path: &str) -> Table {
+    let schema = Schema::builder()
+        .column("id", DataType::Int(IntType::new()))
+        .column(
+            "items",
+            DataType::Array(ArrayType::new(DataType::Blob(BlobType::new()))),
+        )
+        .column(
+            "named",
+            DataType::Map(MapType::new(
+                DataType::VarChar(VarCharType::string_type()),
+                DataType::Blob(BlobType::new()),
+            )),
+        )
+        .primary_key(["id"])
+        .option("bucket", "1")
+        .build()
+        .unwrap();
+    Table::new(
+        file_io.clone(),
+        Identifier::new("default", "managed_blob_nested"),
+        path.to_string(),
+        TableSchema::new(0, &schema),
+        None,
+    )
+}
+
+fn nested_batch(table: &Table) -> RecordBatch {
+    let schema = 
crate::arrow::build_target_arrow_schema(table.schema().fields()).unwrap();
+    let ArrowDataType::List(element) = schema.field(1).data_type() else {
+        panic!("ARRAY<BLOB> must use Arrow List");
+    };
+    let array = ListArray::try_new(
+        element.clone(),
+        OffsetBuffer::new(ScalarBuffer::from(vec![0, 3, 4, 4, 5])),
+        Arc::new(LargeBinaryArray::from(vec![
+            Some(b"alpha".as_slice()),
+            None,
+            Some(b"beta".as_slice()),
+            Some(b"hidden-array-child".as_slice()),
+            Some(b"gamma".as_slice()),
+        ])),
+        Some(NullBuffer::from(vec![true, false, true, true])),
+    )
+    .unwrap();
+    let ArrowDataType::Map(entries_field, ordered) = 
schema.field(2).data_type() else {
+        panic!("MAP<X, BLOB> must use Arrow Map");
+    };
+    let ArrowDataType::Struct(entry_fields) = entries_field.data_type() else {
+        panic!("MAP<X, BLOB> entries must be Struct");
+    };
+    let entries = StructArray::try_new(
+        entry_fields.clone(),
+        vec![
+            Arc::new(StringArray::from(vec!["one", "two", "hidden", "three"])),
+            Arc::new(LargeBinaryArray::from(vec![
+                Some(b"first".as_slice()),
+                None,
+                Some(b"hidden-map-child".as_slice()),
+                Some(b"third".as_slice()),
+            ])),
+        ],
+        None,
+    )
+    .unwrap();
+    let map = MapArray::try_new(
+        entries_field.clone(),
+        OffsetBuffer::new(ScalarBuffer::from(vec![0, 2, 3, 3, 4])),
+        entries,
+        Some(NullBuffer::from(vec![true, false, true, true])),
+        *ordered,
+    )
+    .unwrap();
+    RecordBatch::try_new(
+        schema,
+        vec![
+            Arc::new(Int32Array::from(vec![1, 2, 3, 4])),
+            Arc::new(array),
+            Arc::new(map),
+        ],
+    )
+    .unwrap()
+}
+
+fn nonnullable_nested_table(file_io: &FileIO, path: &str) -> Table {
+    let nonnullable_blob = DataType::Blob(BlobType::with_nullable(false));
+    let schema = Schema::builder()
+        .column("id", DataType::Int(IntType::new()))
+        .column(
+            "items",
+            DataType::Array(ArrayType::new(nonnullable_blob.clone())),
+        )
+        .column(
+            "named",
+            DataType::Map(MapType::new(
+                DataType::VarChar(VarCharType::string_type()),
+                nonnullable_blob,
+            )),
+        )
+        .primary_key(["id"])
+        .option("bucket", "1")
+        .build()
+        .unwrap();
+    Table::new(
+        file_io.clone(),
+        Identifier::new("default", "managed_blob_nonnullable_children"),
+        path.to_string(),
+        TableSchema::new(0, &schema),
+        None,
+    )
+}
+
+fn nonnullable_nested_batch(table: &Table) -> RecordBatch {
+    let schema = 
crate::arrow::build_target_arrow_schema(table.schema().fields()).unwrap();
+    let ArrowDataType::List(element) = schema.field(1).data_type() else {
+        panic!("ARRAY<BLOB NOT NULL> must use Arrow List");
+    };
+    assert!(!element.is_nullable());
+    let array = ListArray::try_new(
+        element.clone(),
+        OffsetBuffer::new(ScalarBuffer::from(vec![0, 1, 2])),
+        Arc::new(LargeBinaryArray::from(vec![
+            Some(b"array-one".as_slice()),
+            Some(b"array-two".as_slice()),
+        ])),
+        None,
+    )
+    .unwrap();
+    let ArrowDataType::Map(entries_field, ordered) = 
schema.field(2).data_type() else {
+        panic!("MAP<STRING, BLOB NOT NULL> must use Arrow Map");
+    };
+    let ArrowDataType::Struct(entry_fields) = entries_field.data_type() else {
+        panic!("MAP entries must be Struct");
+    };
+    assert!(!entry_fields[1].is_nullable());
+    let entries = StructArray::try_new(
+        entry_fields.clone(),
+        vec![
+            Arc::new(StringArray::from(vec!["one", "two"])),
+            Arc::new(LargeBinaryArray::from(vec![
+                Some(b"map-one".as_slice()),
+                Some(b"map-two".as_slice()),
+            ])),
+        ],
+        None,
+    )
+    .unwrap();
+    let map = MapArray::try_new(
+        entries_field.clone(),
+        OffsetBuffer::new(ScalarBuffer::from(vec![0, 1, 2])),
+        entries,
+        None,
+        *ordered,
+    )
+    .unwrap();
+    RecordBatch::try_new(
+        schema,
+        vec![
+            Arc::new(Int32Array::from(vec![1, 2])),
+            Arc::new(array),
+            Arc::new(map),
+        ],
+    )
+    .unwrap()
+}
+
+fn array_values(array: &ListArray, row: usize) -> Option<Vec<Option<Vec<u8>>>> 
{
+    if array.is_null(row) {
+        return None;
+    }
+    let values = array.value(row);
+    let values = values.as_any().downcast_ref::<LargeBinaryArray>().unwrap();
+    Some(
+        (0..values.len())
+            .map(|index| values.is_valid(index).then(|| 
values.value(index).to_vec()))
+            .collect(),
+    )
+}
+
+fn map_values(array: &MapArray, row: usize) -> Option<Vec<(String, 
Option<Vec<u8>>)>> {
+    if array.is_null(row) {
+        return None;
+    }
+    let entries = array.value(row);
+    let keys = entries
+        .column(0)
+        .as_any()
+        .downcast_ref::<StringArray>()
+        .unwrap();
+    let values = entries
+        .column(1)
+        .as_any()
+        .downcast_ref::<LargeBinaryArray>()
+        .unwrap();
+    Some(
+        (0..entries.len())
+            .map(|index| {
+                (
+                    keys.value(index).to_string(),
+                    values.is_valid(index).then(|| 
values.value(index).to_vec()),
+                )
+            })
+            .collect(),
+    )
+}
+
+fn scalar_batch(rows: &[(i32, Option<&[u8]>)]) -> RecordBatch {
+    let schema = Arc::new(ArrowSchema::new(vec![
+        ArrowField::new("id", ArrowDataType::Int32, false),
+        ArrowField::new("payload", ArrowDataType::LargeBinary, true),
+    ]));
+    RecordBatch::try_new(
+        schema,
+        vec![
+            Arc::new(Int32Array::from(
+                rows.iter().map(|(id, _)| *id).collect::<Vec<_>>(),
+            )),
+            Arc::new(LargeBinaryArray::from(
+                rows.iter().map(|(_, value)| *value).collect::<Vec<_>>(),
+            )),
+        ],
+    )
+    .unwrap()
+}
+
+async fn read_scalar_rows(table: &Table) -> Vec<(i32, Option<Vec<u8>>)> {
+    let builder = table.new_read_builder();
+    let plan = builder.new_scan().plan().await.unwrap();
+    let reader = builder.new_read().unwrap();
+    let batches: Vec<RecordBatch> = reader
+        .to_arrow(plan.splits())
+        .unwrap()
+        .try_collect()
+        .await
+        .unwrap();
+    let mut rows = batches
+        .iter()
+        .flat_map(|batch| {
+            let ids = batch
+                .column(0)
+                .as_any()
+                .downcast_ref::<Int32Array>()
+                .unwrap();
+            let payloads = batch
+                .column(1)
+                .as_any()
+                .downcast_ref::<LargeBinaryArray>()
+                .unwrap();
+            (0..batch.num_rows()).map(|row| {
+                (
+                    ids.value(row),
+                    payloads.is_valid(row).then(|| 
payloads.value(row).to_vec()),
+                )
+            })
+        })
+        .collect::<Vec<_>>();
+    rows.sort_by_key(|(id, _)| *id);
+    rows
+}
+
+#[tokio::test]
+async fn primary_key_blob_is_externalized_and_read_back() {
+    let file_io = test_file_io();
+    let path = "memory:/managed_blob_pk_scalar";
+    setup_dirs(&file_io, path).await;
+    let table = scalar_table(&file_io, path, &[]);
+
+    let mut writer = TableWrite::new(&table, "test-user".to_string()).unwrap();
+    writer
+        .write_arrow_batch(&scalar_batch(&[
+            (2, Some(b"world")),
+            (1, Some(b"hello")),
+            (3, None),
+        ]))
+        .await
+        .unwrap();
+    let messages = writer.prepare_commit().await.unwrap();
+    assert_eq!(messages.len(), 1);
+    assert_eq!(messages[0].new_files.len(), 1);
+    let file = &messages[0].new_files[0];
+    assert!(file.file_name.ends_with(".parquet"));
+    assert_eq!(
+        file.extra_files,
+        vec![format!("{}.blobref", file.file_name)]
+    );
+    let bucket_dir = bucket_path_under(path, "", messages[0].bucket);
+    let file_path = format!("{bucket_dir}/{}", file.file_name);
+    let sidecar = file_io
+        .new_input(&format!("{file_path}.blobref"))
+        .unwrap()
+        .read()
+        .await
+        .unwrap();
+    assert_eq!(&sidecar[..4], &0x50424c52_i32.to_be_bytes());
+    assert_eq!(sidecar[4], 1);
+    assert_eq!(i32::from_be_bytes(sidecar[5..9].try_into().unwrap()), 1);
+
+    // The physical Parquet column contains descriptors, not the BLOB payload.
+    let fields = vec![
+        DataField::new(
+            SEQUENCE_NUMBER_FIELD_ID,
+            SEQUENCE_NUMBER_FIELD_NAME.to_string(),
+            DataType::BigInt(BigIntType::new()),
+        ),
+        DataField::new(
+            VALUE_KIND_FIELD_ID,
+            VALUE_KIND_FIELD_NAME.to_string(),
+            DataType::TinyInt(TinyIntType::new()),
+        ),
+        table.schema().fields()[0].clone(),
+        table.schema().fields()[1].clone(),
+    ];
+    let format_reader = create_format_reader(&file_path, false, 
&fields).unwrap();
+    let input = file_io.new_input(&file_path).unwrap();
+    let batches: Vec<RecordBatch> = format_reader
+        .read_batch_stream(
+            Box::new(input.reader().await.unwrap()),
+            file.file_size as u64,
+            &fields,
+            None,
+            None,
+            None,
+        )
+        .await
+        .unwrap()
+        .try_collect()
+        .await
+        .unwrap();
+    let physical = batches[0]
+        .column(3)
+        .as_any()
+        .downcast_ref::<LargeBinaryArray>()
+        .unwrap();
+    let mut descriptors = Vec::new();
+    for row in 0..physical.len() {
+        if physical.is_valid(row) {
+            let descriptor = 
BlobDescriptor::deserialize(physical.value(row)).unwrap();
+            assert!(descriptor.uri().ends_with(".managed.blob"));
+            descriptors.push(descriptor);
+        }
+    }
+    assert_eq!(descriptors.len(), 2);
+
+    TableCommit::new(table.clone(), "test-user".to_string())
+        .commit(messages)
+        .await
+        .unwrap();
+    assert_eq!(
+        read_scalar_rows(&table).await,
+        vec![
+            (1, Some(b"hello".to_vec())),
+            (2, Some(b"world".to_vec())),
+            (3, None),
+        ]
+    );
+}
+
+#[tokio::test]
+async fn 
primary_key_array_and_map_blob_values_round_trip_through_managed_packs() {
+    let file_io = test_file_io();
+    let path = "memory:/managed_blob_pk_nested";
+    setup_dirs(&file_io, path).await;
+    let table = nested_table(&file_io, path);
+    let mut writer = TableWrite::new(&table, "test-user".to_string()).unwrap();
+    writer
+        .write_arrow_batch(&nested_batch(&table))
+        .await
+        .unwrap();
+    let messages = writer.prepare_commit().await.unwrap();
+    let data_file = &messages[0].new_files[0];
+    assert_eq!(
+        data_file.extra_files,
+        vec![format!("{}.blobref", data_file.file_name)]
+    );
+    let sidecar_path = format!(
+        "{}/{}.blobref",
+        bucket_path_under(path, "", messages[0].bucket),
+        data_file.file_name
+    );
+    let sidecar = file_io
+        .new_input(&sidecar_path)
+        .unwrap()
+        .read()
+        .await
+        .unwrap();
+    // One pack per managed field, with multiple values inside each pack.
+    assert_eq!(i32::from_be_bytes(sidecar[5..9].try_into().unwrap()), 2);
+    TableCommit::new(table.clone(), "test-user".to_string())
+        .commit(messages)
+        .await
+        .unwrap();
+
+    let builder = table.new_read_builder();
+    let plan = builder.new_scan().plan().await.unwrap();
+    let read = builder.new_read().unwrap();
+    let batches: Vec<RecordBatch> = read
+        .to_arrow(plan.splits())
+        .unwrap()
+        .try_collect()
+        .await
+        .unwrap();
+    assert_eq!(batches.iter().map(RecordBatch::num_rows).sum::<usize>(), 4);
+    let mut rows = Vec::new();
+    for batch in &batches {
+        let ids = batch
+            .column(0)
+            .as_any()
+            .downcast_ref::<Int32Array>()
+            .unwrap();
+        let arrays = batch
+            .column(1)
+            .as_any()
+            .downcast_ref::<ListArray>()
+            .unwrap();
+        let maps = 
batch.column(2).as_any().downcast_ref::<MapArray>().unwrap();
+        for row in 0..batch.num_rows() {
+            rows.push((
+                ids.value(row),
+                array_values(arrays, row),
+                map_values(maps, row),
+            ));
+        }
+    }
+    rows.sort_by_key(|(id, _, _)| *id);
+    assert_eq!(
+        rows,
+        vec![
+            (
+                1,
+                Some(vec![Some(b"alpha".to_vec()), None, 
Some(b"beta".to_vec())]),
+                Some(vec![
+                    ("one".to_string(), Some(b"first".to_vec())),
+                    ("two".to_string(), None)
+                ]),
+            ),
+            (2, None, None),
+            (3, Some(Vec::new()), Some(Vec::new())),
+            (
+                4,
+                Some(vec![Some(b"gamma".to_vec())]),
+                Some(vec![("three".to_string(), Some(b"third".to_vec()))]),
+            ),
+        ]
+    );
+}
+
+#[tokio::test]
+async fn 
sliced_primary_key_blob_collections_only_externalize_visible_children() {
+    for (suffix, columns) in [
+        ("array", vec![0, 1]),
+        ("map", vec![0, 2]),
+        ("both", vec![0, 1, 2]),
+    ] {
+        for start in [0, 3] {
+            let file_io = test_file_io();
+            let path = format!("memory:/managed_blob_sliced_{suffix}_{start}");
+            setup_dirs(&file_io, &path).await;
+            let template = nested_table(&file_io, &path);
+            let table = if columns.len() == 3 {
+                template.clone()
+            } else {
+                let field = &template.schema().fields()[columns[1]];
+                let schema = Schema::builder()
+                    .column("id", DataType::Int(IntType::new()))
+                    .column(field.name(), field.data_type().clone())
+                    .primary_key(["id"])
+                    .option("bucket", "1")
+                    .build()
+                    .unwrap();
+                Table::new(
+                    file_io.clone(),
+                    Identifier::new("default", "managed_blob_sliced"),
+                    path.clone(),
+                    TableSchema::new(0, &schema),
+                    None,
+                )
+            };
+            let input = nested_batch(&template)
+                .project(&columns)
+                .unwrap()
+                .slice(start, 1);
+            let mut writer = TableWrite::new(&table, 
"test-user".to_string()).unwrap();
+            writer.write_arrow_batch(&input).await.unwrap();
+            let messages = writer.prepare_commit().await.unwrap();
+            TableCommit::new(table.clone(), "test-user".to_string())
+                .commit(messages)
+                .await
+                .unwrap();
+            let mut builder = table.new_read_builder();
+            let projection = if columns == [0, 1] {
+                vec!["id", "items"]
+            } else if columns == [0, 2] {
+                vec!["id", "named"]
+            } else {
+                vec!["id", "items", "named"]
+            };
+            builder.with_projection(&projection).unwrap();
+            let plan = builder.new_scan().plan().await.unwrap();
+            let batches: Vec<RecordBatch> = builder
+                .new_read()
+                .unwrap()
+                .to_arrow(plan.splits())
+                .unwrap()
+                .try_collect()
+                .await
+                .unwrap();
+            
assert_eq!(batches.iter().map(RecordBatch::num_rows).sum::<usize>(), 1);
+            let batch = &batches[0];
+            assert_eq!(
+                batch
+                    .column(0)
+                    .as_any()
+                    .downcast_ref::<Int32Array>()
+                    .unwrap()
+                    .value(0),
+                start as i32 + 1
+            );
+            if let Some(index) = projection.iter().position(|name| *name == 
"items") {
+                let array = batch
+                    .column(index)
+                    .as_any()
+                    .downcast_ref::<ListArray>()
+                    .unwrap();
+                let expected = if start == 0 {
+                    vec![Some(b"alpha".to_vec()), None, Some(b"beta".to_vec())]
+                } else {
+                    vec![Some(b"gamma".to_vec())]
+                };
+                assert_eq!(array_values(array, 0), Some(expected));
+            }
+            if let Some(index) = projection.iter().position(|name| *name == 
"named") {
+                let map = batch
+                    .column(index)
+                    .as_any()
+                    .downcast_ref::<MapArray>()
+                    .unwrap();
+                let expected = if start == 0 {
+                    vec![
+                        ("one".to_string(), Some(b"first".to_vec())),
+                        ("two".to_string(), None),
+                    ]
+                } else {
+                    vec![("three".to_string(), Some(b"third".to_vec()))]
+                };
+                assert_eq!(map_values(map, 0), Some(expected));
+            }
+        }
+    }
+}
+
+#[tokio::test]
+async fn nonnullable_blob_collection_children_survive_sliced_writes() {
+    for start in [0, 1] {
+        let file_io = test_file_io();
+        let path = format!("memory:/managed_blob_nonnullable_sliced_{start}");
+        setup_dirs(&file_io, &path).await;
+        let table = nonnullable_nested_table(&file_io, &path);
+        let batch = nonnullable_nested_batch(&table).slice(start, 1);
+        let mut writer = TableWrite::new(&table, 
"test-user".to_string()).unwrap();
+        writer.write_arrow_batch(&batch).await.unwrap();
+        let messages = writer.prepare_commit().await.unwrap();
+        TableCommit::new(table.clone(), "test-user".to_string())
+            .commit(messages)
+            .await
+            .unwrap();
+
+        let builder = table.new_read_builder();
+        let plan = builder.new_scan().plan().await.unwrap();
+        let batches: Vec<RecordBatch> = builder
+            .new_read()
+            .unwrap()
+            .to_arrow(plan.splits())
+            .unwrap()
+            .try_collect()
+            .await
+            .unwrap();
+        assert_eq!(batches.iter().map(RecordBatch::num_rows).sum::<usize>(), 
1);
+        let batch = &batches[0];
+        let items = batch
+            .column(1)
+            .as_any()
+            .downcast_ref::<ListArray>()
+            .unwrap();
+        let named = 
batch.column(2).as_any().downcast_ref::<MapArray>().unwrap();
+        let expected_array = if start == 0 {
+            b"array-one"
+        } else {
+            b"array-two"
+        };
+        let expected_map = if start == 0 { b"map-one" } else { b"map-two" };
+        assert_eq!(
+            array_values(items, 0),
+            Some(vec![Some(expected_array.to_vec())])
+        );
+        assert_eq!(
+            map_values(named, 0),
+            Some(vec![(
+                if start == 0 { "one" } else { "two" }.to_string(),
+                Some(expected_map.to_vec()),
+            )])
+        );
+    }
+}
+
+#[tokio::test]
+async fn nonnullable_blob_collection_children_survive_retractions() {
+    for start in [0, 1] {
+        let file_io = test_file_io();
+        let path = format!("memory:/managed_blob_nonnullable_delete_{start}");
+        setup_dirs(&file_io, &path).await;
+        let table = nonnullable_nested_table(&file_io, &path);
+        let input = nonnullable_nested_batch(&table).slice(start, 1);
+        let mut fields = 
input.schema().fields().iter().cloned().collect::<Vec<_>>();
+        fields[1] = Arc::new(fields[1].as_ref().clone().with_nullable(false));
+        fields[2] = Arc::new(fields[2].as_ref().clone().with_nullable(false));
+        fields.push(Arc::new(ArrowField::new(
+            VALUE_KIND_FIELD_NAME,
+            ArrowDataType::Int8,
+            false,
+        )));
+        let mut columns = input.columns().to_vec();
+        columns.push(Arc::new(arrow_array::Int8Array::from(vec![3])));
+        let delete = RecordBatch::try_new(Arc::new(ArrowSchema::new(fields)), 
columns).unwrap();
+        let mut writer = TableWrite::new(&table, 
"test-user".to_string()).unwrap();
+        writer.write_arrow_batch(&delete).await.unwrap();
+        let messages = writer.prepare_commit().await.unwrap();
+        TableCommit::new(table.clone(), "test-user".to_string())
+            .commit(messages)
+            .await
+            .unwrap();
+    }
+}
+
+#[tokio::test]
+async fn primary_key_blob_payload_filter_runs_after_resolution() {
+    let file_io = test_file_io();
+    let path = "memory:/managed_blob_payload_filter";
+    setup_dirs(&file_io, path).await;
+    let table = scalar_table(&file_io, path, &[]);
+    let mut writer = TableWrite::new(&table, "test-user".to_string()).unwrap();
+    writer
+        .write_arrow_batch(&scalar_batch(&[(1, Some(b"hello")), (2, 
Some(b"world"))]))
+        .await
+        .unwrap();
+    let messages = writer.prepare_commit().await.unwrap();
+    TableCommit::new(table.clone(), "test-user".to_string())
+        .commit(messages)
+        .await
+        .unwrap();
+
+    let predicate = PredicateBuilder::new(table.schema().fields())
+        .equal("payload", Datum::Bytes(b"hello".to_vec()))
+        .unwrap();
+    for projection in [vec!["id", "payload"], vec!["id"]] {
+        let mut builder = table.new_read_builder();
+        builder.with_projection(&projection).unwrap();
+        builder.with_filter(predicate.clone());
+        let plan = builder.new_scan().plan().await.unwrap();
+        assert!(!plan.splits().is_empty());
+        let batches: Vec<RecordBatch> = builder
+            .new_read()
+            .unwrap()
+            .to_arrow(plan.splits())
+            .unwrap()
+            .try_collect()
+            .await
+            .unwrap();
+        assert_eq!(batches.iter().map(RecordBatch::num_rows).sum::<usize>(), 
1);
+        assert_eq!(batches[0].num_columns(), projection.len());
+        assert_eq!(
+            batches[0]
+                .column(0)
+                .as_any()
+                .downcast_ref::<Int32Array>()
+                .unwrap()
+                .value(0),
+            1
+        );
+    }
+}
+
+#[tokio::test]
+async fn primary_key_blob_retract_accepts_non_nullable_arrow_inputs() {
+    let file_io = test_file_io();
+    let scalar_path = "memory:/managed_blob_nonnullable_retract_scalar";
+    setup_dirs(&file_io, scalar_path).await;
+    let scalar_table = scalar_table(&file_io, scalar_path, &[]);
+    let scalar = scalar_batch_with_kinds(&[(1, Some(b"hello"), 3)]);
+    let scalar_schema = Arc::new(ArrowSchema::new(vec![
+        scalar.schema().field(0).clone(),
+        ArrowField::new("payload", ArrowDataType::LargeBinary, false),
+        scalar.schema().field(2).clone(),
+    ]));
+    let scalar = RecordBatch::try_new(scalar_schema, 
scalar.columns().to_vec()).unwrap();
+    let mut writer = TableWrite::new(&scalar_table, 
"test-user".to_string()).unwrap();
+    writer.write_arrow_batch(&scalar).await.unwrap();
+    writer.prepare_commit().await.unwrap();
+
+    let nested_path = "memory:/managed_blob_nonnullable_retract_nested";
+    setup_dirs(&file_io, nested_path).await;
+    let nested_table = nested_table(&file_io, nested_path);
+    let nested = nested_batch(&nested_table).slice(0, 1);
+    let nested_schema = Arc::new(ArrowSchema::new(vec![
+        nested.schema().field(0).clone(),
+        nested.schema().field(1).clone().with_nullable(false),
+        nested.schema().field(2).clone().with_nullable(false),
+        ArrowField::new(VALUE_KIND_FIELD_NAME, ArrowDataType::Int8, false),
+    ]));
+    let mut columns = nested.columns().to_vec();
+    columns.push(Arc::new(arrow_array::Int8Array::from(vec![3])));
+    let nested = RecordBatch::try_new(nested_schema, columns).unwrap();
+    let mut writer = TableWrite::new(&nested_table, 
"test-user".to_string()).unwrap();
+    writer.write_arrow_batch(&nested).await.unwrap();
+    writer.prepare_commit().await.unwrap();
+}
+
+#[tokio::test]
+async fn primary_key_blob_updates_and_deletes_follow_merge_order() {
+    let file_io = test_file_io();
+    let path = "memory:/managed_blob_pk_updates";
+    setup_dirs(&file_io, path).await;
+    let table = scalar_table(&file_io, path, &[]);
+
+    for batch in [
+        scalar_batch(&[(1, Some(b"first")), (2, Some(b"survivor"))]),
+        scalar_batch(&[(1, Some(b"second"))]),
+    ] {
+        let mut writer = TableWrite::new(&table, 
"test-user".to_string()).unwrap();
+        writer.write_arrow_batch(&batch).await.unwrap();
+        let messages = writer.prepare_commit().await.unwrap();
+        TableCommit::new(table.clone(), "test-user".to_string())
+            .commit(messages)
+            .await
+            .unwrap();
+    }
+    assert_eq!(
+        read_scalar_rows(&table).await,
+        vec![
+            (1, Some(b"second".to_vec())),
+            (2, Some(b"survivor".to_vec()))
+        ]
+    );
+
+    let mut writer = TableWrite::new(&table, "test-user".to_string()).unwrap();
+    writer
+        .write_arrow_batch(&scalar_batch_with_kinds(&[(1, Some(b"ignored"), 
3)]))
+        .await
+        .unwrap();
+    let messages = writer.prepare_commit().await.unwrap();
+    assert_eq!(messages.len(), 1);
+    assert_eq!(messages[0].new_files.len(), 1);
+    let delete_file = &messages[0].new_files[0];
+    assert_eq!(delete_file.delete_row_count, Some(1));
+    let sidecar_path = format!(
+        "{}/{}.blobref",
+        bucket_path_under(path, "", messages[0].bucket),
+        delete_file.file_name
+    );
+    let sidecar = file_io
+        .new_input(&sidecar_path)
+        .unwrap()
+        .read()
+        .await
+        .unwrap();
+    assert_eq!(i32::from_be_bytes(sidecar[5..9].try_into().unwrap()), 0);
+    TableCommit::new(table.clone(), "test-user".to_string())
+        .commit(messages)
+        .await
+        .unwrap();
+    assert_eq!(
+        read_scalar_rows(&table).await,
+        vec![(2, Some(b"survivor".to_vec()))]
+    );
+}
+
+#[tokio::test]
+async fn primary_key_blob_first_row_keeps_the_earliest_value() {
+    let file_io = test_file_io();
+    let path = "memory:/managed_blob_pk_first_row";
+    setup_dirs(&file_io, path).await;
+    let table = scalar_table(&file_io, path, &[("merge-engine", "first-row")]);
+    let mut writer = TableWrite::new(&table, "test-user".to_string()).unwrap();
+    writer
+        .write_arrow_batch(&scalar_batch(&[
+            (2, Some(b"other")),
+            (1, Some(b"first")),
+            (1, Some(b"second")),
+        ]))
+        .await
+        .unwrap();
+    let messages = writer.prepare_commit().await.unwrap();
+    assert_eq!(messages[0].new_files[0].row_count, 2);
+    TableCommit::new(table.clone(), "test-user".to_string())
+        .commit(messages)
+        .await
+        .unwrap();
+    // First-row scans normally skip level-0 files until compaction, like Java.
+    // Inspect the just-written files explicitly to test the writer's merge.
+    let builder = table.new_read_builder();
+    let plan = builder
+        .new_scan()
+        .with_scan_all_files()
+        .plan()
+        .await
+        .unwrap();
+    let batches: Vec<RecordBatch> = builder
+        .new_read()
+        .unwrap()
+        .to_arrow(plan.splits())
+        .unwrap()
+        .try_collect()
+        .await
+        .unwrap();
+    let mut rows = batches
+        .iter()
+        .flat_map(|batch| {
+            let ids = batch
+                .column(0)
+                .as_any()
+                .downcast_ref::<Int32Array>()
+                .unwrap();
+            let payloads = batch
+                .column(1)
+                .as_any()
+                .downcast_ref::<LargeBinaryArray>()
+                .unwrap();
+            (0..batch.num_rows()).map(|row| (ids.value(row), 
payloads.value(row).to_vec()))
+        })
+        .collect::<Vec<_>>();
+    rows.sort_by_key(|(id, _)| *id);
+    assert_eq!(rows, vec![(1, b"first".to_vec()), (2, b"other".to_vec())]);
+}
+
+#[tokio::test]
+async fn primary_key_blob_partial_update_keeps_the_last_non_null_value() {
+    let file_io = test_file_io();
+    let path = "memory:/managed_blob_pk_partial_update";
+    setup_dirs(&file_io, path).await;
+    let table = scalar_table(&file_io, path, &[("merge-engine", 
"partial-update")]);
+    let mut writer = TableWrite::new(&table, "test-user".to_string()).unwrap();
+    writer
+        .write_arrow_batch(&scalar_batch(&[
+            (1, Some(b"first")),
+            (1, None),
+            (2, Some(b"other")),
+        ]))
+        .await
+        .unwrap();
+    let messages = writer.prepare_commit().await.unwrap();
+    assert_eq!(messages[0].new_files[0].row_count, 2);
+    TableCommit::new(table.clone(), "test-user".to_string())
+        .commit(messages)
+        .await
+        .unwrap();
+    assert_eq!(
+        read_scalar_rows(&table).await,
+        vec![(1, Some(b"first".to_vec())), (2, Some(b"other".to_vec()))]
+    );
+}
+
+#[tokio::test]
+async fn primary_key_blob_projection_and_descriptor_mode() {
+    let file_io = test_file_io();
+    let path = "memory:/managed_blob_pk_projection";
+    setup_dirs(&file_io, path).await;
+    let table = scalar_table(&file_io, path, &[]);
+    let mut writer = TableWrite::new(&table, "test-user".to_string()).unwrap();
+    writer
+        .write_arrow_batch(&scalar_batch(&[(1, Some(b"one")), (2, None)]))
+        .await
+        .unwrap();
+    let messages = writer.prepare_commit().await.unwrap();
+    TableCommit::new(table.clone(), "test-user".to_string())
+        .commit(messages)
+        .await
+        .unwrap();
+
+    let mut projected = table.new_read_builder();
+    projected.with_projection(&["payload"]).unwrap();
+    let plan = projected.new_scan().plan().await.unwrap();
+    let batches: Vec<RecordBatch> = projected
+        .new_read()
+        .unwrap()
+        .to_arrow(plan.splits())
+        .unwrap()
+        .try_collect()
+        .await
+        .unwrap();
+    assert_eq!(batches.iter().map(RecordBatch::num_rows).sum::<usize>(), 2);
+    let values = batches
+        .iter()
+        .flat_map(|batch| {
+            let col = batch
+                .column(0)
+                .as_any()
+                .downcast_ref::<LargeBinaryArray>()
+                .unwrap();
+            (0..col.len()).map(|row| col.is_valid(row).then(|| 
col.value(row).to_vec()))
+        })
+        .collect::<Vec<_>>();
+    assert_eq!(values, vec![Some(b"one".to_vec()), None]);
+
+    let descriptor_table = scalar_table(&file_io, path, 
&[("blob-as-descriptor", "true")]);
+    let descriptors = read_scalar_rows(&descriptor_table).await;
+    assert_eq!(descriptors[0].0, 1);
+    let descriptor = 
BlobDescriptor::deserialize(descriptors[0].1.as_ref().unwrap()).unwrap();
+    assert!(descriptor.uri().ends_with(".managed.blob"));
+    assert_eq!(descriptors[1], (2, None));
+}
+
+#[tokio::test]
+async fn primary_key_blob_rolls_packs_and_records_every_reachable_pack() {
+    let file_io = test_file_io();
+    let path = "memory:/managed_blob_pk_roll";
+    setup_dirs(&file_io, path).await;
+    let table = scalar_table(&file_io, path, &[("blob.target-file-size", 
"1b")]);
+    let mut writer = TableWrite::new(&table, "test-user".to_string()).unwrap();
+    writer
+        .write_arrow_batch(&scalar_batch(&[
+            (1, Some(b"a")),
+            (2, Some(b"bb")),
+            (3, Some(b"ccc")),
+        ]))
+        .await
+        .unwrap();
+    let messages = writer.prepare_commit().await.unwrap();
+    let file = &messages[0].new_files[0];
+    let sidecar_path = format!(
+        "{}/{}.blobref",
+        bucket_path_under(path, "", messages[0].bucket),
+        file.file_name
+    );
+    let sidecar = file_io
+        .new_input(&sidecar_path)
+        .unwrap()
+        .read()
+        .await
+        .unwrap();
+    assert_eq!(i32::from_be_bytes(sidecar[5..9].try_into().unwrap()), 3);
+    let managed_packs = file_io
+        .list_status(&format!("{}/", bucket_path_under(path, "", 0)))
+        .await
+        .unwrap()
+        .into_iter()
+        .filter(|entry| entry.path.ends_with(".managed.blob"))
+        .collect::<Vec<_>>();
+    assert_eq!(managed_packs.len(), 3);
+    TableCommit::new(table.clone(), "test-user".to_string())
+        .commit(messages)
+        .await
+        .unwrap();
+    assert_eq!(
+        read_scalar_rows(&table).await,
+        vec![
+            (1, Some(b"a".to_vec())),
+            (2, Some(b"bb".to_vec())),
+            (3, Some(b"ccc".to_vec())),
+        ]
+    );
+}
+
+#[tokio::test]
+async fn primary_key_blob_copies_input_descriptor_into_managed_pack() {
+    let file_io = test_file_io();
+    let path = "memory:/managed_blob_pk_descriptor_copy";
+    setup_dirs(&file_io, path).await;
+    let source = "memory:/managed_blob_source.bin";
+    file_io
+        .new_output(source)
+        .unwrap()
+        .write(bytes::Bytes::from_static(b"prefixPAYLOADsuffix"))
+        .await
+        .unwrap();
+    let input = BlobDescriptor::new(source.to_string(), 6, 7).serialize();
+    let table = scalar_table(&file_io, path, &[]);
+    let mut writer = TableWrite::new(&table, "test-user".to_string()).unwrap();
+    writer
+        .write_arrow_batch(&scalar_batch(&[(1, Some(&input))]))
+        .await
+        .unwrap();
+    let messages = writer.prepare_commit().await.unwrap();
+    TableCommit::new(table.clone(), "test-user".to_string())
+        .commit(messages)
+        .await
+        .unwrap();
+    // Removing the source proves the table now refers to its managed copy.
+    file_io.delete_file(source).await.unwrap();
+    assert_eq!(
+        read_scalar_rows(&table).await,
+        vec![(1, Some(b"PAYLOAD".to_vec()))]
+    );
+}
+
+#[tokio::test]
+async fn closing_uncommitted_primary_key_writer_removes_managed_packs() {
+    let file_io = test_file_io();
+    let path = "memory:/managed_blob_pk_abort";
+    setup_dirs(&file_io, path).await;
+    let table = scalar_table(&file_io, path, &[("blob.target-file-size", 
"1b")]);
+    let mut writer = TableWrite::new(&table, "test-user".to_string()).unwrap();
+    writer
+        .write_arrow_batch(&scalar_batch(&[(1, Some(b"a")), (2, Some(b"b"))]))
+        .await
+        .unwrap();
+    writer.close().await;
+    let files = file_io
+        .list_status(&format!("{}/", bucket_path_under(path, "", 0)))
+        .await
+        .unwrap();
+    assert!(files
+        .iter()
+        .all(|entry| !entry.path.ends_with(".managed.blob")));
+}
+
+#[test]
+fn blobref_serialization_matches_java_modified_utf_and_checksum() {
+    let mut refs = ManagedBlobReferenceCollector::default();
+    refs.references_for_test([
+        ManagedBlobReference {
+            storage_root: "memory:/é/😀".to_string(),
+            file_name: "a.managed.blob".to_string(),
+        },
+        ManagedBlobReference {
+            storage_root: "memory:/a".to_string(),
+            file_name: "z.managed.blob".to_string(),
+        },
+    ]);
+    let bytes = refs.serialize().unwrap();
+    assert_eq!(&bytes[..4], &0x50424c52_i32.to_be_bytes());
+    assert_eq!(bytes[4], 1);
+    assert_eq!(i32::from_be_bytes(bytes[5..9].try_into().unwrap()), 2);
+    let payload_end = bytes.len() - 4;
+    assert_eq!(
+        &bytes[payload_end..],
+        &crc32fast::hash(&bytes[4..payload_end]).to_be_bytes()
+    );
+    // Java's writeUTF encodes supplementary code points as two surrogate 
units.
+    assert!(bytes
+        .windows(6)
+        .any(|window| window == [0xed, 0xa0, 0xbd, 0xed, 0xb8, 0x80]));
+}
diff --git a/crates/paimon/src/table/managed_blob_writer.rs 
b/crates/paimon/src/table/managed_blob_writer.rs
new file mode 100644
index 00000000..886a58af
--- /dev/null
+++ b/crates/paimon/src/table/managed_blob_writer.rs
@@ -0,0 +1,467 @@
+// 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.
+
+//! Externalize primary-key BLOB values before they enter the merge buffer.
+//!
+//! Java's `PrimaryKeyBlobExternalizer` writes scalar, array-element, and map-
+//! value BLOBs into per-field `.managed.blob` packs. The buffered row carries
+//! descriptors, so a merge or compaction never copies large payloads into
+//! Parquet. A separate `.blobref` sidecar records the packs retained by each
+//! physical data file.
+
+use super::kv_file_writer::KeyValueWriteConfig;
+use crate::arrow::format::blob::BlobFormatWriter;
+use crate::arrow::format::FormatFileWriter;
+use crate::io::FileIO;
+use crate::spec::{bucket_path_under, BlobDescriptor, CoreOptions, DataField, 
DataType, RowKind};
+use crate::Result;
+use arrow_array::builder::LargeBinaryBuilder;
+use arrow_array::{
+    Array, ArrayRef, Int8Array, LargeBinaryArray, ListArray, MapArray, 
RecordBatch, StructArray,
+    UInt32Array,
+};
+use arrow_buffer::{NullBuffer, OffsetBuffer, ScalarBuffer};
+use arrow_schema::{DataType as ArrowDataType, Schema as ArrowSchema};
+use arrow_select::take::take;
+use std::sync::Arc;
+
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+pub(crate) enum ManagedBlobKind {
+    Scalar,
+    Array,
+    Map,
+}
+
+pub(crate) fn managed_blob_kind(data_type: &DataType) -> 
Option<ManagedBlobKind> {
+    match data_type {
+        DataType::Blob(_) => Some(ManagedBlobKind::Scalar),
+        DataType::Array(array) if matches!(array.element_type(), 
DataType::Blob(_)) => {
+            Some(ManagedBlobKind::Array)
+        }
+        DataType::Map(map) if matches!(map.value_type(), DataType::Blob(_)) => 
{
+            Some(ManagedBlobKind::Map)
+        }
+        _ => None,
+    }
+}
+
+pub(crate) fn managed_blob_fields(
+    fields: &[DataField],
+    options: &CoreOptions<'_>,
+) -> Vec<(usize, ManagedBlobKind)> {
+    let inline = options.blob_inline_fields();
+    fields
+        .iter()
+        .enumerate()
+        .filter(|(_, field)| !inline.contains(field.name()))
+        .filter_map(|(index, field)| 
managed_blob_kind(field.data_type()).map(|kind| (index, kind)))
+        .collect()
+}
+
+struct ManagedBlobField {
+    index: usize,
+    kind: ManagedBlobKind,
+    current: Option<ManagedBlobPack>,
+}
+
+struct ManagedBlobPack {
+    path: String,
+    writer: Box<BlobFormatWriter>,
+}
+
+pub(crate) struct ManagedBlobWriter {
+    file_io: FileIO,
+    bucket_dir: String,
+    file_prefix: String,
+    target_file_size: u64,
+    fields: Vec<ManagedBlobField>,
+    uncommitted_paths: Vec<String>,
+}
+
+/// Optional managed-pack state kept off the ordinary key-value write path.
+pub(crate) struct ManagedBlobWriteState {
+    writer: Option<tokio::sync::Mutex<Box<ManagedBlobWriter>>>,
+}
+
+impl ManagedBlobWriteState {
+    pub(crate) fn new(file_io: &FileIO, config: &KeyValueWriteConfig) -> 
Result<Self> {
+        let options = CoreOptions::new(&config.table_options);
+        let writer = ManagedBlobWriter::new(
+            file_io.clone(),
+            &config.table_location,
+            &config.partition_path,
+            config.bucket,
+            &config.data_file_prefix,
+            options.blob_target_file_size(),
+            managed_blob_fields(&config.value_fields, &options),
+        )?;
+        Ok(Self {
+            writer: writer.map(|writer| 
tokio::sync::Mutex::new(Box::new(writer))),
+        })
+    }
+
+    pub(crate) fn enabled(&self) -> bool {
+        self.writer.is_some()
+    }
+
+    pub(crate) async fn externalize(&mut self, batch: RecordBatch) -> 
Result<RecordBatch> {
+        match &mut self.writer {
+            Some(writer) => writer.get_mut().externalize(&batch).await,
+            None => Ok(batch),
+        }
+    }
+
+    pub(crate) async fn prepare_commit(&mut self) -> Result<()> {
+        if let Some(writer) = &mut self.writer {
+            writer.get_mut().prepare_commit().await?;
+        }
+        Ok(())
+    }
+
+    pub(crate) async fn abort(&mut self) {
+        if let Some(writer) = &mut self.writer {
+            writer.get_mut().abort().await;
+        }
+    }
+}
+
+impl ManagedBlobWriter {
+    pub(crate) fn new(
+        file_io: FileIO,
+        table_location: &str,
+        partition_path: &str,
+        bucket: i32,
+        file_prefix: &str,
+        target_file_size: i64,
+        fields: Vec<(usize, ManagedBlobKind)>,
+    ) -> Result<Option<Self>> {
+        if fields.is_empty() {
+            return Ok(None);
+        }
+        let target_file_size = u64::try_from(target_file_size)
+            .ok()
+            .filter(|size| *size > 0)
+            .ok_or_else(|| crate::Error::DataInvalid {
+                message: "Managed BLOB target file size must be 
positive".to_string(),
+                source: None,
+            })?;
+        Ok(Some(Self {
+            file_io,
+            bucket_dir: bucket_path_under(table_location, partition_path, 
bucket),
+            file_prefix: file_prefix.to_string(),
+            target_file_size,
+            fields: fields
+                .into_iter()
+                .map(|(index, kind)| ManagedBlobField {
+                    index,
+                    kind,
+                    current: None,
+                })
+                .collect(),
+            uncommitted_paths: Vec::new(),
+        }))
+    }
+
+    pub(crate) async fn externalize(&mut self, batch: &RecordBatch) -> 
Result<RecordBatch> {
+        let kind_index = batch
+            .schema()
+            .fields()
+            .iter()
+            .position(|field| field.name() == 
crate::spec::VALUE_KIND_FIELD_NAME);
+        let kinds = kind_index
+            .map(|index| {
+                batch
+                    .column(index)
+                    .as_any()
+                    .downcast_ref::<Int8Array>()
+                    .ok_or_else(|| crate::Error::DataInvalid {
+                        message: "_VALUE_KIND column must be Int8".to_string(),
+                        source: None,
+                    })
+            })
+            .transpose()?;
+        let mut retract = Vec::with_capacity(batch.num_rows());
+        for row in 0..batch.num_rows() {
+            let kind = kinds
+                .filter(|column| column.is_valid(row))
+                .map_or(RowKind::Insert.to_value(), |column| 
column.value(row));
+            retract.push(RowKind::from_value(kind)?.is_retract());
+        }
+
+        let mut columns = batch.columns().to_vec();
+        for field_index in 0..self.fields.len() {
+            let column_index = self.fields[field_index].index;
+            let column = match self.fields[field_index].kind {
+                ManagedBlobKind::Scalar => {
+                    let values = 
downcast_blob_column(columns[column_index].as_ref())?;
+                    Arc::new(
+                        self.externalize_values(field_index, values, &retract)
+                            .await?,
+                    ) as ArrayRef
+                }
+                ManagedBlobKind::Array => {
+                    let array = columns[column_index]
+                        .as_any()
+                        .downcast_ref::<ListArray>()
+                        .ok_or_else(|| invalid_blob_column("ARRAY<BLOB> 
requires ListArray"))?;
+                    let values = 
downcast_blob_column(array.values().as_ref())?;
+                    let (offsets, indices) = visible_child_indices(
+                        array.value_offsets(),
+                        array,
+                        values.len(),
+                        &retract,
+                    )?;
+                    let values = Arc::new(
+                        self.externalize_selected_values(field_index, values, 
&indices)
+                            .await?,
+                    );
+                    let ArrowDataType::List(element) = array.data_type() else {
+                        unreachable!()
+                    };
+                    Arc::new(
+                        ListArray::try_new(
+                            element.clone(),
+                            offsets,
+                            values,
+                            combined_nulls(array, &retract),
+                        )
+                        .map_err(|error| 
invalid_blob_column(&error.to_string()))?,
+                    )
+                }
+                ManagedBlobKind::Map => {
+                    let map = columns[column_index]
+                        .as_any()
+                        .downcast_ref::<MapArray>()
+                        .ok_or_else(|| invalid_blob_column("MAP<X, BLOB> 
requires MapArray"))?;
+                    let values = 
downcast_blob_column(map.entries().column(1).as_ref())?;
+                    let (offsets, indices) =
+                        visible_child_indices(map.value_offsets(), map, 
values.len(), &retract)?;
+                    let values = Arc::new(
+                        self.externalize_selected_values(field_index, values, 
&indices)
+                            .await?,
+                    );
+                    let ArrowDataType::Map(entries_field, ordered) = 
map.data_type() else {
+                        unreachable!()
+                    };
+                    let ArrowDataType::Struct(entry_fields) = 
entries_field.data_type() else {
+                        unreachable!()
+                    };
+                    let entries = StructArray::try_new(
+                        entry_fields.clone(),
+                        vec![
+                            take(map.entries().column(0).as_ref(), &indices, 
None)
+                                .map_err(|error| 
invalid_blob_column(&error.to_string()))?,
+                            values,
+                        ],
+                        None,
+                    )
+                    .map_err(|error| invalid_blob_column(&error.to_string()))?;
+                    Arc::new(
+                        MapArray::try_new(
+                            entries_field.clone(),
+                            offsets,
+                            entries,
+                            combined_nulls(map, &retract),
+                            *ordered,
+                        )
+                        .map_err(|error| 
invalid_blob_column(&error.to_string()))?,
+                    )
+                }
+            };
+            columns[column_index] = column;
+        }
+        // Retractions discard BLOB payloads even when the caller supplied a
+        // stricter, non-nullable Arrow field than the table schema. Keep the
+        // internal descriptor batch nullable for every managed BLOB field.
+        let mut fields = batch.schema().fields().to_vec();
+        for field in &self.fields {
+            let input = &fields[field.index];
+            if !input.is_nullable() {
+                fields[field.index] = 
Arc::new(input.as_ref().clone().with_nullable(true));
+            }
+        }
+        let schema = Arc::new(ArrowSchema::new_with_metadata(
+            fields,
+            batch.schema().metadata().clone(),
+        ));
+        RecordBatch::try_new(schema, columns)
+            .map_err(|error| invalid_blob_column(&error.to_string()))
+    }
+
+    async fn externalize_values(
+        &mut self,
+        field_index: usize,
+        values: &LargeBinaryArray,
+        retract: &[bool],
+    ) -> Result<LargeBinaryArray> {
+        if values.len() != retract.len() {
+            return Err(invalid_blob_column(
+                "Managed BLOB retract mask length mismatch",
+            ));
+        }
+        let mut builder = LargeBinaryBuilder::new();
+        for (index, is_retract) in retract.iter().enumerate() {
+            if *is_retract || values.is_null(index) {
+                builder.append_null();
+                continue;
+            }
+            let descriptor = self.write_value(field_index, 
values.value(index)).await?;
+            builder.append_value(descriptor.serialize());
+        }
+        Ok(builder.finish())
+    }
+
+    async fn externalize_selected_values(
+        &mut self,
+        field_index: usize,
+        values: &LargeBinaryArray,
+        indices: &UInt32Array,
+    ) -> Result<LargeBinaryArray> {
+        let mut builder = LargeBinaryBuilder::new();
+        for &index in indices.values().iter() {
+            let index = index as usize;
+            if values.is_null(index) {
+                builder.append_null();
+            } else {
+                let descriptor = self.write_value(field_index, 
values.value(index)).await?;
+                builder.append_value(descriptor.serialize());
+            }
+        }
+        Ok(builder.finish())
+    }
+
+    async fn write_value(&mut self, field_index: usize, value: &[u8]) -> 
Result<BlobDescriptor> {
+        if self.fields[field_index].current.is_none() {
+            self.file_io
+                .mkdirs(&format!("{}/", self.bucket_dir))
+                .await?;
+            let path = format!(
+                "{}/{}{}.managed.blob",
+                self.bucket_dir,
+                self.file_prefix,
+                uuid::Uuid::new_v4()
+            );
+            let output = self.file_io.new_output(&path)?;
+            let writer =
+                Box::new(BlobFormatWriter::new(&output, 
Some(self.file_io.clone())).await?);
+            self.uncommitted_paths.push(path.clone());
+            self.fields[field_index].current = Some(ManagedBlobPack { path, 
writer });
+        }
+        let pack = self.fields[field_index].current.as_mut().unwrap();
+        let (offset, length) = pack.writer.write_managed_value(value).await?;
+        let descriptor = BlobDescriptor::new(pack.path.clone(), offset, 
length);
+        if pack.writer.num_bytes() as u64 >= self.target_file_size {
+            self.close_pack(field_index).await?;
+        }
+        Ok(descriptor)
+    }
+
+    async fn close_pack(&mut self, field_index: usize) -> Result<()> {
+        if let Some(pack) = self.fields[field_index].current.take() {
+            pack.writer.close().await?;
+        }
+        Ok(())
+    }
+
+    pub(crate) async fn prepare_commit(&mut self) -> Result<()> {
+        for index in 0..self.fields.len() {
+            self.close_pack(index).await?;
+        }
+        self.uncommitted_paths.clear();
+        Ok(())
+    }
+
+    pub(crate) async fn abort(&mut self) {
+        for field in &mut self.fields {
+            field.current.take();
+        }
+        for path in self.uncommitted_paths.drain(..) {
+            let _ = self.file_io.delete_file(&path).await;
+        }
+    }
+}
+
+fn downcast_blob_column(array: &dyn Array) -> Result<&LargeBinaryArray> {
+    array
+        .as_any()
+        .downcast_ref::<LargeBinaryArray>()
+        .ok_or_else(|| invalid_blob_column("BLOB values require 
LargeBinaryArray"))
+}
+
+fn invalid_blob_column(message: &str) -> crate::Error {
+    crate::Error::DataInvalid {
+        message: message.to_string(),
+        source: None,
+    }
+}
+
+fn visible_child_indices(
+    offsets: &[i32],
+    array: &dyn Array,
+    child_len: usize,
+    retract: &[bool],
+) -> Result<(OffsetBuffer<i32>, UInt32Array)> {
+    if offsets.len() != array.len() + 1 || retract.len() != array.len() {
+        return Err(invalid_blob_column(
+            "Managed BLOB collection offsets do not match parent rows",
+        ));
+    }
+    // Slices retain backing children outside their visible rows. A retract or
+    // null parent also has no logical children. Rebase offsets and retain only
+    // visible children so non-nullable element/value fields stay valid.
+    let mut indices = Vec::new();
+    let mut rebased = Vec::with_capacity(array.len() + 1);
+    rebased.push(0);
+    for (row, is_retract) in retract.iter().enumerate() {
+        if !*is_retract && array.is_valid(row) {
+            let start = usize::try_from(offsets[row])
+                .map_err(|_| invalid_blob_column("Negative BLOB child 
offset"))?;
+            let end = usize::try_from(offsets[row + 1])
+                .map_err(|_| invalid_blob_column("Negative BLOB child 
offset"))?;
+            if start > end || end > child_len {
+                return Err(invalid_blob_column(
+                    "BLOB collection offset exceeds child array",
+                ));
+            }
+            for index in start..end {
+                indices.push(
+                    u32::try_from(index).map_err(|_| {
+                        invalid_blob_column("BLOB collection child index 
exceeds u32")
+                    })?,
+                );
+            }
+        }
+        rebased.push(
+            i32::try_from(indices.len())
+                .map_err(|_| invalid_blob_column("BLOB collection offset 
exceeds i32"))?,
+        );
+    }
+    Ok((
+        OffsetBuffer::new(ScalarBuffer::from(rebased)),
+        UInt32Array::from(indices),
+    ))
+}
+
+fn combined_nulls(array: &dyn Array, retract: &[bool]) -> Option<NullBuffer> {
+    let valid = (0..array.len())
+        .map(|index| array.is_valid(index) && !retract[index])
+        .collect::<Vec<_>>();
+    valid
+        .iter()
+        .any(|valid| !valid)
+        .then(|| NullBuffer::from(valid))
+}
diff --git a/crates/paimon/src/table/mod.rs b/crates/paimon/src/table/mod.rs
index 4afd1520..0b343e36 100644
--- a/crates/paimon/src/table/mod.rs
+++ b/crates/paimon/src/table/mod.rs
@@ -72,6 +72,11 @@ pub(crate) mod index_file_path;
 mod kv_file_reader;
 mod kv_file_writer;
 mod lumina_index_build_builder;
+mod managed_blob_reader;
+mod managed_blob_reference;
+#[cfg(test)]
+mod managed_blob_table_tests;
+mod managed_blob_writer;
 pub(crate) mod merge_tree_split_generator;
 #[cfg(test)]
 mod mosaic_table_write_tests;
diff --git a/crates/paimon/src/table/read_builder.rs 
b/crates/paimon/src/table/read_builder.rs
index 3561616a..75183365 100644
--- a/crates/paimon/src/table/read_builder.rs
+++ b/crates/paimon/src/table/read_builder.rs
@@ -530,10 +530,12 @@ impl<'a> PaimonReadBuilder<'a> {
             PartitionFilter::from_predicate(pred, 
&self.table.schema().partition_fields())
         });
         let read_type = self.resolve_read_type().unwrap_or(None);
+        let scan_predicates =
+            super::managed_blob_reader::scan_predicates(self.table, 
&self.filter.data_predicates);
         TableScan::new(
             self.table,
             partition_filter,
-            self.filter.data_predicates.clone(),
+            scan_predicates,
             self.filter.bucket_predicate.clone(),
             self.limit,
             self.effective_row_ranges(),
diff --git a/crates/paimon/src/table/table_read.rs 
b/crates/paimon/src/table/table_read.rs
index 81a5b200..c7516df4 100644
--- a/crates/paimon/src/table/table_read.rs
+++ b/crates/paimon/src/table/table_read.rs
@@ -872,7 +872,7 @@ impl<'a> PaimonTableRead<'a> {
                     | MergeEngine::Aggregation
             )
         {
-            return self.read_pk(data_splits, &core_options);
+            return self.read_pk_with_blob(data_splits, &core_options);
         }
 
         if core_options.data_evolution_enabled() {
@@ -882,6 +882,41 @@ impl<'a> PaimonTableRead<'a> {
         }
     }
 
+    fn read_pk_with_blob(
+        &self,
+        data_splits: &[DataSplit],
+        core_options: &CoreOptions<'_>,
+    ) -> crate::Result<ArrowRecordBatchStream> {
+        use super::managed_blob_reader::{resolve_primary_key_blob_stream, 
ManagedBlobReadPlan};
+
+        if let Some(plan) = ManagedBlobReadPlan::new(
+            self.read_type(),
+            &self.data_predicates,
+            self.table.schema().fields(),
+            core_options,
+        ) {
+            let mut inner = self.clone();
+            inner.read_type = plan.scan_fields().to_vec();
+            inner.data_predicates.clear();
+            let stream = inner.read_pk(data_splits, core_options)?;
+            return Ok(plan.finish(
+                stream,
+                core_options,
+                self.table.file_io.clone(),
+                self.blob_parallelism,
+            ));
+        }
+
+        let stream = self.read_pk(data_splits, core_options)?;
+        Ok(resolve_primary_key_blob_stream(
+            stream,
+            self.read_type(),
+            core_options,
+            self.table.file_io.clone(),
+            self.blob_parallelism,
+        ))
+    }
+
     /// Read PK table. For `Deduplicate` and `FirstRow`, raw-convertible 
splits from scan
     /// planning (mirrors Java `DataSplit#convertToRawFiles`) use the faster
     /// DataFileReader; the rest go through KeyValueFileReader for sort-merge

Reply via email to