leaves12138 commented on code in PR #954:
URL: https://github.com/apache/paimon-rust/pull/954#discussion_r4101518524


##########
crates/paimon/src/table/managed_blob_writer.rs:
##########
@@ -0,0 +1,372 @@
+// 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 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,
+};
+use arrow_buffer::NullBuffer;
+use arrow_schema::{DataType as ArrowDataType, Schema as ArrowSchema};
+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>,
+}
+
+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 child_retract =
+                        child_retract_mask(array.value_offsets(), array, 
values.len(), &retract)?;
+                    let values = Arc::new(
+                        self.externalize_values(field_index, values, 
&child_retract)
+                            .await?,
+                    );
+                    let ArrowDataType::List(element) = array.data_type() else {
+                        unreachable!()
+                    };
+                    Arc::new(
+                        ListArray::try_new(
+                            element.clone(),
+                            array.offsets().clone(),
+                            values,

Review Comment:
   Rechecked the follow-up refactor at 07aab2b966e081df08f66bb78fca61b449fed75d 
as well. All four original reproductions now pass, but all four non-nullable 
collection-child reproductions above still fail (ARRAY/MAP, each with a sliced 
insert or DELETE). Full paimon library test run including the diagnostic cases: 
3261 passed, 4 failed, 6 ignored. The remaining failures are exclusively this 
child-nullability issue; the changes-requested conclusion still applies to this 
revision.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to