diff --git a/.licenserc.yaml b/.licenserc.yaml index e147040a5..9ee03354c 100644 --- a/.licenserc.yaml +++ b/.licenserc.yaml @@ -28,6 +28,7 @@ header: - ".github/PULL_REQUEST_TEMPLATE.md" - "crates/paimon/tests/**/*.json" - "crates/paimon/testdata/**" + - "bindings/go/tests/testdata/**" - "third-party-licenses/jieba-rs-0.10.3.LICENSE" - "third-party-licenses/openssl-1.1.1.LICENSE" - "**/go.sum" diff --git a/bindings/go/blob_reader.go b/bindings/go/blob_reader.go index 7f0071e51..529ae9b03 100644 --- a/bindings/go/blob_reader.go +++ b/bindings/go/blob_reader.go @@ -23,9 +23,12 @@ import ( "context" "fmt" "runtime" + "strings" "sync" "unsafe" + "github.com/apache/arrow-go/v18/arrow" + "github.com/apache/arrow-go/v18/arrow/array" "github.com/jupiterrider/ffi" ) @@ -83,6 +86,40 @@ func (r *BlobReader) ReadBlobs(descriptors [][]byte) ([][]byte, error) { return ffiBlobReaderReadBlobs.symbol(r.ctx)(r.inner, descriptors) } +// StringBlobMapDescriptors returns one MAP row. +func StringBlobMapDescriptors(column arrow.Array, row int) (map[string][]byte, error) { + m, ok := column.(*array.Map) + if !ok { + return nil, fmt.Errorf("paimon: BLOB map column is %T, want *array.Map", column) + } + if row < 0 || row >= m.Len() { + return nil, fmt.Errorf("paimon: BLOB map row %d is out of range", row) + } + if m.IsNull(row) { + return nil, nil + } + keys, ok := m.Keys().(*array.String) + if !ok { + return nil, fmt.Errorf("paimon: BLOB map keys are %T, want *array.String", m.Keys()) + } + descriptors, ok := m.Items().(*array.Binary) + if !ok { + return nil, fmt.Errorf("paimon: BLOB map values are %T, want *array.Binary", m.Items()) + } + start, end := m.ValueOffsets(row) + result := make(map[string][]byte, end-start) + for index := start; index < end; index++ { + i := int(index) + key := strings.Clone(keys.Value(i)) + if descriptors.IsNull(i) { + result[key] = nil + continue + } + result[key] = append([]byte(nil), descriptors.Value(i)...) + } + return result, nil +} + // Close releases the reader and is idempotent. func (r *BlobReader) Close() { r.mu.Lock() diff --git a/bindings/go/table.go b/bindings/go/table.go index e6008a228..d41a84bdd 100644 --- a/bindings/go/table.go +++ b/bindings/go/table.go @@ -21,6 +21,7 @@ package paimon import ( "context" + "runtime" "sync" "unsafe" @@ -54,8 +55,20 @@ func (t *Table) NewReadBuilder() (*ReadBuilder, error) { if t.inner == nil { return nil, ErrClosed } - createFn := ffiTableNewReadBuilder.symbol(t.ctx) - inner, err := createFn(t.inner) + inner, err := ffiTableNewReadBuilder.symbol(t.ctx)(t.inner) + if err != nil { + return nil, err + } + t.lib.acquire() + return &ReadBuilder{ctx: t.ctx, lib: t.lib, inner: inner}, nil +} + +// NewReadBuilderWithOptions creates a ReadBuilder with per-read options. +func (t *Table) NewReadBuilderWithOptions(options map[string]string) (*ReadBuilder, error) { + if t.inner == nil { + return nil, ErrClosed + } + inner, err := ffiTableNewReadBuilderWithOptions.symbol(t.ctx)(t.inner, options) if err != nil { return nil, err } @@ -93,3 +106,49 @@ var ffiTableNewReadBuilder = newFFI(ffiOpts{ return result.readBuilder, nil } }) + +var ffiTableNewReadBuilderWithOptions = newFFI(ffiOpts{ + sym: "paimon_table_new_read_builder_with_options", + rType: &typeResultReadBuilder, + aTypes: []*ffi.Type{ + &ffi.TypePointer, + &ffi.TypePointer, + &ffi.TypePointer, + }, +}, func(ctx context.Context, ffiCall ffiCall) func(*paimonTable, map[string]string) (*paimonReadBuilder, error) { + return func(table *paimonTable, options map[string]string) (*paimonReadBuilder, error) { + type paimonOption struct { + key *byte + value *byte + } + opts := make([]paimonOption, 0, len(options)) + for key, value := range options { + keyPtr, err := bytePtrFromString(key) + if err != nil { + return nil, err + } + valuePtr, err := bytePtrFromString(value) + if err != nil { + return nil, err + } + opts = append(opts, paimonOption{key: keyPtr, value: valuePtr}) + } + var optsPtr unsafe.Pointer + if len(opts) > 0 { + optsPtr = unsafe.Pointer(&opts[0]) + } + optsLen := uintptr(len(opts)) + var result resultReadBuilder + ffiCall( + unsafe.Pointer(&result), + unsafe.Pointer(&table), + unsafe.Pointer(&optsPtr), + unsafe.Pointer(&optsLen), + ) + runtime.KeepAlive(opts) + if result.error != nil { + return nil, parseError(ctx, result.error) + } + return result.readBuilder, nil + } +}) diff --git a/bindings/go/tests/blob_reader_test.go b/bindings/go/tests/blob_reader_test.go index 14e3ea39b..191d854c7 100644 --- a/bindings/go/tests/blob_reader_test.go +++ b/bindings/go/tests/blob_reader_test.go @@ -30,6 +30,7 @@ import ( "strings" "testing" + "github.com/apache/arrow-go/v18/arrow/array" paimon "github.com/apache/paimon-rust/bindings/go" ) @@ -101,6 +102,95 @@ func TestBlobReaderReadBlobAndBatch(t *testing.T) { } } +func TestStringBlobMapDescriptors(t *testing.T) { + source := filepath.Join("testdata", "map_blob_table") + warehouse := t.TempDir() + if err := copyDirectory(source, filepath.Join(warehouse, "default.db", "map_blob_table")); err != nil { + t.Fatal(err) + } + table := openTableAt(t, warehouse, "map_blob_table") + builder, err := table.NewReadBuilderWithOptions(map[string]string{ + "blob-as-descriptor": "true", + }) + if err != nil { + t.Fatal(err) + } + defer builder.Close() + if err := builder.WithProjection([]string{"id", "assets"}); err != nil { + t.Fatal(err) + } + scan, err := builder.NewScan() + if err != nil { + t.Fatal(err) + } + defer scan.Close() + plan, err := scan.Plan() + if err != nil { + t.Fatal(err) + } + defer plan.Close() + read, err := builder.NewRead() + if err != nil { + t.Fatal(err) + } + defer read.Close() + batches, err := read.NewRecordBatchReader(plan.Splits()) + if err != nil { + t.Fatal(err) + } + defer batches.Close() + + rows := make(map[int32]map[string][]byte) + for { + record, err := batches.NextRecord() + if errors.Is(err, io.EOF) { + break + } + if err != nil { + t.Fatal(err) + } + ids := record.Column(0).(*array.Int32) + for row := 0; row < int(record.NumRows()); row++ { + rows[ids.Value(row)], err = paimon.StringBlobMapDescriptors(record.Column(1), row) + if err != nil { + record.Release() + t.Fatal(err) + } + } + record.Release() + } + + descriptors := rows[1] + if len(rows) != 3 { + t.Fatalf("read %d rows, want 3", len(rows)) + } + if len(descriptors["first"]) == 0 || len(descriptors["tail"]) == 0 { + t.Fatalf("descriptor map is invalid after Arrow release: %#v", descriptors) + } + if rows[2] != nil || len(rows[3]) != 0 { + t.Fatalf("unexpected null or empty maps: %#v", rows) + } + reader, err := table.NewBlobReader() + if err != nil { + t.Fatal(err) + } + defer reader.Close() + resolved, err := reader.ReadBlobs([][]byte{ + descriptors["tail"], + descriptors["first"], + descriptors["empty"], + }) + if err != nil { + t.Fatal(err) + } + if string(resolved[0]) != "ghij" || string(resolved[1]) != "abc" || len(resolved[2]) != 0 { + t.Fatalf("unexpected values: %q", resolved) + } + if descriptors["null"] != nil { + t.Fatalf("null BLOB returned %#v", descriptors["null"]) + } +} + func TestBlobReaderFromTableOutlivesTable(t *testing.T) { file := writeBlobFile(t, "table", "abcdefghij") diff --git a/bindings/go/tests/testdata/map_blob_table/bucket-0/data-5d9ffa8c-bb8e-4beb-bec0-20417639ec99-0.parquet b/bindings/go/tests/testdata/map_blob_table/bucket-0/data-5d9ffa8c-bb8e-4beb-bec0-20417639ec99-0.parquet new file mode 100644 index 000000000..135180d74 Binary files /dev/null and b/bindings/go/tests/testdata/map_blob_table/bucket-0/data-5d9ffa8c-bb8e-4beb-bec0-20417639ec99-0.parquet differ diff --git a/bindings/go/tests/testdata/map_blob_table/bucket-0/data-feedc1e6-e063-4f9b-8fb7-e70912f4e374-0.blob b/bindings/go/tests/testdata/map_blob_table/bucket-0/data-feedc1e6-e063-4f9b-8fb7-e70912f4e374-0.blob new file mode 100644 index 000000000..93a623bad Binary files /dev/null and b/bindings/go/tests/testdata/map_blob_table/bucket-0/data-feedc1e6-e063-4f9b-8fb7-e70912f4e374-0.blob differ diff --git a/bindings/go/tests/testdata/map_blob_table/manifest/manifest-aa99cf7f-6ec4-420a-a2fb-091e7a31c5a2-0 b/bindings/go/tests/testdata/map_blob_table/manifest/manifest-aa99cf7f-6ec4-420a-a2fb-091e7a31c5a2-0 new file mode 100644 index 000000000..095355684 Binary files /dev/null and b/bindings/go/tests/testdata/map_blob_table/manifest/manifest-aa99cf7f-6ec4-420a-a2fb-091e7a31c5a2-0 differ diff --git a/bindings/go/tests/testdata/map_blob_table/manifest/manifest-list-d854a2ec-c453-4ece-bcbe-5ad5411a8273-0 b/bindings/go/tests/testdata/map_blob_table/manifest/manifest-list-d854a2ec-c453-4ece-bcbe-5ad5411a8273-0 new file mode 100644 index 000000000..78521e820 Binary files /dev/null and b/bindings/go/tests/testdata/map_blob_table/manifest/manifest-list-d854a2ec-c453-4ece-bcbe-5ad5411a8273-0 differ diff --git a/bindings/go/tests/testdata/map_blob_table/manifest/manifest-list-d854a2ec-c453-4ece-bcbe-5ad5411a8273-1 b/bindings/go/tests/testdata/map_blob_table/manifest/manifest-list-d854a2ec-c453-4ece-bcbe-5ad5411a8273-1 new file mode 100644 index 000000000..2bdd4a875 Binary files /dev/null and b/bindings/go/tests/testdata/map_blob_table/manifest/manifest-list-d854a2ec-c453-4ece-bcbe-5ad5411a8273-1 differ diff --git a/bindings/go/tests/testdata/map_blob_table/schema/schema-0 b/bindings/go/tests/testdata/map_blob_table/schema/schema-0 new file mode 100644 index 000000000..acdb0d855 --- /dev/null +++ b/bindings/go/tests/testdata/map_blob_table/schema/schema-0 @@ -0,0 +1,30 @@ +{ + "version": 3, + "id": 0, + "fields": [ + { + "id": 0, + "name": "id", + "type": "INT" + }, + { + "id": 1, + "name": "assets", + "type": { + "type": "MAP", + "key": "STRING NOT NULL", + "value": "BLOB", + "nullable": true + } + } + ], + "highestFieldId": 1, + "partitionKeys": [], + "primaryKeys": [], + "options": { + "row-tracking.enabled": "true", + "data-evolution.enabled": "true" + }, + "comment": null, + "timeMillis": 1788359448372 +} diff --git a/bindings/go/tests/testdata/map_blob_table/snapshot/LATEST b/bindings/go/tests/testdata/map_blob_table/snapshot/LATEST new file mode 100644 index 000000000..56a6051ca --- /dev/null +++ b/bindings/go/tests/testdata/map_blob_table/snapshot/LATEST @@ -0,0 +1 @@ +1 \ No newline at end of file diff --git a/bindings/go/tests/testdata/map_blob_table/snapshot/snapshot-1 b/bindings/go/tests/testdata/map_blob_table/snapshot/snapshot-1 new file mode 100644 index 000000000..0e069762d --- /dev/null +++ b/bindings/go/tests/testdata/map_blob_table/snapshot/snapshot-1 @@ -0,0 +1,15 @@ +{ + "version": 3, + "id": 1, + "schemaId": 0, + "baseManifestList": "manifest-list-d854a2ec-c453-4ece-bcbe-5ad5411a8273-0", + "deltaManifestList": "manifest-list-d854a2ec-c453-4ece-bcbe-5ad5411a8273-1", + "totalRecordCount": 6, + "deltaRecordCount": 6, + "commitUser": "ab9bd0e7-c0bc-4695-84cf-44f36df42a6a", + "commitIdentifier": 9223372036854775807, + "commitKind": "APPEND", + "timeMillis": 1788359448379, + "nextRowId": 3, + "uuid": "19a3dfb1-eae3-4d7a-8329-26b85143edcc" +} \ No newline at end of file diff --git a/crates/paimon/src/arrow/format/blob.rs b/crates/paimon/src/arrow/format/blob.rs index a7c7d7ed3..6002d4657 100644 --- a/crates/paimon/src/arrow/format/blob.rs +++ b/crates/paimon/src/arrow/format/blob.rs @@ -22,7 +22,13 @@ use crate::spec::{BlobDescriptor, DataField, DataType}; use crate::table::{ArrowRecordBatchStream, RowRange}; use crate::Error; use arrow_array::builder::{BinaryBuilder, ListBuilder}; -use arrow_array::{Array, ArrayRef, RecordBatch, RecordBatchOptions}; +use arrow_array::{ + Array, ArrayRef, BinaryArray, BooleanArray, Date32Array, Decimal128Array, Int16Array, + Int32Array, Int64Array, Int8Array, MapArray, RecordBatch, RecordBatchOptions, StringArray, + StructArray, Time32MillisecondArray, +}; +use arrow_buffer::{BooleanBuffer, NullBuffer, OffsetBuffer, ScalarBuffer}; +use arrow_schema::DataType as ArrowDataType; use async_stream::try_stream; use async_trait::async_trait; use bytes::Bytes; @@ -96,12 +102,29 @@ impl IndexedBlobReader { ) .await } + + pub(crate) async fn read_map_positions( + &self, + positions: &[usize], + key_type: &DataType, + ) -> crate::Result> { + let planned_reads = plan_blob_array_reads(&self.index, positions)?; + fetch_blob_map_values( + self.reader.as_ref(), + planned_reads, + &self.file_path, + self.descriptor_mode, + key_type, + ) + .await + } } #[derive(Debug)] pub(crate) enum BlobReadValue { Value(Bytes), Array(Vec>), + Map(Vec<(Bytes, Option)>), Null, Placeholder, } @@ -121,11 +144,18 @@ const BLOB_ARRAY_HEADER_SIZE: u64 = 9; const BLOB_ARRAY_INDEX_LENGTH_SIZE: u64 = 4; const BLOB_ARRAY_MIN_PAYLOAD_SIZE: u64 = BLOB_ARRAY_HEADER_SIZE + BLOB_ARRAY_INDEX_LENGTH_SIZE; const BLOB_ARRAY_NULL_ELEMENT_LENGTH: i64 = -1; +const BLOB_MAP_MAGIC_NUMBER: i32 = 0x4D424342; +const BLOB_MAP_VERSION: u8 = 1; +const BLOB_MAP_HEADER_SIZE: u64 = 9; +const BLOB_MAP_INDEX_LENGTHS_SIZE: u64 = 8; +const BLOB_MAP_MIN_PAYLOAD_SIZE: u64 = BLOB_MAP_HEADER_SIZE + BLOB_MAP_INDEX_LENGTHS_SIZE; +const BLOB_MAP_NULL_LENGTH: i64 = -1; -#[derive(Debug, Clone, Copy)] -enum BlobFieldKind { +#[derive(Debug, Clone)] +pub(crate) enum BlobFieldKind { Scalar, Array, + Map(DataType), } #[async_trait] @@ -159,7 +189,7 @@ impl FormatFileReader for BlobFormatReader { Ok(try_stream! { while let Some(positions) = selection.next_batch(batch_size) { - let batch = match field_kind { + let batch = match &field_kind { Some(BlobFieldKind::Scalar) => { let values = blob_reader.read_positions(&positions).await?; build_blob_batch(&target_schema, values)? @@ -168,6 +198,10 @@ impl FormatFileReader for BlobFormatReader { let values = blob_reader.read_array_positions(&positions).await?; build_blob_array_batch(&target_schema, values)? } + Some(BlobFieldKind::Map(key_type)) => { + let values = blob_reader.read_map_positions(&positions, key_type).await?; + build_blob_map_batch(&target_schema, values, key_type)? + } None => RecordBatch::try_new_with_options( target_schema.clone(), Vec::new(), @@ -203,9 +237,12 @@ fn validate_read_fields(read_fields: &[DataField]) -> crate::Result { Ok(BlobFieldKind::Array) } + DataType::Map(map) if matches!(map.value_type(), DataType::Blob(_)) => { + Ok(BlobFieldKind::Map(map.key_type().clone())) + } other => Err(Error::DataInvalid { message: format!( - ".blob format requires a Blob or Array field, got {:?} for column '{}'", + ".blob format requires a Blob, Array, or Map field, got {:?} for column '{}'", other, field.name() ), @@ -258,7 +295,7 @@ pub(crate) fn build_blob_batch( match value { BlobReadValue::Value(bytes) => builder.append_value(bytes.as_ref()), BlobReadValue::Null | BlobReadValue::Placeholder => builder.append_null(), - BlobReadValue::Array(_) => { + BlobReadValue::Array(_) | BlobReadValue::Map(_) => { return Err(Error::UnexpectedError { message: "Scalar BLOB reader produced an ARRAY value".to_string(), source: None, @@ -302,7 +339,7 @@ pub(crate) fn build_blob_array_batch( builder.append(true); } BlobReadValue::Null | BlobReadValue::Placeholder => builder.append(false), - BlobReadValue::Value(_) => { + BlobReadValue::Value(_) | BlobReadValue::Map(_) => { return Err(Error::UnexpectedError { message: "ARRAY reader produced a scalar BLOB value".to_string(), source: None, @@ -318,6 +355,275 @@ pub(crate) fn build_blob_array_batch( }) } +pub(crate) fn build_blob_map_batch( + target_schema: &Arc, + values: Vec, + key_type: &DataType, +) -> crate::Result { + let ArrowDataType::Map(entries_field, ordered) = target_schema.field(0).data_type() else { + return Err(Error::UnexpectedError { + message: "Expected MAP to map to Arrow Map".to_string(), + source: None, + }); + }; + let ArrowDataType::Struct(entry_fields) = entries_field.data_type() else { + return Err(Error::UnexpectedError { + message: "Expected MAP entries to be an Arrow Struct".to_string(), + source: None, + }); + }; + + let mut keys = Vec::new(); + let mut blobs = Vec::new(); + let mut blob_data_length = 0u64; + let mut offsets = vec![0i32]; + let mut validity = Vec::with_capacity(values.len()); + for value in values { + match value { + BlobReadValue::Map(entries) => { + validity.push(true); + let next = offsets + .last() + .copied() + .unwrap() + .checked_add( + i32::try_from(entries.len()).map_err(|e| Error::DataInvalid { + message: "MAP entry count exceeds Arrow i32 offsets" + .to_string(), + source: Some(Box::new(e)), + })?, + ) + .ok_or_else(|| Error::DataInvalid { + message: "MAP batch exceeds Arrow i32 offsets".to_string(), + source: None, + })?; + for (key, blob) in entries { + if let Some(blob) = &blob { + blob_data_length = checked_arrow_binary_data_length( + blob_data_length, + blob.len() as u64, + "MAP batch value data", + )?; + } + keys.push(key); + blobs.push(blob); + } + offsets.push(next); + } + BlobReadValue::Null | BlobReadValue::Placeholder => { + validity.push(false); + offsets.push(*offsets.last().unwrap()); + } + BlobReadValue::Value(_) | BlobReadValue::Array(_) => { + return Err(Error::UnexpectedError { + message: "MAP reader produced a non-map value".to_string(), + source: None, + }); + } + } + } + + let key_array = decode_blob_map_keys(&keys, key_type)?; + let value_array = Arc::new(BinaryArray::from_iter( + blobs.iter().map(|value| value.as_deref()), + )) as ArrayRef; + let entries = StructArray::try_new(entry_fields.clone(), vec![key_array, value_array], None) + .map_err(|e| Error::UnexpectedError { + message: format!("Failed to build MAP entries: {e}"), + source: Some(Box::new(e)), + })?; + let map = MapArray::try_new( + entries_field.clone(), + OffsetBuffer::new(ScalarBuffer::from(offsets)), + entries, + Some(NullBuffer::new(BooleanBuffer::from(validity))), + *ordered, + ) + .map_err(|e| Error::UnexpectedError { + message: format!("Failed to build MAP array: {e}"), + source: Some(Box::new(e)), + })?; + RecordBatch::try_new(target_schema.clone(), vec![Arc::new(map)]).map_err(|e| { + Error::UnexpectedError { + message: format!("Failed to build MAP RecordBatch: {e}"), + source: Some(Box::new(e)), + } + }) +} + +fn decode_blob_map_keys(keys: &[Bytes], key_type: &DataType) -> crate::Result { + macro_rules! fixed_keys { + ($array:ty, $type:ty, $size:expr) => {{ + let values = keys + .iter() + .map(|key| { + let bytes: [u8; $size] = + key.as_ref().try_into().map_err(|_| Error::DataInvalid { + message: format!( + "Invalid MAP fixed-width key length: {}", + key.len() + ), + source: None, + })?; + Ok(<$type>::from_le_bytes(bytes)) + }) + .collect::>>()?; + Ok(Arc::new(<$array>::from(values)) as ArrayRef) + }}; + } + + for key in keys { + validate_blob_map_key_length(key_type, key.len() as u64)?; + } + if blob_map_key_uses_binary_offsets(key_type) { + keys.iter().try_fold(0u64, |total, key| { + checked_arrow_binary_data_length(total, key.len() as u64, "MAP batch key data") + })?; + } + + match key_type { + DataType::TinyInt(_) => fixed_keys!(Int8Array, i8, 1), + DataType::SmallInt(_) => fixed_keys!(Int16Array, i16, 2), + DataType::Int(_) => fixed_keys!(Int32Array, i32, 4), + DataType::BigInt(_) => fixed_keys!(Int64Array, i64, 8), + DataType::Date(_) => fixed_keys!(Date32Array, i32, 4), + DataType::Time(_) => fixed_keys!(Time32MillisecondArray, i32, 4), + DataType::Boolean(_) => { + let values = keys + .iter() + .map(|key| match key.as_ref() { + [0] => Ok(false), + [1] => Ok(true), + _ => Err(Error::DataInvalid { + message: "Invalid MAP boolean key".to_string(), + source: None, + }), + }) + .collect::>>()?; + Ok(Arc::new(BooleanArray::from(values))) + } + DataType::Char(_) | DataType::VarChar(_) => { + let values = keys + .iter() + .map(|key| { + std::str::from_utf8(key).map_err(|e| Error::DataInvalid { + message: "Invalid MAP string key".to_string(), + source: Some(Box::new(e)), + }) + }) + .collect::>>()?; + Ok(Arc::new(StringArray::from(values))) + } + DataType::Binary(_) | DataType::VarBinary(_) => Ok(Arc::new( + BinaryArray::from_iter_values(keys.iter().map(|key| key.as_ref())), + )), + DataType::Decimal(decimal) => { + let values = keys + .iter() + .map(|key| decode_blob_map_decimal(key, decimal.precision())) + .collect::>>()?; + let array = Decimal128Array::from(values) + .with_precision_and_scale(decimal.precision() as u8, decimal.scale() as i8) + .map_err(|e| Error::DataInvalid { + message: format!("Invalid MAP decimal key: {e}"), + source: Some(Box::new(e)), + })?; + Ok(Arc::new(array)) + } + other => Err(Error::Unsupported { + message: format!("Unsupported key type for MAP: {other:?}"), + }), + } +} + +fn blob_map_key_uses_binary_offsets(key_type: &DataType) -> bool { + matches!( + key_type, + DataType::Char(_) | DataType::VarChar(_) | DataType::Binary(_) | DataType::VarBinary(_) + ) +} + +fn validate_blob_map_key_length(key_type: &DataType, length: u64) -> crate::Result<()> { + let fixed_length = match key_type { + DataType::TinyInt(_) | DataType::Boolean(_) => Some(1), + DataType::SmallInt(_) => Some(2), + DataType::Int(_) | DataType::Date(_) | DataType::Time(_) => Some(4), + DataType::BigInt(_) => Some(8), + DataType::Decimal(decimal) if decimal.precision() <= 18 => Some(8), + DataType::Decimal(_) => { + if !(1..=16).contains(&length) { + return Err(Error::DataInvalid { + message: "Invalid MAP decimal key".to_string(), + source: None, + }); + } + return Ok(()); + } + DataType::Char(_) | DataType::VarChar(_) | DataType::Binary(_) | DataType::VarBinary(_) => { + return Ok(()) + } + other => { + return Err(Error::Unsupported { + message: format!("Unsupported key type for MAP: {other:?}"), + }); + } + }; + if fixed_length != Some(length) { + return Err(Error::DataInvalid { + message: format!("Invalid MAP fixed-width key length: {length}"), + source: None, + }); + } + Ok(()) +} + +fn checked_arrow_binary_data_length( + current: u64, + additional: u64, + context: &str, +) -> crate::Result { + let total = current + .checked_add(additional) + .filter(|total| *total <= i32::MAX as u64) + .ok_or_else(|| Error::DataInvalid { + message: format!("{context} is too large for Arrow Binary"), + source: None, + })?; + Ok(total) +} + +fn decode_blob_map_decimal(bytes: &[u8], precision: u32) -> crate::Result { + let value = if precision <= 18 { + let bytes: [u8; 8] = bytes.try_into().map_err(|_| Error::DataInvalid { + message: format!( + "Invalid MAP fixed-width key length: {}", + bytes.len() + ), + source: None, + })?; + i64::from_le_bytes(bytes) as i128 + } else { + if bytes.is_empty() || bytes.len() > 16 { + return Err(Error::DataInvalid { + message: "Invalid MAP decimal key".to_string(), + source: None, + }); + } + let fill = if bytes[0] & 0x80 == 0 { 0 } else { 0xff }; + let mut extended = [fill; 16]; + extended[16 - bytes.len()..].copy_from_slice(bytes); + i128::from_be_bytes(extended) + }; + let digits = value.unsigned_abs().to_string().len() as u32; + if digits > precision { + return Err(Error::DataInvalid { + message: "MAP decimal key exceeds declared precision".to_string(), + source: None, + }); + } + Ok(value) +} + fn plan_blob_reads( blob_index: &BlobFileIndex, positions: &[usize], @@ -764,6 +1070,287 @@ fn build_blob_array_descriptors( Ok(BlobReadValue::Array(elements)) } +async fn fetch_blob_map_values( + reader: &dyn FileRead, + planned_reads: Vec, + file_path: &str, + descriptor_mode: bool, + key_type: &DataType, +) -> crate::Result> { + futures::stream::iter(planned_reads.into_iter().map(|planned_read| async move { + match planned_read { + PlannedBlobArrayRead::Null => Ok(BlobReadValue::Null), + PlannedBlobArrayRead::Placeholder => Ok(BlobReadValue::Placeholder), + PlannedBlobArrayRead::Read(payload_range) => { + read_blob_map_entry(reader, payload_range, file_path, descriptor_mode, key_type) + .await + } + } + })) + .buffered(BLOB_READ_CONCURRENCY) + .try_collect() + .await +} + +async fn read_blob_map_entry( + reader: &dyn FileRead, + payload_range: Range, + file_path: &str, + descriptor_mode: bool, + key_type: &DataType, +) -> crate::Result { + let payload_length = payload_range + .end + .checked_sub(payload_range.start) + .ok_or_else(|| Error::DataInvalid { + message: format!("Invalid MAP payload range: {payload_range:?}"), + source: None, + })?; + if payload_length < BLOB_MAP_MIN_PAYLOAD_SIZE { + return Err(Error::DataInvalid { + message: format!( + "MAP payload is too small: expected at least {BLOB_MAP_MIN_PAYLOAD_SIZE} bytes, got {payload_length}" + ), + source: None, + }); + } + + let header = read_blob_map_range( + reader, + payload_range.start..payload_range.start + BLOB_MAP_HEADER_SIZE, + "header", + ) + .await?; + let magic = i32::from_le_bytes(header[..4].try_into().unwrap()); + if magic != BLOB_MAP_MAGIC_NUMBER { + return Err(Error::DataInvalid { + message: format!( + "Invalid MAP payload magic number: expected {BLOB_MAP_MAGIC_NUMBER}, got {magic}" + ), + source: None, + }); + } + if header[4] != BLOB_MAP_VERSION { + return Err(Error::Unsupported { + message: format!( + "Unsupported MAP payload version: expected {BLOB_MAP_VERSION}, got {}", + header[4] + ), + }); + } + let entry_count = i32::from_le_bytes(header[5..9].try_into().unwrap()); + if entry_count < 0 { + return Err(Error::DataInvalid { + message: format!("Invalid MAP entry count: {entry_count}"), + source: None, + }); + } + let entry_count = entry_count as usize; + + let index_lengths_start = payload_range.end - BLOB_MAP_INDEX_LENGTHS_SIZE; + let index_lengths = read_blob_map_range( + reader, + index_lengths_start..payload_range.end, + "index lengths", + ) + .await?; + let key_index_length = i32::from_le_bytes(index_lengths[..4].try_into().unwrap()); + let value_index_length = i32::from_le_bytes(index_lengths[4..8].try_into().unwrap()); + let max_indexes = payload_length - BLOB_MAP_MIN_PAYLOAD_SIZE; + if key_index_length < 0 || key_index_length as u64 > max_indexes { + return Err(Error::DataInvalid { + message: format!("Invalid MAP key index length: {key_index_length}"), + source: None, + }); + } + if value_index_length < 0 || value_index_length as u64 > max_indexes { + return Err(Error::DataInvalid { + message: format!("Invalid MAP value index length: {value_index_length}"), + source: None, + }); + } + let key_index_length = key_index_length as u64; + let value_index_length = value_index_length as u64; + if key_index_length + value_index_length > max_indexes + || entry_count as u64 > key_index_length + || entry_count as u64 > value_index_length + { + return Err(Error::DataInvalid { + message: "MAP indexes do not match the payload".to_string(), + source: None, + }); + } + + let value_index_start = index_lengths_start - value_index_length; + let key_index_start = value_index_start - key_index_length; + let key_index = + read_blob_map_range(reader, key_index_start..value_index_start, "key index").await?; + let value_index = read_blob_map_range( + reader, + value_index_start..index_lengths_start, + "value index", + ) + .await?; + let key_lengths = decode_delta_varints(&key_index).map_err(|e| Error::DataInvalid { + message: format!("Invalid MAP key index: {e}"), + source: Some(Box::new(e)), + })?; + let value_lengths = decode_delta_varints(&value_index).map_err(|e| Error::DataInvalid { + message: format!("Invalid MAP value index: {e}"), + source: Some(Box::new(e)), + })?; + if key_lengths.len() != entry_count || value_lengths.len() != entry_count { + return Err(Error::DataInvalid { + message: "MAP entry count does not match index lengths".to_string(), + source: None, + }); + } + + let data_start = payload_range.start + BLOB_MAP_HEADER_SIZE; + let data_length = key_index_start - data_start; + let mut key_data_length = 0u64; + for &length in &key_lengths { + if length == BLOB_MAP_NULL_LENGTH { + return Err(Error::DataInvalid { + message: "MAP null keys cannot be represented by Arrow".to_string(), + source: None, + }); + } + let length = u64::try_from(length).map_err(|e| Error::DataInvalid { + message: format!("Invalid MAP key length: {length}"), + source: Some(Box::new(e)), + })?; + validate_blob_map_key_length(key_type, length)?; + key_data_length = key_data_length + .checked_add(length) + .filter(|total| *total <= data_length) + .ok_or_else(|| Error::DataInvalid { + message: "MAP key lengths exceed the payload data length".to_string(), + source: None, + })?; + } + let value_data_length = data_length - key_data_length; + let mut total_value_length = 0u64; + for &length in &value_lengths { + if length == BLOB_MAP_NULL_LENGTH { + continue; + } + let length = u64::try_from(length).map_err(|e| Error::DataInvalid { + message: format!("Invalid MAP value length: {length}"), + source: Some(Box::new(e)), + })?; + total_value_length = total_value_length + .checked_add(length) + .filter(|total| *total <= value_data_length) + .ok_or_else(|| Error::DataInvalid { + message: "MAP value lengths exceed the payload data length".to_string(), + source: None, + })?; + } + if total_value_length != value_data_length { + return Err(Error::DataInvalid { + message: "MAP key/value lengths do not match the payload data length" + .to_string(), + source: None, + }); + } + if !descriptor_mode { + checked_arrow_binary_data_length(0, total_value_length, "MAP inline value data")?; + } + if blob_map_key_uses_binary_offsets(key_type) { + checked_arrow_binary_data_length(0, key_data_length, "MAP key data")?; + } + + let key_data = + read_blob_map_range(reader, data_start..data_start + key_data_length, "key data").await?; + let mut keys = Vec::with_capacity(entry_count); + let mut cursor = 0usize; + let mut unique = std::collections::HashSet::with_capacity(entry_count); + for length in key_lengths { + let length = length as usize; + let end = cursor + length; + let key = key_data.slice(cursor..end); + if !unique.insert(key.clone()) { + return Err(Error::DataInvalid { + message: "Invalid MAP payload: duplicate key".to_string(), + source: None, + }); + } + keys.push(key); + cursor = end; + } + + let mut value_offset = data_start + key_data_length; + let mut reads = Vec::with_capacity(entry_count); + for length in value_lengths { + if length == BLOB_MAP_NULL_LENGTH { + reads.push(None); + } else { + let length = length as u64; + reads.push(Some(value_offset..value_offset + length)); + value_offset += length; + } + } + let values = if descriptor_mode { + reads + .into_iter() + .map(|range| { + range + .map(|range| { + let offset = + i64::try_from(range.start).map_err(|e| Error::DataInvalid { + message: "MAP descriptor offset exceeds i64".to_string(), + source: Some(Box::new(e)), + })?; + let length = i64::try_from(range.end - range.start).map_err(|e| { + Error::DataInvalid { + message: "MAP descriptor length exceeds i64".to_string(), + source: Some(Box::new(e)), + } + })?; + Ok(Bytes::from( + BlobDescriptor::new(file_path.to_string(), offset, length).serialize(), + )) + }) + .transpose() + }) + .collect::>>()? + } else { + let payload = read_blob_entry(reader, blob_entry_range(&payload_range)).await?; + reads + .into_iter() + .map(|range| { + range.map(|range| { + payload.slice( + (range.start - payload_range.start) as usize + ..(range.end - payload_range.start) as usize, + ) + }) + }) + .collect() + }; + Ok(BlobReadValue::Map(keys.into_iter().zip(values).collect())) +} + +async fn read_blob_map_range( + reader: &dyn FileRead, + range: Range, + part: &str, +) -> crate::Result { + let expected = range.end - range.start; + let bytes = reader.read(range.clone()).await?; + if bytes.len() as u64 != expected { + return Err(Error::DataInvalid { + message: format!( + "Short read for MAP {part} range {range:?}: expected {expected} bytes, got {}", + bytes.len() + ), + source: None, + }); + } + Ok(bytes) +} + #[derive(Debug, Clone)] enum PlannedBlobRead { Null, @@ -1368,7 +1955,7 @@ fn encode_varint(value: i64, out: &mut Vec) { mod tests { use super::*; use crate::btree::test_util::BytesFileRead; - use crate::spec::{ArrayType, BlobType}; + use crate::spec::{ArrayType, BlobType, MapType, VarCharType}; use arrow_array::Array; use bytes::Bytes; use futures::TryStreamExt; @@ -1504,6 +2091,190 @@ mod tests { ); } + #[tokio::test] + async fn test_blob_map_reader_returns_inline_values_and_descriptors() { + let payload = build_blob_map_payload(&[ + ("video", Some(b"alpha")), + ("thumbnail", None), + ("empty", Some(b"")), + ]); + let file_bytes = blob_test_utils::build_blob_file_bytes(&[Some(payload.as_slice()), None]); + let fields = blob_map_read_fields(); + + let inline = BlobFormatReader::new("file:///tmp/map.blob".to_string(), false) + .read_batch_stream( + Box::new(BytesFileRead(Bytes::from(file_bytes.clone()))), + file_bytes.len() as u64, + &fields, + None, + None, + None, + ) + .await + .unwrap() + .try_collect::>() + .await + .unwrap(); + assert_eq!( + collect_blob_map_values(&inline[0]), + vec![ + Some(vec![ + ("video".to_string(), Some(b"alpha".to_vec())), + ("thumbnail".to_string(), None), + ("empty".to_string(), Some(Vec::new())), + ]), + None, + ] + ); + + let descriptors = BlobFormatReader::new("file:///tmp/map.blob".to_string(), true) + .read_batch_stream( + Box::new(BytesFileRead(Bytes::from(file_bytes.clone()))), + file_bytes.len() as u64, + &fields, + None, + None, + None, + ) + .await + .unwrap() + .try_collect::>() + .await + .unwrap(); + let rows = collect_blob_map_values(&descriptors[0]); + let entries = rows[0].as_ref().unwrap(); + let video = BlobDescriptor::deserialize(entries[0].1.as_ref().unwrap()).unwrap(); + assert_eq!(video.uri(), "file:///tmp/map.blob"); + assert_eq!(video.length(), 5); + assert!(entries[1].1.is_none()); + let empty = BlobDescriptor::deserialize(entries[2].1.as_ref().unwrap()).unwrap(); + assert_eq!(empty.length(), 0); + } + + #[tokio::test] + async fn test_blob_map_descriptor_read_skips_values() { + let file_path = "file:///tmp/map.blob"; + let payload = + build_blob_map_payload(&[("first", Some(b"alpha")), ("second", Some(b"beta"))]); + let file_bytes = blob_test_utils::build_blob_file_bytes(&[Some(payload.as_slice())]); + let reader = TrackingFileRead::new(Bytes::from(file_bytes.clone())); + let batches = BlobFormatReader::new(file_path.to_string(), true) + .read_batch_stream( + Box::new(reader.clone()), + file_bytes.len() as u64, + &blob_map_read_fields(), + None, + None, + None, + ) + .await + .unwrap() + .try_collect::>() + .await + .unwrap(); + + for (_, descriptor) in collect_blob_map_values(&batches[0])[0].as_ref().unwrap() { + let descriptor = BlobDescriptor::deserialize(descriptor.as_ref().unwrap()).unwrap(); + let value_range = + descriptor.offset() as u64..(descriptor.offset() + descriptor.length()) as u64; + assert!(reader + .ranges() + .iter() + .all(|range| range.end <= value_range.start || range.start >= value_range.end)); + } + } + + #[tokio::test] + async fn test_inline_blob_map_reader_rejects_crc_mismatch() { + let payload = build_blob_map_payload(&[("key", Some(b"value"))]); + let mut file_bytes = blob_test_utils::build_blob_file_bytes(&[Some(payload.as_slice())]); + let value_offset = + (BLOB_INLINE_HEADER_SIZE + BLOB_MAP_HEADER_SIZE + "key".len() as u64) as usize; + file_bytes[value_offset] ^= 0xff; + + let stream = BlobFormatReader::new(String::new(), false) + .read_batch_stream( + Box::new(BytesFileRead(Bytes::from(file_bytes.clone()))), + file_bytes.len() as u64, + &blob_map_read_fields(), + None, + None, + None, + ) + .await + .unwrap(); + let error = stream.try_collect::>().await.unwrap_err(); + assert_data_invalid(error, "CRC32 mismatch"); + } + + #[tokio::test] + async fn test_inline_blob_map_reader_rejects_oversized_data_before_entry_read() { + let value_length = i32::MAX as u64 + 1; + let (reader, payload_range) = sparse_blob_map_entry(&[1], &[value_length as i64]); + let key_type = DataType::VarChar(VarCharType::new(VarCharType::MAX_LENGTH).unwrap()); + + let error = read_blob_map_entry(&reader, payload_range.clone(), "", false, &key_type) + .await + .unwrap_err(); + + assert!( + !reader.ranges().contains(&blob_entry_range(&payload_range)), + "oversized inline MAP must be rejected before reading the complete entry" + ); + assert_data_invalid(error, "too large"); + } + + #[tokio::test] + async fn test_blob_map_reader_rejects_oversized_key_before_data_read() { + let key_length = i32::MAX as u64 + 1; + let (reader, payload_range) = sparse_blob_map_entry(&[key_length as i64], &[0]); + let key_type = DataType::VarChar(VarCharType::new(VarCharType::MAX_LENGTH).unwrap()); + + let error = read_blob_map_entry(&reader, payload_range.clone(), "", false, &key_type) + .await + .unwrap_err(); + + assert_eq!(reader.ranges().len(), 4); + assert!(!reader.ranges().contains(&blob_entry_range(&payload_range))); + assert_data_invalid(error, "too large"); + } + + #[tokio::test] + async fn test_blob_map_reader_rejects_invalid_fixed_key_before_data_read() { + let key_length = i32::MAX as u64 + 1; + let (reader, payload_range) = sparse_blob_map_entry(&[key_length as i64], &[0]); + let key_type = DataType::Int(crate::spec::IntType::new()); + + let error = read_blob_map_entry(&reader, payload_range.clone(), "", false, &key_type) + .await + .unwrap_err(); + + assert_eq!(reader.ranges().len(), 4); + assert!(!reader.ranges().contains(&blob_entry_range(&payload_range))); + assert_data_invalid(error, "fixed-width key length"); + } + + #[tokio::test] + async fn test_blob_map_reader_rejects_null_key_for_arrow() { + let (reader, payload_range) = sparse_blob_map_entry(&[-1], &[0]); + let key_type = DataType::VarChar(VarCharType::new(VarCharType::MAX_LENGTH).unwrap()); + + let error = read_blob_map_entry(&reader, payload_range, "", false, &key_type) + .await + .unwrap_err(); + + assert_data_invalid(error, "null keys cannot be represented by Arrow"); + } + + #[test] + fn test_blob_map_batch_rejects_oversized_binary_data() { + let error = + checked_arrow_binary_data_length(i32::MAX as u64, 1, "MAP batch value data") + .unwrap_err(); + + assert_data_invalid(error, "too large"); + } + #[tokio::test] async fn test_inline_blob_array_reader_rejects_payload_crc_mismatch() { let payload = build_blob_array_payload(b"helloworld", &[5, -1, 5]); @@ -1534,8 +2305,20 @@ mod tests { let payload_length = BLOB_ARRAY_MIN_PAYLOAD_SIZE + element_data_length + element_index.len() as u64; let payload_range = BLOB_INLINE_HEADER_SIZE..BLOB_INLINE_HEADER_SIZE + payload_length; - let reader = - BlobArrayPreflightFileRead::new(payload_range.clone(), 1, element_index.len() as i32); + let mut header = Vec::with_capacity(BLOB_ARRAY_HEADER_SIZE as usize); + header.extend_from_slice(&BLOB_ARRAY_MAGIC_NUMBER.to_le_bytes()); + header.push(BLOB_ARRAY_VERSION); + header.extend_from_slice(&1i32.to_le_bytes()); + let reader = SparseFileRead::new(vec![ + ( + payload_range.start..payload_range.start + BLOB_ARRAY_HEADER_SIZE, + Bytes::from(header), + ), + ( + payload_range.end - BLOB_ARRAY_INDEX_LENGTH_SIZE..payload_range.end, + Bytes::copy_from_slice(&(element_index.len() as i32).to_le_bytes()), + ), + ]); let error = read_inline_blob_array_entry(&reader, payload_range.clone()) .await @@ -1848,7 +2631,7 @@ mod tests { .await; assert!( - matches!(result, Err(Error::DataInvalid { message, .. }) if message.contains("Blob or Array field")) + matches!(result, Err(Error::DataInvalid { message, .. }) if message.contains("Blob, Array, or Map field")) ); } @@ -1875,7 +2658,7 @@ mod tests { .await; assert!( - matches!(result, Err(Error::DataInvalid { message, .. }) if message.contains("Blob or Array")) + matches!(result, Err(Error::DataInvalid { message, .. }) if message.contains("Blob, Array, or Map")) ); } @@ -2150,6 +2933,17 @@ mod tests { )] } + fn blob_map_read_fields() -> Vec { + vec![DataField::new( + 0, + "payloads".to_string(), + DataType::Map(MapType::new( + DataType::VarChar(VarCharType::new(VarCharType::MAX_LENGTH).unwrap()), + DataType::Blob(BlobType::new()), + )), + )] + } + fn build_blob_array_payload(element_data: &[u8], element_lengths: &[i64]) -> Vec { let index = encode_delta_varints_write(element_lengths); let mut payload = Vec::with_capacity( @@ -2164,6 +2958,36 @@ mod tests { payload } + fn build_blob_map_payload(entries: &[(&str, Option<&[u8]>)]) -> Vec { + let key_lengths = entries + .iter() + .map(|(key, _)| key.len() as i64) + .collect::>(); + let value_lengths = entries + .iter() + .map(|(_, value)| value.map_or(-1, |value| value.len() as i64)) + .collect::>(); + let key_index = encode_delta_varints_write(&key_lengths); + let value_index = encode_delta_varints_write(&value_lengths); + let mut payload = Vec::new(); + payload.extend_from_slice(&BLOB_MAP_MAGIC_NUMBER.to_le_bytes()); + payload.push(BLOB_MAP_VERSION); + payload.extend_from_slice(&(entries.len() as i32).to_le_bytes()); + for (key, _) in entries { + payload.extend_from_slice(key.as_bytes()); + } + for (_, value) in entries { + if let Some(value) = value { + payload.extend_from_slice(value); + } + } + payload.extend_from_slice(&key_index); + payload.extend_from_slice(&value_index); + payload.extend_from_slice(&(key_index.len() as i32).to_le_bytes()); + payload.extend_from_slice(&(value_index.len() as i32).to_le_bytes()); + payload + } + fn set_blob_array_index_length(payload: &mut [u8], index_length: i32) { let index_length_position = payload.len() - BLOB_ARRAY_INDEX_LENGTH_SIZE as usize; payload[index_length_position..].copy_from_slice(&index_length.to_le_bytes()); @@ -2259,6 +3083,38 @@ mod tests { .collect() } + type BlobMapRows = Vec>)>>>; + + fn collect_blob_map_values(batch: &RecordBatch) -> BlobMapRows { + let array = batch.column(0).as_any().downcast_ref::().unwrap(); + let keys = array.keys().as_any().downcast_ref::().unwrap(); + let values = array + .values() + .as_any() + .downcast_ref::() + .unwrap(); + (0..array.len()) + .map(|row| { + if array.is_null(row) { + return None; + } + let start = array.value_offsets()[row]; + let end = array.value_offsets()[row + 1]; + Some( + (start..end) + .map(|index| { + let index = index as usize; + ( + keys.value(index).to_string(), + (!values.is_null(index)).then(|| values.value(index).to_vec()), + ) + }) + .collect(), + ) + }) + .collect() + } + fn load_blob_fixture(name: &str) -> Vec { let path = format!("{}/testdata/blob/{name}", env!("CARGO_MANIFEST_DIR")); std::fs::read(&path).unwrap_or_else(|e| panic!("Failed to read {path}: {e}")) @@ -2303,19 +3159,15 @@ mod tests { } } - struct BlobArrayPreflightFileRead { - payload_range: Range, - element_count: i32, - index_length: i32, + struct SparseFileRead { + responses: Vec<(Range, Bytes)>, ranges: Mutex>>, } - impl BlobArrayPreflightFileRead { - fn new(payload_range: Range, element_count: i32, index_length: i32) -> Self { + impl SparseFileRead { + fn new(responses: Vec<(Range, Bytes)>) -> Self { Self { - payload_range, - element_count, - index_length, + responses, ranges: Mutex::new(Vec::new()), } } @@ -2325,25 +3177,55 @@ mod tests { } } + fn sparse_blob_map_entry( + key_lengths: &[i64], + value_lengths: &[i64], + ) -> (SparseFileRead, Range) { + let key_index = encode_delta_varints_write(key_lengths); + let value_index = encode_delta_varints_write(value_lengths); + let data_length = key_lengths + .iter() + .chain(value_lengths) + .filter(|length| **length >= 0) + .map(|length| *length as u64) + .sum::(); + let payload_length = BLOB_MAP_MIN_PAYLOAD_SIZE + + data_length + + key_index.len() as u64 + + value_index.len() as u64; + let payload_range = BLOB_INLINE_HEADER_SIZE..BLOB_INLINE_HEADER_SIZE + payload_length; + let mut header = Vec::with_capacity(BLOB_MAP_HEADER_SIZE as usize); + header.extend_from_slice(&BLOB_MAP_MAGIC_NUMBER.to_le_bytes()); + header.push(BLOB_MAP_VERSION); + header.extend_from_slice(&(key_lengths.len() as i32).to_le_bytes()); + let lengths_start = payload_range.end - BLOB_MAP_INDEX_LENGTHS_SIZE; + let value_index_start = lengths_start - value_index.len() as u64; + let key_index_start = value_index_start - key_index.len() as u64; + let mut index_lengths = Vec::with_capacity(BLOB_MAP_INDEX_LENGTHS_SIZE as usize); + index_lengths.extend_from_slice(&(key_index.len() as i32).to_le_bytes()); + index_lengths.extend_from_slice(&(value_index.len() as i32).to_le_bytes()); + let reader = SparseFileRead::new(vec![ + ( + payload_range.start..payload_range.start + BLOB_MAP_HEADER_SIZE, + Bytes::from(header), + ), + (lengths_start..payload_range.end, Bytes::from(index_lengths)), + (key_index_start..value_index_start, Bytes::from(key_index)), + (value_index_start..lengths_start, Bytes::from(value_index)), + ]); + (reader, payload_range) + } + #[async_trait::async_trait] - impl FileRead for BlobArrayPreflightFileRead { + impl FileRead for SparseFileRead { async fn read(&self, range: Range) -> crate::Result { self.ranges.lock().unwrap().push(range.clone()); - - let header_range = - self.payload_range.start..self.payload_range.start + BLOB_ARRAY_HEADER_SIZE; - if range == header_range { - let mut header = Vec::with_capacity(BLOB_ARRAY_HEADER_SIZE as usize); - header.extend_from_slice(&BLOB_ARRAY_MAGIC_NUMBER.to_le_bytes()); - header.push(BLOB_ARRAY_VERSION); - header.extend_from_slice(&self.element_count.to_le_bytes()); - return Ok(Bytes::from(header)); - } - - let index_length_range = - self.payload_range.end - BLOB_ARRAY_INDEX_LENGTH_SIZE..self.payload_range.end; - if range == index_length_range { - return Ok(Bytes::copy_from_slice(&self.index_length.to_le_bytes())); + if let Some((_, bytes)) = self + .responses + .iter() + .find(|(expected, _)| expected == &range) + { + return Ok(bytes.clone()); } Err(Error::UnexpectedError { diff --git a/crates/paimon/src/spec/schema.rs b/crates/paimon/src/spec/schema.rs index ece7992ca..1e020effb 100644 --- a/crates/paimon/src/spec/schema.rs +++ b/crates/paimon/src/spec/schema.rs @@ -938,7 +938,7 @@ fn append_csv_field(existing: Option<&str>, field_name: &str) -> String { fn normalize_blob_field_type(field_name: &str, data_type: DataType) -> crate::Result { let nullable = data_type.is_nullable(); match data_type { - DataType::Blob(_) => Ok(data_type), + ref value if value.is_blob_file_field() => Ok(data_type), DataType::Binary(_) | DataType::VarBinary(_) => { Ok(DataType::Blob(BlobType::with_nullable(nullable))) } diff --git a/crates/paimon/src/spec/types.rs b/crates/paimon/src/spec/types.rs index 46d664dbd..4835ce1df 100644 --- a/crates/paimon/src/spec/types.rs +++ b/crates/paimon/src/spec/types.rs @@ -133,6 +133,7 @@ impl DataType { match self { DataType::Blob(_) => true, DataType::Array(array) => array.element_type().is_blob_type(), + DataType::Map(map) => map.value_type().is_blob_type(), _ => false, } } @@ -1976,10 +1977,12 @@ mod tests { fn test_blob_file_field_classification() { let blob = DataType::Blob(BlobType::new()); let array_blob = DataType::Array(ArrayType::new(blob.clone())); + let map_blob = DataType::Map(MapType::new(DataType::Int(IntType::new()), blob.clone())); let nested_array_blob = DataType::Array(ArrayType::new(array_blob.clone())); assert!(blob.is_blob_file_field()); assert!(array_blob.is_blob_file_field()); + assert!(map_blob.is_blob_file_field()); assert!(!nested_array_blob.is_blob_file_field()); assert!(!DataType::Int(IntType::new()).is_blob_file_field()); } diff --git a/crates/paimon/src/table/data_evolution_reader/blob_fallback.rs b/crates/paimon/src/table/data_evolution_reader/blob_fallback.rs index bc538df16..dbd6a0ae2 100644 --- a/crates/paimon/src/table/data_evolution_reader/blob_fallback.rs +++ b/crates/paimon/src/table/data_evolution_reader/blob_fallback.rs @@ -21,7 +21,8 @@ use super::{ }; use crate::arrow::build_target_arrow_schema; use crate::arrow::format::blob::{ - build_blob_array_batch, build_blob_batch, BlobReadValue, IndexedBlobReader, + build_blob_array_batch, build_blob_batch, build_blob_map_batch, BlobFieldKind, BlobReadValue, + IndexedBlobReader, }; use crate::io::FileIO; use crate::spec::{DataField, DataType}; @@ -50,7 +51,7 @@ impl LazyBlobFile { positions: &[usize], file_io: &FileIO, blob_as_descriptor: bool, - array_field: bool, + field_kind: &BlobFieldKind, ) -> crate::Result> { if self.reader.is_none() { let file_size = u64::try_from(self.file_size).map_err(|e| Error::DataInvalid { @@ -94,10 +95,10 @@ impl LazyBlobFile { .reader .as_ref() .expect("blob reader is initialized above"); - if array_field { - reader.read_array_positions(positions).await - } else { - reader.read_positions(positions).await + match field_kind { + BlobFieldKind::Scalar => reader.read_positions(positions).await, + BlobFieldKind::Array => reader.read_array_positions(positions).await, + BlobFieldKind::Map(key_type) => reader.read_map_positions(positions, key_type).await, } } @@ -119,13 +120,20 @@ pub(super) fn read( ) -> crate::Result { if read_fields.len() != 1 || !read_fields[0].data_type().is_blob_file_field() { return Err(Error::DataInvalid { - message: "Blob bunch should provide exactly one BLOB or ARRAY field".to_string(), + message: + "Blob bunch should provide exactly one BLOB, ARRAY, or MAP field" + .to_string(), source: None, }); } let target_schema = build_target_arrow_schema(&read_fields)?; - let array_field = matches!(read_fields[0].data_type(), DataType::Array(_)); + let field_kind = match read_fields[0].data_type() { + DataType::Blob(_) => BlobFieldKind::Scalar, + DataType::Array(_) => BlobFieldKind::Array, + DataType::Map(map) => BlobFieldKind::Map(map.key_type().clone()), + _ => unreachable!("validated as a blob file field"), + }; let batch_size = batch_size.unwrap_or(BATCH_SIZE).max(1); let split = split.clone(); @@ -172,7 +180,7 @@ pub(super) fn read( target_schema.clone(), &file_io, blob_as_descriptor, - array_field, + &field_kind, ).await?; } } @@ -185,7 +193,7 @@ async fn resolve_batch( target_schema: Arc, file_io: &FileIO, blob_as_descriptor: bool, - array_field: bool, + field_kind: &BlobFieldKind, ) -> crate::Result { let mut resolved = (0..row_ids.len()) .map(|_| BlobReadValue::Placeholder) @@ -233,7 +241,7 @@ async fn resolve_batch( if !file_positions.is_empty() { let values = file - .read_positions(&file_positions, file_io, blob_as_descriptor, array_field) + .read_positions(&file_positions, file_io, blob_as_descriptor, field_kind) .await?; for (output_position, value) in output_positions.into_iter().zip(values) { if !matches!(&value, BlobReadValue::Placeholder) { @@ -262,10 +270,10 @@ async fn resolve_batch( } } - if array_field { - build_blob_array_batch(&target_schema, resolved) - } else { - build_blob_batch(&target_schema, resolved) + match field_kind { + BlobFieldKind::Scalar => build_blob_batch(&target_schema, resolved), + BlobFieldKind::Array => build_blob_array_batch(&target_schema, resolved), + BlobFieldKind::Map(key_type) => build_blob_map_batch(&target_schema, resolved, key_type), } } @@ -424,9 +432,16 @@ mod tests { VecDeque::from([oldest]), ]; let file_io = crate::io::FileIOBuilder::new("file").build().unwrap(); - let batch = resolve_batch(&mut groups, &[0, 1, 2, 3], schema, &file_io, false, false) - .await - .unwrap(); + let batch = resolve_batch( + &mut groups, + &[0, 1, 2, 3], + schema, + &file_io, + false, + &BlobFieldKind::Scalar, + ) + .await + .unwrap(); let values = batch .column(0) .as_any() diff --git a/docs/src/go-binding.md b/docs/src/go-binding.md index 5a82c4c66..40acda2d2 100644 --- a/docs/src/go-binding.md +++ b/docs/src/go-binding.md @@ -172,6 +172,83 @@ before expiry. Set `dlf.oss-endpoint` when the server-provided endpoint is not reachable from the application. Static options passed to `paimon.NewBlobReader` are not refreshed. +Read a `MAP` column as descriptors, then stream one value: + +Rows with null map keys cannot be represented as Arrow maps and return an error. + +```go +readBuilder, err := table.NewReadBuilderWithOptions(map[string]string{ + "blob-as-descriptor": "true", +}) +if err != nil { + log.Fatal(err) +} +defer readBuilder.Close() +if err := readBuilder.WithProjection([]string{"assets"}); err != nil { + log.Fatal(err) +} + +scan, err := readBuilder.NewScan() +if err != nil { + log.Fatal(err) +} +defer scan.Close() +plan, err := scan.Plan() +if err != nil { + log.Fatal(err) +} +defer plan.Close() +read, err := readBuilder.NewRead() +if err != nil { + log.Fatal(err) +} +defer read.Close() +batches, err := read.NewRecordBatchReader(plan.Splits()) +if err != nil { + log.Fatal(err) +} +defer batches.Close() + +record, err := batches.NextRecord() +if err != nil { + log.Fatal(err) +} +descriptors, err := paimon.StringBlobMapDescriptors(record.Column(0), 0) +if err != nil { + log.Fatal(err) +} +record.Release() // the map owns its keys and descriptors +for key, descriptor := range descriptors { + if descriptor == nil { // null BLOB + continue + } + stream, err := reader.OpenBlob(descriptor) + if err != nil { + log.Fatal(err) + } + if _, err := io.Copy(destinationFor(key), stream); err != nil { + stream.Close() + log.Fatal(err) + } + stream.Close() +} +``` + +`StringBlobMapDescriptors` returns an ordinary Go map and remains valid after +releasing the Arrow record. To materialize small values in one merged batch: + +```go +batch := make([][]byte, 0, len(descriptors)) +for _, descriptor := range descriptors { + if descriptor != nil { + batch = append(batch, descriptor) + } +} +values, err := reader.ReadBlobs(batch) +``` + +Use `OpenBlob` for large values. + Reads are grouped by URI and nearby ranges are merged. The fixed limits are a 64 KiB merge gap, 8 MiB merged span, 8 concurrent requests, and a 64 MiB per-reader admission budget. One larger range runs alone but may exceed that