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 0f6a9a4c feat(go): add BlobDescriptor batch and stream readers (#761)
0f6a9a4c is described below
commit 0f6a9a4cc2a95fda095a0290338302c59776cdef
Author: XiaoHongbo <[email protected]>
AuthorDate: Mon Aug 31 19:25:32 2026 +0800
feat(go): add BlobDescriptor batch and stream readers (#761)
---
Cargo.lock | 2 +
bindings/c/Cargo.toml | 10 +-
bindings/c/DEPENDENCIES.rust.tsv | 36 ++++
bindings/c/src/blob_reader.rs | 317 ++++++++++++++++++++++++++++++++++
bindings/c/src/lib.rs | 1 +
bindings/c/src/result.rs | 30 ++++
bindings/c/src/tests.rs | 220 ++++++++++++++++++++++-
bindings/c/src/types.rs | 59 +++++++
bindings/go/DEPENDENCIES.rust.tsv | 36 ++++
bindings/go/blob_reader.go | 227 ++++++++++++++++++++++++
bindings/go/blob_stream.go | 204 ++++++++++++++++++++++
bindings/go/tests/blob_reader_test.go | 261 ++++++++++++++++++++++++++++
bindings/go/types.go | 92 ++++++++++
docs/src/go-binding.md | 138 +++++++++++++++
14 files changed, 1631 insertions(+), 2 deletions(-)
diff --git a/Cargo.lock b/Cargo.lock
index 29e93ca5..6af8279c 100644
--- a/Cargo.lock
+++ b/Cargo.lock
@@ -4732,7 +4732,9 @@ dependencies = [
"paimon",
"paimon-vindex-core",
"serde_json",
+ "tempfile",
"tokio",
+ "url",
]
[[package]]
diff --git a/bindings/c/Cargo.toml b/bindings/c/Cargo.toml
index 8320f39b..949044a0 100644
--- a/bindings/c/Cargo.toml
+++ b/bindings/c/Cargo.toml
@@ -31,7 +31,13 @@ crate-type = ["cdylib", "staticlib", "rlib"]
doc = false
[dependencies]
-paimon = { path = "../../crates/paimon" }
+paimon = { path = "../../crates/paimon", features = [
+ "storage-azdls",
+ "storage-cos",
+ "storage-gcs",
+ "storage-obs",
+ "storage-s3",
+] }
tokio = { workspace = true, features = ["rt-multi-thread"] }
futures = "0.3"
arrow = { workspace = true }
@@ -45,3 +51,5 @@ bytes = "1.7.1"
# Test-only: the vector-search integration tests build a real primary-key
vindex
# IVF-flat ANN segment fixture in-process. Versions match crates/paimon.
paimon-vindex-core = "0.4.0"
+tempfile = "3"
+url = "2.5.2"
diff --git a/bindings/c/DEPENDENCIES.rust.tsv b/bindings/c/DEPENDENCIES.rust.tsv
index 9a40a09f..f50b9243 100644
--- a/bindings/c/DEPENDENCIES.rust.tsv
+++ b/bindings/c/DEPENDENCIES.rust.tsv
@@ -1,5 +1,6 @@
crate 0BSD Apache-2.0 Apache-2.0 WITH LLVM-exception BSD-2-Clause
BSD-3-Clause BSL-1.0 CC0-1.0 CDLA-Permissive-2.0 ISC
LGPL-2.1-or-later MIT MIT-0 Unicode-3.0 Unlicense Zlib
[email protected] X X
X
[email protected] X
X
[email protected] X
X
[email protected]
X X
[email protected] X
@@ -35,10 +36,12 @@ [email protected] X X
X X X
[email protected] X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
@@ -47,15 +50,18 @@ [email protected] X
X
[email protected] X
X X
[email protected]
X X
[email protected]
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected]
X
[email protected]
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
@@ -63,6 +69,7 @@ [email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
@@ -77,6 +84,7 @@ [email protected] X
X
[email protected]
X
[email protected]
X
[email protected]
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
@@ -143,6 +151,7 @@ [email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected]
X
[email protected] X
X
[email protected] X
X
@@ -157,6 +166,7 @@ [email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
@@ -186,6 +196,7 @@ [email protected] X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
@@ -195,8 +206,14 @@ [email protected] X
X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
X
[email protected] X
X
@@ -210,10 +227,16 @@ [email protected] X
[email protected] X
[email protected] X
[email protected] X
X
[email protected] X
X
[email protected]
X
[email protected] X
X
[email protected] X
X
[email protected]
X
[email protected]
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
@@ -243,11 +266,18 @@ [email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected]
X
[email protected] X
X
[email protected] X X
X
@@ -260,8 +290,10 @@ [email protected]
X
[email protected] X
X
[email protected] X X
[email protected] X
X X
[email protected] X
X
[email protected]
X X
[email protected]
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
@@ -280,6 +312,7 @@ [email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
[email protected]
X
[email protected] X
X
@@ -293,6 +326,9 @@ [email protected] X
X
[email protected] X
X
[email protected] X
[email protected] X
X
[email protected]
X
[email protected]
X
[email protected] X
X
[email protected] X
X
[email protected]
X
[email protected]
X
diff --git a/bindings/c/src/blob_reader.rs b/bindings/c/src/blob_reader.rs
new file mode 100644
index 00000000..7c17b70c
--- /dev/null
+++ b/bindings/c/src/blob_reader.rs
@@ -0,0 +1,317 @@
+// 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 std::collections::HashMap;
+use std::ffi::c_void;
+use std::io::SeekFrom;
+
+use paimon::{BlobReader, BlobStream};
+
+use crate::error::{check_non_null, paimon_error, validate_cstr,
PaimonErrorCode};
+use crate::result::{
+ paimon_result_blob_reader, paimon_result_blob_stream,
paimon_result_blob_stream_read,
+ paimon_result_blob_stream_seek, paimon_result_read_blobs,
+};
+use crate::runtime;
+use crate::types::{
+ paimon_blob_reader, paimon_blob_stream, paimon_byte_slice,
paimon_bytes_array, paimon_option,
+ paimon_table,
+};
+
+fn new_reader(reader: BlobReader) -> paimon_result_blob_reader {
+ let reader = Box::new(reader);
+ let wrapper = Box::new(paimon_blob_reader {
+ inner: Box::into_raw(reader) as *mut c_void,
+ });
+ paimon_result_blob_reader {
+ reader: Box::into_raw(wrapper),
+ error: std::ptr::null_mut(),
+ }
+}
+
+fn read_error(error: *mut paimon_error) -> paimon_result_read_blobs {
+ paimon_result_read_blobs {
+ blobs: paimon_bytes_array::empty(),
+ error,
+ }
+}
+
+fn reader_error(error: *mut paimon_error) -> paimon_result_blob_reader {
+ paimon_result_blob_reader {
+ reader: std::ptr::null_mut(),
+ error,
+ }
+}
+
+fn stream_error(error: *mut paimon_error) -> paimon_result_blob_stream {
+ paimon_result_blob_stream {
+ stream: std::ptr::null_mut(),
+ error,
+ }
+}
+
+/// # Safety
+/// `options` is null for zero length or points to valid UTF-8 C-string pairs.
+#[no_mangle]
+pub unsafe extern "C" fn paimon_blob_reader_new(
+ options: *const paimon_option,
+ options_len: usize,
+) -> paimon_result_blob_reader {
+ if options_len > 0 && options.is_null() {
+ return reader_error(paimon_error::new(
+ PaimonErrorCode::InvalidInput,
+ "null pointer passed for `options`".to_string(),
+ ));
+ }
+
+ let mut storage_options = HashMap::with_capacity(options_len);
+ if options_len > 0 {
+ for option in std::slice::from_raw_parts(options, options_len) {
+ let key = match validate_cstr(option.key, "option key") {
+ Ok(value) => value,
+ Err(error) => return reader_error(error),
+ };
+ let value = match validate_cstr(option.value, "option value") {
+ Ok(value) => value,
+ Err(error) => return reader_error(error),
+ };
+ storage_options.insert(key, value);
+ }
+ }
+
+ new_reader(BlobReader::new(storage_options))
+}
+
+/// Create a reader using a table's FileIO.
+///
+/// # Safety
+/// `table` is a valid handle returned by the Paimon C API.
+#[no_mangle]
+pub unsafe extern "C" fn paimon_table_new_blob_reader(
+ table: *const paimon_table,
+) -> paimon_result_blob_reader {
+ if let Err(error) = check_non_null(table, "table") {
+ return reader_error(error);
+ }
+
+ let table = &*((*table).inner as *const paimon::Table);
+ new_reader(BlobReader::from_file_io(table.file_io().clone()))
+}
+
+/// # Safety
+/// The handle and input slices are valid for this call. Free the output with
+/// `paimon_bytes_array_free`.
+#[no_mangle]
+pub unsafe extern "C" fn paimon_blob_reader_read_blobs(
+ reader: *const paimon_blob_reader,
+ descriptors: *const paimon_byte_slice,
+ descriptors_len: usize,
+) -> paimon_result_read_blobs {
+ if let Err(error) = check_non_null(reader, "blob reader") {
+ return read_error(error);
+ }
+ if descriptors_len > 0 && descriptors.is_null() {
+ return read_error(paimon_error::new(
+ PaimonErrorCode::InvalidInput,
+ "null pointer passed for `descriptors`".to_string(),
+ ));
+ }
+
+ let mut owned = Vec::with_capacity(descriptors_len);
+ if descriptors_len > 0 {
+ for (index, descriptor) in std::slice::from_raw_parts(descriptors,
descriptors_len)
+ .iter()
+ .enumerate()
+ {
+ if descriptor.len > 0 && descriptor.data.is_null() {
+ return read_error(paimon_error::new(
+ PaimonErrorCode::InvalidInput,
+ format!(
+ "null data pointer for BlobDescriptor input index
{index}, URI unavailable"
+ ),
+ ));
+ }
+ let bytes = if descriptor.len == 0 {
+ &[]
+ } else {
+ std::slice::from_raw_parts(descriptor.data, descriptor.len)
+ };
+ owned.push(bytes.to_vec());
+ }
+ }
+
+ let reader = &*((*reader).inner as *const BlobReader);
+ match runtime().block_on(reader.read_blobs(&owned)) {
+ Ok(values) => paimon_result_read_blobs {
+ blobs: paimon_bytes_array::new(values),
+ error: std::ptr::null_mut(),
+ },
+ Err(error) => read_error(paimon_error::from_paimon(error)),
+ }
+}
+
+/// Open one descriptor for incremental reads.
+///
+/// # Safety
+/// `reader` is valid and `descriptor` points to `descriptor_len` bytes.
+#[no_mangle]
+pub unsafe extern "C" fn paimon_blob_reader_open_blob(
+ reader: *const paimon_blob_reader,
+ descriptor: *const u8,
+ descriptor_len: usize,
+) -> paimon_result_blob_stream {
+ if let Err(error) = check_non_null(reader, "blob reader") {
+ return stream_error(error);
+ }
+ if descriptor_len > 0 && descriptor.is_null() {
+ return stream_error(paimon_error::new(
+ PaimonErrorCode::InvalidInput,
+ "null pointer passed for `descriptor`".to_string(),
+ ));
+ }
+ let bytes = if descriptor_len == 0 {
+ &[]
+ } else {
+ std::slice::from_raw_parts(descriptor, descriptor_len)
+ };
+ let reader = &*((*reader).inner as *const BlobReader);
+ match reader.open_blob(bytes) {
+ Ok(stream) => {
+ let stream = Box::new(stream);
+ let wrapper = Box::new(paimon_blob_stream {
+ inner: Box::into_raw(stream) as *mut c_void,
+ });
+ paimon_result_blob_stream {
+ stream: Box::into_raw(wrapper),
+ error: std::ptr::null_mut(),
+ }
+ }
+ Err(error) => stream_error(paimon_error::from_paimon(error)),
+ }
+}
+
+/// Read at most `buffer_len` bytes into caller-owned memory.
+///
+/// A zero `bytes_read` result means end of stream when `buffer_len` is
nonzero.
+///
+/// # Safety
+/// `stream` is valid and `buffer` points to `buffer_len` writable bytes.
+#[no_mangle]
+pub unsafe extern "C" fn paimon_blob_stream_read(
+ stream: *mut paimon_blob_stream,
+ buffer: *mut u8,
+ buffer_len: usize,
+) -> paimon_result_blob_stream_read {
+ if let Err(error) = check_non_null(stream, "blob stream") {
+ return paimon_result_blob_stream_read {
+ bytes_read: 0,
+ error,
+ };
+ }
+ if buffer_len > 0 && buffer.is_null() {
+ return paimon_result_blob_stream_read {
+ bytes_read: 0,
+ error: paimon_error::new(
+ PaimonErrorCode::InvalidInput,
+ "null pointer passed for `buffer`".to_string(),
+ ),
+ };
+ }
+
+ let stream = &mut *((*stream).inner as *mut BlobStream);
+ match runtime().block_on(stream.read(buffer_len)) {
+ Ok(bytes) => {
+ if !bytes.is_empty() {
+ std::ptr::copy_nonoverlapping(bytes.as_ptr(), buffer,
bytes.len());
+ }
+ paimon_result_blob_stream_read {
+ bytes_read: bytes.len(),
+ error: std::ptr::null_mut(),
+ }
+ }
+ Err(error) => paimon_result_blob_stream_read {
+ bytes_read: 0,
+ error: paimon_error::from_paimon(error),
+ },
+ }
+}
+
+/// Seek within the descriptor's range. `whence` uses the standard 0, 1, 2
values.
+///
+/// # Safety
+/// `stream` is valid.
+#[no_mangle]
+pub unsafe extern "C" fn paimon_blob_stream_seek(
+ stream: *mut paimon_blob_stream,
+ offset: i64,
+ whence: i32,
+) -> paimon_result_blob_stream_seek {
+ if let Err(error) = check_non_null(stream, "blob stream") {
+ return paimon_result_blob_stream_seek { position: 0, error };
+ }
+ let from = match whence {
+ 0 if offset >= 0 => SeekFrom::Start(offset as u64),
+ 1 => SeekFrom::Current(offset),
+ 2 => SeekFrom::End(offset),
+ _ => {
+ return paimon_result_blob_stream_seek {
+ position: 0,
+ error: paimon_error::new(
+ PaimonErrorCode::InvalidInput,
+ "invalid blob stream seek".to_string(),
+ ),
+ };
+ }
+ };
+ let stream = &mut *((*stream).inner as *mut BlobStream);
+ match runtime().block_on(stream.seek(from)) {
+ Ok(position) => paimon_result_blob_stream_seek {
+ position,
+ error: std::ptr::null_mut(),
+ },
+ Err(error) => paimon_result_blob_stream_seek {
+ position: 0,
+ error: paimon_error::from_paimon(error),
+ },
+ }
+}
+
+/// # Safety
+/// `stream` is null or was returned by `paimon_blob_reader_open_blob`.
+#[no_mangle]
+pub unsafe extern "C" fn paimon_blob_stream_free(stream: *mut
paimon_blob_stream) {
+ if stream.is_null() {
+ return;
+ }
+ let stream = Box::from_raw(stream);
+ if !stream.inner.is_null() {
+ drop(Box::from_raw(stream.inner as *mut BlobStream));
+ }
+}
+
+/// # Safety
+/// `reader` is null or was returned by `paimon_blob_reader_new`.
+#[no_mangle]
+pub unsafe extern "C" fn paimon_blob_reader_free(reader: *mut
paimon_blob_reader) {
+ if reader.is_null() {
+ return;
+ }
+ let reader = Box::from_raw(reader);
+ if !reader.inner.is_null() {
+ drop(Box::from_raw(reader.inner as *mut BlobReader));
+ }
+}
diff --git a/bindings/c/src/lib.rs b/bindings/c/src/lib.rs
index 4974413a..0a5710cc 100644
--- a/bindings/c/src/lib.rs
+++ b/bindings/c/src/lib.rs
@@ -19,6 +19,7 @@
// So it's type node can't meet camel case.
#![allow(non_camel_case_types)]
+mod blob_reader;
mod catalog;
mod error;
mod file_io;
diff --git a/bindings/c/src/result.rs b/bindings/c/src/result.rs
index 667ba85b..94317572 100644
--- a/bindings/c/src/result.rs
+++ b/bindings/c/src/result.rs
@@ -18,6 +18,36 @@
use crate::error::paimon_error;
use crate::types::*;
+#[repr(C)]
+pub struct paimon_result_blob_reader {
+ pub reader: *mut paimon_blob_reader,
+ pub error: *mut paimon_error,
+}
+
+#[repr(C)]
+pub struct paimon_result_blob_stream {
+ pub stream: *mut paimon_blob_stream,
+ pub error: *mut paimon_error,
+}
+
+#[repr(C)]
+pub struct paimon_result_blob_stream_read {
+ pub bytes_read: usize,
+ pub error: *mut paimon_error,
+}
+
+#[repr(C)]
+pub struct paimon_result_blob_stream_seek {
+ pub position: u64,
+ pub error: *mut paimon_error,
+}
+
+#[repr(C)]
+pub struct paimon_result_read_blobs {
+ pub blobs: paimon_bytes_array,
+ pub error: *mut paimon_error,
+}
+
#[repr(C)]
pub struct paimon_result_catalog_new {
pub catalog: *mut paimon_catalog,
diff --git a/bindings/c/src/tests.rs b/bindings/c/src/tests.rs
index 80ce4384..565c3f5a 100644
--- a/bindings/c/src/tests.rs
+++ b/bindings/c/src/tests.rs
@@ -40,9 +40,12 @@ use arrow_array::{Array, Int32Array, RecordBatch,
StringArray, StructArray};
use arrow_schema::{DataType as ArrowDataType, Field as ArrowField, Schema as
ArrowSchema};
use paimon::catalog::Identifier;
use paimon::io::FileIOBuilder;
-use paimon::spec::{CommitKind, DataType, IntType, Schema, TableSchema,
VarCharType};
+use paimon::spec::{
+ BlobDescriptor, CommitKind, DataType, IntType, Schema, TableSchema,
VarCharType,
+};
use paimon::table::{SnapshotManager, Table};
+use crate::blob_reader::*;
use crate::error::*;
use crate::file_io::*;
use crate::table::*;
@@ -3610,3 +3613,218 @@ fn
vector_search_projection_unknown_column_errors_at_execute_read() {
unwrap_table(handle);
}
}
+
+#[test]
+fn blob_reader_reads_batch_and_owns_output_buffers() {
+ let file = tempfile::NamedTempFile::new().unwrap();
+ std::fs::write(file.path(), b"abcdefghij").unwrap();
+ let uri = url::Url::from_file_path(file.path()).unwrap().to_string();
+ let mut descriptors = vec![
+ BlobDescriptor::new(uri.clone(), 3, -1).serialize(),
+ BlobDescriptor::new(uri.clone(), 1, 3).serialize(),
+ BlobDescriptor::new(uri, 5, 0).serialize(),
+ ];
+ let slices = descriptors
+ .iter()
+ .map(|value| paimon_byte_slice {
+ data: value.as_ptr(),
+ len: value.len(),
+ })
+ .collect::<Vec<_>>();
+
+ unsafe {
+ let created = paimon_blob_reader_new(ptr::null(), 0);
+ assert!(created.error.is_null());
+ assert!(!created.reader.is_null());
+
+ let result = paimon_blob_reader_read_blobs(created.reader,
slices.as_ptr(), slices.len());
+ assert!(result.error.is_null());
+ assert_eq!(result.blobs.len, 3);
+
+ descriptors.clear();
+ paimon_blob_reader_free(created.reader);
+ let values = std::slice::from_raw_parts(result.blobs.data,
result.blobs.len)
+ .iter()
+ .map(|value| std::slice::from_raw_parts(value.data,
value.len).to_vec())
+ .collect::<Vec<_>>();
+ assert_eq!(
+ values,
+ vec![b"defghij".to_vec(), b"bcd".to_vec(), Vec::new()]
+ );
+ paimon_bytes_array_free(result.blobs);
+ }
+}
+
+#[test]
+fn blob_reader_from_table_keeps_file_io_alive() {
+ let file_io = memory_file_io();
+ let uri = "memory:/blob_reader_from_table";
+ crate::runtime().block_on(async {
+ file_io
+ .new_output(uri)
+ .unwrap()
+ .write(bytes::Bytes::from_static(b"abcdefghij"))
+ .await
+ .unwrap();
+ });
+ let table = Table::new(
+ file_io,
+ Identifier::new("default", "blob_table"),
+ "memory:/blob_table".to_string(),
+ simple_table_schema(),
+ None,
+ );
+ let table = unsafe { wrap_table(table) };
+ let descriptor = BlobDescriptor::new(uri.to_string(), 2, 4).serialize();
+ let descriptor_slice = paimon_byte_slice {
+ data: descriptor.as_ptr(),
+ len: descriptor.len(),
+ };
+
+ unsafe {
+ let created = paimon_table_new_blob_reader(table);
+ assert!(created.error.is_null());
+ assert!(!created.reader.is_null());
+ unwrap_table(table);
+
+ let result = paimon_blob_reader_read_blobs(created.reader,
&descriptor_slice, 1);
+ assert!(result.error.is_null());
+ let values = std::slice::from_raw_parts(result.blobs.data,
result.blobs.len);
+ assert_eq!(
+ std::slice::from_raw_parts(values[0].data, values[0].len),
+ b"cdef"
+ );
+
+ paimon_bytes_array_free(result.blobs);
+ paimon_blob_reader_free(created.reader);
+ }
+}
+
+#[test]
+fn blob_reader_handles_empty_and_error_batches() {
+ unsafe {
+ let null_table = paimon_table_new_blob_reader(ptr::null());
+ assert!(null_table.reader.is_null());
+ assert!(!null_table.error.is_null());
+ paimon_error_free(null_table.error);
+
+ let created = paimon_blob_reader_new(ptr::null(), 0);
+ assert!(created.error.is_null());
+
+ let empty = paimon_blob_reader_read_blobs(created.reader, ptr::null(),
0);
+ assert!(empty.error.is_null());
+ assert!(empty.blobs.data.is_null());
+ assert_eq!(empty.blobs.len, 0);
+ paimon_bytes_array_free(empty.blobs);
+
+ let invalid_bytes = [0_u8; 1];
+ let invalid_slice = paimon_byte_slice {
+ data: invalid_bytes.as_ptr(),
+ len: invalid_bytes.len(),
+ };
+ let invalid = paimon_blob_reader_read_blobs(created.reader,
&invalid_slice, 1);
+ assert!(!invalid.error.is_null());
+ assert!(invalid.blobs.data.is_null());
+ paimon_error_free(invalid.error);
+
+ let null_slice = paimon_byte_slice {
+ data: ptr::null(),
+ len: 1,
+ };
+ let null_data = paimon_blob_reader_read_blobs(created.reader,
&null_slice, 1);
+ assert!(!null_data.error.is_null());
+ assert!(null_data.blobs.data.is_null());
+ paimon_error_free(null_data.error);
+
+ paimon_blob_reader_free(created.reader);
+ paimon_blob_reader_free(ptr::null_mut());
+ }
+}
+
+#[test]
+fn blob_stream_reads_chunks_and_outlives_reader() {
+ let file = tempfile::NamedTempFile::new().unwrap();
+ std::fs::write(file.path(), b"abcdefghij").unwrap();
+ let uri = url::Url::from_file_path(file.path()).unwrap().to_string();
+ let descriptor = BlobDescriptor::new(uri, 2, 5).serialize();
+
+ unsafe {
+ let created = paimon_blob_reader_new(ptr::null(), 0);
+ assert!(created.error.is_null());
+ let opened =
+ paimon_blob_reader_open_blob(created.reader, descriptor.as_ptr(),
descriptor.len());
+ assert!(opened.error.is_null());
+ assert!(!opened.stream.is_null());
+ paimon_blob_reader_free(created.reader);
+
+ let mut buffer = [0xFF_u8; 3];
+ let first = paimon_blob_stream_read(opened.stream,
buffer.as_mut_ptr(), buffer.len());
+ assert!(first.error.is_null());
+ assert_eq!(first.bytes_read, 3);
+ assert_eq!(&buffer, b"cde");
+
+ let seek = paimon_blob_stream_seek(opened.stream, -2, 2);
+ assert!(seek.error.is_null());
+ assert_eq!(seek.position, 3);
+
+ buffer.fill(0xFF);
+ let second = paimon_blob_stream_read(opened.stream,
buffer.as_mut_ptr(), buffer.len());
+ assert!(second.error.is_null());
+ assert_eq!(second.bytes_read, 2);
+ assert_eq!(&buffer[..2], b"fg");
+ assert_eq!(buffer[2], 0xFF);
+
+ let end = paimon_blob_stream_read(opened.stream, buffer.as_mut_ptr(),
buffer.len());
+ assert!(end.error.is_null());
+ assert_eq!(end.bytes_read, 0);
+
+ paimon_blob_stream_free(opened.stream);
+ paimon_blob_stream_free(ptr::null_mut());
+ }
+}
+
+#[test]
+fn blob_stream_validates_handles_and_buffers() {
+ unsafe {
+ let null_reader = paimon_blob_reader_open_blob(ptr::null(),
ptr::null(), 0);
+ assert!(null_reader.stream.is_null());
+ assert!(!null_reader.error.is_null());
+ paimon_error_free(null_reader.error);
+
+ let created = paimon_blob_reader_new(ptr::null(), 0);
+ let invalid = paimon_blob_reader_open_blob(created.reader,
ptr::null(), 0);
+ assert!(invalid.stream.is_null());
+ assert!(!invalid.error.is_null());
+ paimon_error_free(invalid.error);
+
+ let file = tempfile::NamedTempFile::new().unwrap();
+ let uri = url::Url::from_file_path(file.path()).unwrap().to_string();
+ let descriptor = BlobDescriptor::new(uri, 0, 0).serialize();
+ let opened =
+ paimon_blob_reader_open_blob(created.reader, descriptor.as_ptr(),
descriptor.len());
+ assert!(opened.error.is_null());
+
+ let null_buffer = paimon_blob_stream_read(opened.stream,
ptr::null_mut(), 1);
+ assert!(!null_buffer.error.is_null());
+ paimon_error_free(null_buffer.error);
+
+ let zero = paimon_blob_stream_read(opened.stream, ptr::null_mut(), 0);
+ assert!(zero.error.is_null());
+ assert_eq!(zero.bytes_read, 0);
+
+ let null_stream = paimon_blob_stream_read(ptr::null_mut(),
ptr::null_mut(), 0);
+ assert!(!null_stream.error.is_null());
+ paimon_error_free(null_stream.error);
+
+ let invalid_seek = paimon_blob_stream_seek(opened.stream, -1, 0);
+ assert!(!invalid_seek.error.is_null());
+ paimon_error_free(invalid_seek.error);
+
+ let null_seek = paimon_blob_stream_seek(ptr::null_mut(), 0, 0);
+ assert!(!null_seek.error.is_null());
+ paimon_error_free(null_seek.error);
+
+ paimon_blob_stream_free(opened.stream);
+ paimon_blob_reader_free(created.reader);
+ }
+}
diff --git a/bindings/c/src/types.rs b/bindings/c/src/types.rs
index 6be3e975..4e6d2709 100644
--- a/bindings/c/src/types.rs
+++ b/bindings/c/src/types.rs
@@ -62,6 +62,55 @@ pub unsafe extern "C" fn paimon_bytes_free(bytes:
paimon_bytes) {
}
}
+#[repr(C)]
+#[derive(Clone, Copy)]
+pub struct paimon_byte_slice {
+ pub data: *const u8,
+ pub len: usize,
+}
+
+#[repr(C)]
+pub struct paimon_bytes_array {
+ pub data: *mut paimon_bytes,
+ pub len: usize,
+}
+
+impl paimon_bytes_array {
+ pub fn empty() -> Self {
+ Self {
+ data: std::ptr::null_mut(),
+ len: 0,
+ }
+ }
+
+ pub fn new(values: Vec<Vec<u8>>) -> Self {
+ if values.is_empty() {
+ return Self::empty();
+ }
+ let boxed = values
+ .into_iter()
+ .map(paimon_bytes::new)
+ .collect::<Vec<_>>()
+ .into_boxed_slice();
+ let len = boxed.len();
+ let data = Box::into_raw(boxed) as *mut paimon_bytes;
+ Self { data, len }
+ }
+}
+
+/// # Safety
+/// `array` was returned by `paimon_blob_reader_read_blobs`.
+#[no_mangle]
+pub unsafe extern "C" fn paimon_bytes_array_free(array: paimon_bytes_array) {
+ if array.data.is_null() {
+ return;
+ }
+ let values = Box::from_raw(std::ptr::slice_from_raw_parts_mut(array.data,
array.len));
+ for value in values.iter().copied() {
+ paimon_bytes_free(value);
+ }
+}
+
/// Opaque wrapper around a heap-allocated Rust object.
#[repr(C)]
pub struct paimon_catalog {
@@ -119,6 +168,16 @@ pub struct paimon_file_cache_callbacks_v1 {
pub destroy: Option<unsafe extern "C" fn(context: *mut c_void)>,
}
+#[repr(C)]
+pub struct paimon_blob_reader {
+ pub inner: *mut c_void,
+}
+
+#[repr(C)]
+pub struct paimon_blob_stream {
+ pub inner: *mut c_void,
+}
+
#[repr(C)]
pub struct paimon_identifier {
pub inner: *mut c_void,
diff --git a/bindings/go/DEPENDENCIES.rust.tsv
b/bindings/go/DEPENDENCIES.rust.tsv
index 9a40a09f..f50b9243 100644
--- a/bindings/go/DEPENDENCIES.rust.tsv
+++ b/bindings/go/DEPENDENCIES.rust.tsv
@@ -1,5 +1,6 @@
crate 0BSD Apache-2.0 Apache-2.0 WITH LLVM-exception BSD-2-Clause
BSD-3-Clause BSL-1.0 CC0-1.0 CDLA-Permissive-2.0 ISC
LGPL-2.1-or-later MIT MIT-0 Unicode-3.0 Unlicense Zlib
[email protected] X X
X
[email protected] X
X
[email protected] X
X
[email protected]
X X
[email protected] X
@@ -35,10 +36,12 @@ [email protected] X X
X X X
[email protected] X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
@@ -47,15 +50,18 @@ [email protected] X
X
[email protected] X
X X
[email protected]
X X
[email protected]
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected]
X
[email protected]
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
@@ -63,6 +69,7 @@ [email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
@@ -77,6 +84,7 @@ [email protected] X
X
[email protected]
X
[email protected]
X
[email protected]
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
@@ -143,6 +151,7 @@ [email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected]
X
[email protected] X
X
[email protected] X
X
@@ -157,6 +166,7 @@ [email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
@@ -186,6 +196,7 @@ [email protected] X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
@@ -195,8 +206,14 @@ [email protected] X
X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
X
[email protected] X
X
@@ -210,10 +227,16 @@ [email protected] X
[email protected] X
[email protected] X
[email protected] X
X
[email protected] X
X
[email protected]
X
[email protected] X
X
[email protected] X
X
[email protected]
X
[email protected]
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
@@ -243,11 +266,18 @@ [email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected]
X
[email protected] X
X
[email protected] X X
X
@@ -260,8 +290,10 @@ [email protected]
X
[email protected] X
X
[email protected] X X
[email protected] X
X X
[email protected] X
X
[email protected]
X X
[email protected]
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
@@ -280,6 +312,7 @@ [email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
X
[email protected] X
[email protected]
X
[email protected] X
X
@@ -293,6 +326,9 @@ [email protected] X
X
[email protected] X
X
[email protected] X
[email protected] X
X
[email protected]
X
[email protected]
X
[email protected] X
X
[email protected] X
X
[email protected]
X
[email protected]
X
diff --git a/bindings/go/blob_reader.go b/bindings/go/blob_reader.go
new file mode 100644
index 00000000..7f0071e5
--- /dev/null
+++ b/bindings/go/blob_reader.go
@@ -0,0 +1,227 @@
+/*
+ * 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.
+ */
+
+package paimon
+
+import (
+ "context"
+ "fmt"
+ "runtime"
+ "sync"
+ "unsafe"
+
+ "github.com/jupiterrider/ffi"
+)
+
+// BlobReader resolves serialized BlobDescriptors.
+type BlobReader struct {
+ ctx context.Context
+ lib *libRef
+ inner *paimonBlobReader
+ mu sync.RWMutex
+}
+
+// NewBlobReader creates a descriptor reader with FileIO options.
+func NewBlobReader(storageOptions map[string]string) (*BlobReader, error) {
+ ctx, lib, err := ensureLoaded()
+ if err != nil {
+ return nil, err
+ }
+ inner, err := ffiBlobReaderNew.symbol(ctx)(storageOptions)
+ if err != nil {
+ return nil, err
+ }
+ lib.acquire()
+ return &BlobReader{ctx: ctx, lib: lib, inner: inner}, nil
+}
+
+// NewBlobReader creates a descriptor reader using this table's FileIO.
+func (t *Table) NewBlobReader() (*BlobReader, error) {
+ if t.inner == nil {
+ return nil, ErrClosed
+ }
+ inner, err := ffiTableNewBlobReader.symbol(t.ctx)(t.inner)
+ if err != nil {
+ return nil, err
+ }
+ t.lib.acquire()
+ return &BlobReader{ctx: t.ctx, lib: t.lib, inner: inner}, nil
+}
+
+// ReadBlob resolves one descriptor.
+func (r *BlobReader) ReadBlob(descriptor []byte) ([]byte, error) {
+ values, err := r.ReadBlobs([][]byte{descriptor})
+ if err != nil {
+ return nil, err
+ }
+ return values[0], nil
+}
+
+// ReadBlobs resolves a batch in input order.
+func (r *BlobReader) ReadBlobs(descriptors [][]byte) ([][]byte, error) {
+ r.mu.RLock()
+ defer r.mu.RUnlock()
+ if r.inner == nil {
+ return nil, ErrClosed
+ }
+ return ffiBlobReaderReadBlobs.symbol(r.ctx)(r.inner, descriptors)
+}
+
+// Close releases the reader and is idempotent.
+func (r *BlobReader) Close() {
+ r.mu.Lock()
+ defer r.mu.Unlock()
+ if r.inner == nil {
+ return
+ }
+ ffiBlobReaderFree.symbol(r.ctx)(r.inner)
+ r.inner = nil
+ r.lib.release()
+}
+
+var ffiBlobReaderNew = newFFI(ffiOpts{
+ sym: "paimon_blob_reader_new",
+ rType: &typeResultBlobReader,
+ aTypes: []*ffi.Type{&ffi.TypePointer, &ffi.TypePointer},
+}, func(ctx context.Context, ffiCall ffiCall) func(map[string]string)
(*paimonBlobReader, error) {
+ return func(options map[string]string) (*paimonBlobReader, 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 resultBlobReader
+ ffiCall(
+ unsafe.Pointer(&result),
+ unsafe.Pointer(&optsPtr),
+ unsafe.Pointer(&optsLen),
+ )
+ runtime.KeepAlive(opts)
+ if result.error != nil {
+ return nil, parseError(ctx, result.error)
+ }
+ return result.reader, nil
+ }
+})
+
+var ffiTableNewBlobReader = newFFI(ffiOpts{
+ sym: "paimon_table_new_blob_reader",
+ rType: &typeResultBlobReader,
+ aTypes: []*ffi.Type{&ffi.TypePointer},
+}, func(ctx context.Context, ffiCall ffiCall) func(*paimonTable)
(*paimonBlobReader, error) {
+ return func(table *paimonTable) (*paimonBlobReader, error) {
+ var result resultBlobReader
+ ffiCall(
+ unsafe.Pointer(&result),
+ unsafe.Pointer(&table),
+ )
+ if result.error != nil {
+ return nil, parseError(ctx, result.error)
+ }
+ return result.reader, nil
+ }
+})
+
+var ffiBlobReaderReadBlobs = newFFI(ffiOpts{
+ sym: "paimon_blob_reader_read_blobs",
+ rType: &typeResultReadBlobs,
+ aTypes: []*ffi.Type{
+ &ffi.TypePointer,
+ &ffi.TypePointer,
+ &ffi.TypePointer,
+ },
+}, func(ctx context.Context, ffiCall ffiCall) func(*paimonBlobReader,
[][]byte) ([][]byte, error) {
+ return func(reader *paimonBlobReader, descriptors [][]byte) ([][]byte,
error) {
+ slices := make([]paimonByteSlice, len(descriptors))
+ for index, descriptor := range descriptors {
+ if len(descriptor) > 0 {
+ slices[index].data = &descriptor[0]
+ }
+ slices[index].len = uintptr(len(descriptor))
+ }
+ var slicesPtr unsafe.Pointer
+ if len(slices) > 0 {
+ slicesPtr = unsafe.Pointer(&slices[0])
+ }
+ slicesLen := uintptr(len(slices))
+ var result resultReadBlobs
+ ffiCall(
+ unsafe.Pointer(&result),
+ unsafe.Pointer(&reader),
+ unsafe.Pointer(&slicesPtr),
+ unsafe.Pointer(&slicesLen),
+ )
+ runtime.KeepAlive(descriptors)
+ runtime.KeepAlive(slices)
+ if result.error != nil {
+ return nil, parseError(ctx, result.error)
+ }
+ defer ffiBytesArrayFree.symbol(ctx)(result.blobs)
+ if result.blobs.len > 0 && result.blobs.data == nil {
+ return nil, fmt.Errorf("paimon: native BlobReader
returned a null result array")
+ }
+
+ values := make([][]byte, result.blobs.len)
+ for index, value := range unsafe.Slice(result.blobs.data,
result.blobs.len) {
+ if value.len == 0 {
+ values[index] = []byte{}
+ } else {
+ values[index] = parseBytes(value)
+ }
+ }
+ return values, nil
+ }
+})
+
+var ffiBlobReaderFree = newFFI(ffiOpts{
+ sym: "paimon_blob_reader_free",
+ rType: &ffi.TypeVoid,
+ aTypes: []*ffi.Type{&ffi.TypePointer},
+}, func(_ context.Context, ffiCall ffiCall) func(*paimonBlobReader) {
+ return func(reader *paimonBlobReader) {
+ ffiCall(nil, unsafe.Pointer(&reader))
+ }
+})
+
+var ffiBytesArrayFree = newFFI(ffiOpts{
+ sym: "paimon_bytes_array_free",
+ rType: &ffi.TypeVoid,
+ aTypes: []*ffi.Type{&typePaimonBytesArray},
+}, func(_ context.Context, ffiCall ffiCall) func(paimonBytesArray) {
+ return func(values paimonBytesArray) {
+ ffiCall(nil, unsafe.Pointer(&values))
+ }
+})
diff --git a/bindings/go/blob_stream.go b/bindings/go/blob_stream.go
new file mode 100644
index 00000000..e6a6b7b9
--- /dev/null
+++ b/bindings/go/blob_stream.go
@@ -0,0 +1,204 @@
+/*
+ * 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.
+ */
+
+package paimon
+
+import (
+ "context"
+ "fmt"
+ "io"
+ "runtime"
+ "sync"
+ "unsafe"
+
+ "github.com/jupiterrider/ffi"
+)
+
+// BlobStream incrementally reads one BlobDescriptor.
+type BlobStream struct {
+ ctx context.Context
+ lib *libRef
+ inner *paimonBlobStream
+ mu sync.Mutex
+}
+
+var _ io.ReadSeekCloser = (*BlobStream)(nil)
+
+// OpenBlob opens one descriptor without reading its contents.
+func (r *BlobReader) OpenBlob(descriptor []byte) (*BlobStream, error) {
+ r.mu.RLock()
+ defer r.mu.RUnlock()
+ if r.inner == nil {
+ return nil, ErrClosed
+ }
+ inner, err := ffiBlobReaderOpenBlob.symbol(r.ctx)(r.inner, descriptor)
+ if err != nil {
+ return nil, err
+ }
+ r.lib.acquire()
+ return &BlobStream{ctx: r.ctx, lib: r.lib, inner: inner}, nil
+}
+
+// Read implements io.Reader.
+func (s *BlobStream) Read(buffer []byte) (int, error) {
+ s.mu.Lock()
+ defer s.mu.Unlock()
+ if s.inner == nil {
+ return 0, ErrClosed
+ }
+ if len(buffer) == 0 {
+ return 0, nil
+ }
+
+ read, err := ffiBlobStreamRead.symbol(s.ctx)(s.inner, buffer)
+ if err != nil {
+ return 0, err
+ }
+ if read > len(buffer) {
+ return 0, fmt.Errorf("paimon: native BlobStream returned %d
bytes for a %d-byte buffer", read, len(buffer))
+ }
+ if read == 0 {
+ return 0, io.EOF
+ }
+ return read, nil
+}
+
+// Seek implements io.Seeker within the descriptor's range.
+func (s *BlobStream) Seek(offset int64, whence int) (int64, error) {
+ s.mu.Lock()
+ defer s.mu.Unlock()
+ if s.inner == nil {
+ return 0, ErrClosed
+ }
+ if whence != io.SeekStart && whence != io.SeekCurrent && whence !=
io.SeekEnd {
+ return 0, fmt.Errorf("paimon: invalid BlobStream whence %d",
whence)
+ }
+ position, err := ffiBlobStreamSeek.symbol(s.ctx)(s.inner, offset,
int32(whence))
+ if err != nil {
+ return 0, err
+ }
+ if position > uint64(^uint64(0)>>1) {
+ return 0, fmt.Errorf("paimon: BlobStream position exceeds
int64")
+ }
+ return int64(position), nil
+}
+
+// Close releases the stream and is idempotent.
+func (s *BlobStream) Close() error {
+ s.mu.Lock()
+ defer s.mu.Unlock()
+ if s.inner == nil {
+ return nil
+ }
+ ffiBlobStreamFree.symbol(s.ctx)(s.inner)
+ s.inner = nil
+ s.lib.release()
+ return nil
+}
+
+var ffiBlobReaderOpenBlob = newFFI(ffiOpts{
+ sym: "paimon_blob_reader_open_blob",
+ rType: &typeResultBlobStream,
+ aTypes: []*ffi.Type{
+ &ffi.TypePointer,
+ &ffi.TypePointer,
+ &ffi.TypePointer,
+ },
+}, func(ctx context.Context, ffiCall ffiCall) func(*paimonBlobReader, []byte)
(*paimonBlobStream, error) {
+ return func(reader *paimonBlobReader, descriptor []byte)
(*paimonBlobStream, error) {
+ var descriptorPtr unsafe.Pointer
+ if len(descriptor) > 0 {
+ descriptorPtr = unsafe.Pointer(&descriptor[0])
+ }
+ descriptorLen := uintptr(len(descriptor))
+ var result resultBlobStream
+ ffiCall(
+ unsafe.Pointer(&result),
+ unsafe.Pointer(&reader),
+ unsafe.Pointer(&descriptorPtr),
+ unsafe.Pointer(&descriptorLen),
+ )
+ runtime.KeepAlive(descriptor)
+ if result.error != nil {
+ return nil, parseError(ctx, result.error)
+ }
+ return result.stream, nil
+ }
+})
+
+var ffiBlobStreamRead = newFFI(ffiOpts{
+ sym: "paimon_blob_stream_read",
+ rType: &typeResultBlobStreamRead,
+ aTypes: []*ffi.Type{
+ &ffi.TypePointer,
+ &ffi.TypePointer,
+ &ffi.TypePointer,
+ },
+}, func(ctx context.Context, ffiCall ffiCall) func(*paimonBlobStream, []byte)
(int, error) {
+ return func(stream *paimonBlobStream, buffer []byte) (int, error) {
+ bufferPtr := unsafe.Pointer(&buffer[0])
+ bufferLen := uintptr(len(buffer))
+ var result resultBlobStreamRead
+ ffiCall(
+ unsafe.Pointer(&result),
+ unsafe.Pointer(&stream),
+ unsafe.Pointer(&bufferPtr),
+ unsafe.Pointer(&bufferLen),
+ )
+ runtime.KeepAlive(buffer)
+ if result.error != nil {
+ return 0, parseError(ctx, result.error)
+ }
+ return int(result.bytesRead), nil
+ }
+})
+
+var ffiBlobStreamSeek = newFFI(ffiOpts{
+ sym: "paimon_blob_stream_seek",
+ rType: &typeResultBlobStreamSeek,
+ aTypes: []*ffi.Type{
+ &ffi.TypePointer,
+ &ffi.TypeSint64,
+ &ffi.TypeSint32,
+ },
+}, func(ctx context.Context, ffiCall ffiCall) func(*paimonBlobStream, int64,
int32) (uint64, error) {
+ return func(stream *paimonBlobStream, offset int64, whence int32)
(uint64, error) {
+ var result resultBlobStreamSeek
+ ffiCall(
+ unsafe.Pointer(&result),
+ unsafe.Pointer(&stream),
+ unsafe.Pointer(&offset),
+ unsafe.Pointer(&whence),
+ )
+ if result.error != nil {
+ return 0, parseError(ctx, result.error)
+ }
+ return result.position, nil
+ }
+})
+
+var ffiBlobStreamFree = newFFI(ffiOpts{
+ sym: "paimon_blob_stream_free",
+ rType: &ffi.TypeVoid,
+ aTypes: []*ffi.Type{&ffi.TypePointer},
+}, func(_ context.Context, ffiCall ffiCall) func(*paimonBlobStream) {
+ return func(stream *paimonBlobStream) {
+ ffiCall(nil, unsafe.Pointer(&stream))
+ }
+})
diff --git a/bindings/go/tests/blob_reader_test.go
b/bindings/go/tests/blob_reader_test.go
new file mode 100644
index 00000000..14e3ea39
--- /dev/null
+++ b/bindings/go/tests/blob_reader_test.go
@@ -0,0 +1,261 @@
+/*
+ * 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.
+ */
+
+package paimon_test
+
+import (
+ "bytes"
+ "encoding/binary"
+ "errors"
+ "io"
+ "net/url"
+ "os"
+ "path/filepath"
+ "strings"
+ "testing"
+
+ paimon "github.com/apache/paimon-rust/bindings/go"
+)
+
+func blobDescriptorV2(uri string, offset, length int64) []byte {
+ result := make([]byte, 0, 29+len(uri))
+ result = append(result, 2)
+ result = binary.LittleEndian.AppendUint64(result, 0x424C4F4244455343)
+ result = binary.LittleEndian.AppendUint32(result, uint32(len(uri)))
+ result = append(result, uri...)
+ result = binary.LittleEndian.AppendUint64(result, uint64(offset))
+ result = binary.LittleEndian.AppendUint64(result, uint64(length))
+ return result
+}
+
+func localFileURI(path string) string {
+ return (&url.URL{Scheme: "file", Path: path}).String()
+}
+
+func writeBlobFile(t *testing.T, name, value string) string {
+ t.Helper()
+ path := filepath.Join(t.TempDir(), name)
+ if err := os.WriteFile(path, []byte(value), 0o600); err != nil {
+ t.Fatal(err)
+ }
+ return path
+}
+
+func TestBlobReaderReadBlobAndBatch(t *testing.T) {
+ first := writeBlobFile(t, "first", "abcdefghij")
+ second := writeBlobFile(t, "second", "UVWXYZ")
+
+ reader, err := paimon.NewBlobReader(nil)
+ if err != nil {
+ t.Fatal(err)
+ }
+ defer reader.Close()
+
+ value, err := reader.ReadBlob(blobDescriptorV2(localFileURI(first), 1,
3))
+ if err != nil {
+ t.Fatal(err)
+ }
+ if string(value) != "bcd" {
+ t.Fatalf("ReadBlob returned %q, want %q", value, "bcd")
+ }
+
+ values, err := reader.ReadBlobs([][]byte{
+ blobDescriptorV2(localFileURI(second), 1, 3),
+ blobDescriptorV2(localFileURI(first), 3, -1),
+ blobDescriptorV2(localFileURI(first), 5, 0),
+ blobDescriptorV2(localFileURI(first), 2, 4),
+ blobDescriptorV2(localFileURI(first), 2, 4),
+ })
+ if err != nil {
+ t.Fatal(err)
+ }
+ want := []string{"VWX", "defghij", "", "cdef", "cdef"}
+ for index, value := range values {
+ if string(value) != want[index] {
+ t.Fatalf("ReadBlobs result %d = %q, want %q", index,
value, want[index])
+ }
+ }
+
+ empty, err := reader.ReadBlobs(nil)
+ if err != nil {
+ t.Fatal(err)
+ }
+ if empty == nil || len(empty) != 0 {
+ t.Fatalf("empty batch returned %#v", empty)
+ }
+}
+
+func TestBlobReaderFromTableOutlivesTable(t *testing.T) {
+ file := writeBlobFile(t, "table", "abcdefghij")
+
+ table := openCopiedTestTable(t)
+ reader, err := table.NewBlobReader()
+ if err != nil {
+ t.Fatal(err)
+ }
+ table.Close()
+ defer reader.Close()
+
+ value, err := reader.ReadBlob(blobDescriptorV2(localFileURI(file), 2,
4))
+ if err != nil {
+ t.Fatal(err)
+ }
+ if string(value) != "cdef" {
+ t.Fatalf("ReadBlob returned %q, want %q", value, "cdef")
+ }
+}
+
+func TestBlobReaderErrorsAndClose(t *testing.T) {
+ reader, err := paimon.NewBlobReader(map[string]string{})
+ if err != nil {
+ t.Fatal(err)
+ }
+
+ if _, err := reader.ReadBlob(nil); err == nil {
+ t.Fatal("expected invalid descriptor error")
+ }
+
+ missingURI := localFileURI(t.TempDir() + "/missing.blob")
+ _, err = reader.ReadBlobs([][]byte{
+ blobDescriptorV2(missingURI, 0, 1),
+ })
+ if err == nil {
+ t.Fatal("expected missing object error")
+ }
+ if !strings.Contains(err.Error(), "input indices [0]") ||
!strings.Contains(err.Error(), missingURI) {
+ t.Fatalf("error lacks descriptor context: %v", err)
+ }
+
+ reader.Close()
+ reader.Close()
+ if _, err := reader.ReadBlob(blobDescriptorV2(missingURI, 0, 0));
!errors.Is(err, paimon.ErrClosed) {
+ t.Fatalf("ReadBlob after Close returned %v, want ErrClosed",
err)
+ }
+ if _, err := reader.ReadBlobs(nil); !errors.Is(err, paimon.ErrClosed) {
+ t.Fatalf("ReadBlobs after Close returned %v, want ErrClosed",
err)
+ }
+}
+
+func TestBlobStreamReadsIncrementally(t *testing.T) {
+ file := writeBlobFile(t, "stream", "abcdefghij")
+
+ reader, err := paimon.NewBlobReader(nil)
+ if err != nil {
+ t.Fatal(err)
+ }
+ stream, err := reader.OpenBlob(blobDescriptorV2(localFileURI(file), 2,
5))
+ if err != nil {
+ t.Fatal(err)
+ }
+ reader.Close()
+ if size, err := stream.Seek(0, io.SeekEnd); err != nil || size != 5 {
+ t.Fatalf("SeekEnd returned (%d, %v), want (5, nil)", size, err)
+ }
+ if position, err := stream.Seek(1, io.SeekStart); err != nil ||
position != 1 {
+ t.Fatalf("SeekStart returned (%d, %v), want (1, nil)",
position, err)
+ }
+ var ranged bytes.Buffer
+ if _, err := io.CopyN(&ranged, stream, 3); err != nil {
+ t.Fatal(err)
+ }
+ if ranged.String() != "def" {
+ t.Fatalf("range returned %q, want %q", ranged.String(), "def")
+ }
+ if _, err := stream.Seek(0, io.SeekStart); err != nil {
+ t.Fatal(err)
+ }
+
+ buffer := make([]byte, 2)
+ var value []byte
+ for {
+ read, err := stream.Read(buffer)
+ value = append(value, buffer[:read]...)
+ if errors.Is(err, io.EOF) {
+ break
+ }
+ if err != nil {
+ t.Fatal(err)
+ }
+ }
+ if string(value) != "cdefg" {
+ t.Fatalf("stream returned %q, want %q", value, "cdefg")
+ }
+ if err := stream.Close(); err != nil {
+ t.Fatal(err)
+ }
+ if err := stream.Close(); err != nil {
+ t.Fatal(err)
+ }
+ if _, err := stream.Read(buffer); !errors.Is(err, paimon.ErrClosed) {
+ t.Fatalf("Read after Close returned %v, want ErrClosed", err)
+ }
+ if _, err := stream.Seek(0, io.SeekStart); !errors.Is(err,
paimon.ErrClosed) {
+ t.Fatalf("Seek after Close returned %v, want ErrClosed", err)
+ }
+}
+
+func TestBlobStreamToEndEmptyAndLazyErrors(t *testing.T) {
+ file := writeBlobFile(t, "tail", "abcdefghij")
+
+ reader, err := paimon.NewBlobReader(nil)
+ if err != nil {
+ t.Fatal(err)
+ }
+ defer reader.Close()
+
+ stream, err := reader.OpenBlob(blobDescriptorV2(localFileURI(file), 4,
-1))
+ if err != nil {
+ t.Fatal(err)
+ }
+ value, err := io.ReadAll(stream)
+ if err != nil {
+ t.Fatal(err)
+ }
+ stream.Close()
+ if string(value) != "efghij" {
+ t.Fatalf("stream returned %q, want %q", value, "efghij")
+ }
+
+ empty, err := reader.OpenBlob(blobDescriptorV2(localFileURI(file), 3,
0))
+ if err != nil {
+ t.Fatal(err)
+ }
+ value, err = io.ReadAll(empty)
+ if err != nil {
+ t.Fatal(err)
+ }
+ empty.Close()
+ if len(value) != 0 {
+ t.Fatalf("empty stream returned %q", value)
+ }
+
+ missing := localFileURI(t.TempDir() + "/missing.blob")
+ lazy, err := reader.OpenBlob(blobDescriptorV2(missing, 0, -1))
+ if err != nil {
+ t.Fatalf("OpenBlob performed eager I/O: %v", err)
+ }
+ defer lazy.Close()
+ if _, err := lazy.Read(make([]byte, 1)); err == nil {
+ t.Fatal("expected missing object error on first Read")
+ }
+
+ if _, err := reader.OpenBlob(nil); err == nil {
+ t.Fatal("expected invalid descriptor error")
+ }
+}
diff --git a/bindings/go/types.go b/bindings/go/types.go
index 04fc57ae..6adbac94 100644
--- a/bindings/go/types.go
+++ b/bindings/go/types.go
@@ -29,6 +29,61 @@ import (
// FFI type definitions mirroring C repr structs from paimon-c.
var (
+ typeResultBlobReader = ffi.Type{
+ Type: ffi.Struct,
+ Elements: &[]*ffi.Type{
+ &ffi.TypePointer,
+ &ffi.TypePointer,
+ nil,
+ }[0],
+ }
+
+ typeResultBlobStream = ffi.Type{
+ Type: ffi.Struct,
+ Elements: &[]*ffi.Type{
+ &ffi.TypePointer,
+ &ffi.TypePointer,
+ nil,
+ }[0],
+ }
+
+ typeResultBlobStreamRead = ffi.Type{
+ Type: ffi.Struct,
+ Elements: &[]*ffi.Type{
+ &ffi.TypePointer,
+ &ffi.TypePointer,
+ nil,
+ }[0],
+ }
+
+ typeResultBlobStreamSeek = ffi.Type{
+ Type: ffi.Struct,
+ Elements: &[]*ffi.Type{
+ &ffi.TypeUint64,
+ &ffi.TypePointer,
+ nil,
+ }[0],
+ }
+
+ typePaimonBytesArray = ffi.Type{
+ Type: ffi.Struct,
+ Elements: &[]*ffi.Type{
+ &ffi.TypePointer,
+ &ffi.TypePointer,
+ nil,
+ }[0],
+ }
+
+ typeResultReadBlobs = ffi.Type{
+ Type: ffi.Struct,
+ Elements: &[]*ffi.Type{
+ &ffi.TypePointer,
+ &ffi.TypePointer,
+ &ffi.TypePointer,
+ nil,
+ }[0],
+ }
+
// Result types: { value, *error }
// paimon_result_catalog_new { catalog: paimon_catalog, error:
*paimon_error }
typeResultCatalogNew = ffi.Type{
@@ -243,6 +298,16 @@ type paimonBytes struct {
len uintptr
}
+type paimonByteSlice struct {
+ data *byte
+ len uintptr
+}
+
+type paimonBytesArray struct {
+ data *paimonBytes
+ len uintptr
+}
+
type paimonError struct {
code int32
message paimonBytes
@@ -250,6 +315,8 @@ type paimonError struct {
// Opaque pointer wrappers
type paimonCatalog struct{}
+type paimonBlobReader struct{}
+type paimonBlobStream struct{}
type paimonIdentifier struct{}
type paimonTable struct{}
type paimonReadBuilder struct{}
@@ -273,6 +340,31 @@ type resultCatalogNew struct {
error *paimonError
}
+type resultBlobReader struct {
+ reader *paimonBlobReader
+ error *paimonError
+}
+
+type resultBlobStream struct {
+ stream *paimonBlobStream
+ error *paimonError
+}
+
+type resultBlobStreamRead struct {
+ bytesRead uintptr
+ error *paimonError
+}
+
+type resultBlobStreamSeek struct {
+ position uint64
+ error *paimonError
+}
+
+type resultReadBlobs struct {
+ blobs paimonBytesArray
+ error *paimonError
+}
+
type resultGetTable struct {
table *paimonTable
error *paimonError
diff --git a/docs/src/go-binding.md b/docs/src/go-binding.md
index 26999d9b..5a82c4c6 100644
--- a/docs/src/go-binding.md
+++ b/docs/src/go-binding.md
@@ -39,6 +39,144 @@ go get github.com/apache/paimon-rust/bindings/go
The native library is embedded and loaded automatically. Build with
`CGO_ENABLED=1`.
+## Reading BlobDescriptor Values
+
+`BlobReader` reads serialized `BlobDescriptor` values without scanning a table.
+`ReadBlobs` resolves a batch in one call; `ReadBlob` handles one descriptor.
+
+```go
+package main
+
+import (
+ "database/sql"
+ "log"
+ "os"
+
+ paimon "github.com/apache/paimon-rust/bindings/go"
+ _ "github.com/go-sql-driver/mysql"
+)
+
+func main() {
+ db, err := sql.Open("mysql", os.Getenv("STARROCKS_DSN"))
+ if err != nil {
+ log.Fatal(err)
+ }
+ defer db.Close()
+
+ rows, err := db.Query("SELECT blob_descriptor FROM catalog.db.my_table")
+ if err != nil {
+ log.Fatal(err)
+ }
+ defer rows.Close()
+
+ var descriptors [][]byte
+ for rows.Next() {
+ var descriptor []byte
+ if err := rows.Scan(&descriptor); err != nil {
+ log.Fatal(err)
+ }
+ descriptors = append(descriptors, append([]byte(nil), descriptor...))
+ }
+ if err := rows.Err(); err != nil {
+ log.Fatal(err)
+ }
+
+ reader, err := paimon.NewBlobReader(map[string]string{
+ "fs.oss.accessKeyId": os.Getenv("OSS_ACCESS_KEY_ID"),
+ "fs.oss.accessKeySecret": os.Getenv("OSS_ACCESS_KEY_SECRET"),
+ "fs.oss.endpoint": os.Getenv("OSS_ENDPOINT"),
+ })
+ if err != nil {
+ log.Fatal(err)
+ }
+ defer reader.Close()
+
+ blobs, err := reader.ReadBlobs(descriptors)
+ if err != nil {
+ log.Fatal(err)
+ }
+ for _, blob := range blobs {
+ log.Printf("read %d bytes", len(blob))
+ }
+}
+```
+
+The descriptor contains only URI, offset, and length. Pass OSS/S3 credentials
+with the same FileIO option names used by catalogs. If StarRocks returns a hex
+or base64 string, decode it to the original descriptor bytes before calling
+`ReadBlobs`.
+
+Stream a large value without holding it all in memory:
+
+```go
+stream, err := reader.OpenBlob(descriptor)
+if err != nil {
+ log.Fatal(err)
+}
+defer stream.Close()
+if _, err := io.Copy(destination, stream); err != nil {
+ log.Fatal(err)
+}
+```
+
+For an HTTP byte range, seek relative to the descriptor and copy only that
range:
+
+```go
+size, err := stream.Seek(0, io.SeekEnd)
+if err != nil || start < 0 || end < start || end >= size {
+ log.Fatal("invalid range")
+}
+if _, err := stream.Seek(start, io.SeekStart); err != nil {
+ log.Fatal(err)
+}
+if _, err := io.CopyN(w, stream, end-start+1); err != nil {
+ log.Fatal(err)
+}
+```
+
+`OpenBlob` is lazy and returns an `io.ReadSeekCloser`. `ReadBlobs` groups and
+merges ranges; separate streams are not merged.
+
+For DLF temporary data tokens, reuse a table's refreshing FileIO:
+
+```go
+catalog, err := paimon.NewCatalog(map[string]string{
+ "metastore": "rest",
+ "uri": os.Getenv("DLF_ENDPOINT"),
+ "warehouse": os.Getenv("DLF_CATALOG"),
+ "token.provider": "dlf",
+ "dlf.region": os.Getenv("DLF_REGION"),
+ "dlf.oss-endpoint": os.Getenv("DLF_OSS_ENDPOINT"),
+ "dlf.token-loader": "ecs",
+ "dlf.token-ecs-role-name": os.Getenv("DLF_ECS_ROLE"),
+ "data-token.enabled": "true",
+})
+if err != nil {
+ log.Fatal(err)
+}
+defer catalog.Close()
+table, err := catalog.GetTable(paimon.NewIdentifier("db", "descriptor_table"))
+if err != nil {
+ log.Fatal(err)
+}
+defer table.Close()
+reader, err := table.NewBlobReader()
+if err != nil {
+ log.Fatal(err)
+}
+defer reader.Close()
+```
+
+The reader and its streams keep the table FileIO and refresh DLF data tokens
+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.
+
+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
+budget; use `OpenBlob` for large values. Results retain descriptor input order.
+
## Creating a Catalog
Use `NewCatalog` with a map of options to create a catalog. The catalog type
is determined by the `metastore` option (default: `filesystem`).