leaves12138 commented on code in PR #954: URL: https://github.com/apache/paimon-rust/pull/954#discussion_r4101432279
########## 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: [P2] Preserve validity for non-nullable BLOB collection children The outer-field nullability change happens too late to handle non-nullable collection children. `externalize_values` inserts nulls for retracted rows and children outside a sliced parent's visible range, but `ListArray::try_new` still receives the original non-nullable element field; the MAP branch similarly reuses the non-nullable value field in `StructArray::try_new`. I reproduced four failures through `TableWrite::write_arrow_batch`: ARRAY<BLOB NOT NULL> and MAP<STRING, BLOB NOT NULL>, each with (a) two valid rows of one non-null payload each followed by `batch.slice(0, 1)`, or (b) DELETE row kinds. ARRAY reports `Non-nullable field of ListArray "element" cannot contain nulls`; MAP reports `Found unmasked nulls for non-nullable StructArray field "value"`. These schemas are accepted by the schema builder. Please handle child nullability/hidden-child representation as well as outer-field nullability, and add these cases to the regression tests. -- 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]
