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