From 148b0d7cdaf6ada836cda551125f84f614b3179c Mon Sep 17 00:00:00 2001 From: JingsongLi Date: Mon, 28 Sep 2026 22:59:24 +0800 Subject: [PATCH 1/2] [core] Align dedicated Blob writes and URI resolution with Java --- Cargo.lock | 5 + crates/paimon/Cargo.toml | 2 +- crates/paimon/src/arrow/format/blob.rs | 7 +- crates/paimon/src/io/mod.rs | 1 + crates/paimon/src/io/uri_reader.rs | 353 +++++++++++++++ crates/paimon/src/spec/core_options.rs | 7 + crates/paimon/src/table/blob_resolver.rs | 55 ++- .../paimon/src/table/data_evolution_reader.rs | 8 +- .../src/table/data_file_path_factory.rs | 11 + crates/paimon/src/table/data_file_writer.rs | 37 +- .../src/table/dedicated_format_file_writer.rs | 129 +++--- crates/paimon/src/table/inline_blob.rs | 71 +++ .../paimon/src/table/managed_blob_reader.rs | 9 +- crates/paimon/src/table/mod.rs | 1 + crates/paimon/src/table/table_write.rs | 13 +- .../paimon/tests/dedicated_blob_write_test.rs | 421 ++++++++++++++++++ 16 files changed, 1026 insertions(+), 104 deletions(-) create mode 100644 crates/paimon/src/io/uri_reader.rs create mode 100644 crates/paimon/src/table/inline_blob.rs create mode 100644 crates/paimon/tests/dedicated_blob_write_test.rs diff --git a/Cargo.lock b/Cargo.lock index 383487cc1..e08d17a42 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -7228,12 +7228,17 @@ version = "0.6.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "4cfcf7e2740e6fc6d4d688b4ef00650406bb94adf4731e43c096c3a19fe40840" dependencies = [ + "async-compression", "bitflags 2.13.1", "bytes", + "futures-core", "futures-util", "http", "http-body", + "http-body-util", "pin-project-lite", + "tokio", + "tokio-util", "tower", "tower-layer", "tower-service", diff --git a/crates/paimon/Cargo.toml b/crates/paimon/Cargo.toml index 7c971dcda..7fb7d914a 100644 --- a/crates/paimon/Cargo.toml +++ b/crates/paimon/Cargo.toml @@ -128,7 +128,7 @@ tokio-util = { workspace = true, features = ["compat", "io-util"] } parquet = { workspace = true, features = ["async", "zstd", "lz4", "snap"] } orc-rust = "0.8.0" async-stream = "0.3.6" -reqwest = { version = "0.12", features = ["json"] } +reqwest = { version = "0.12", features = ["json", "gzip", "deflate"] } # DLF authentication dependencies base64 = "0.22" hex = "0.4" diff --git a/crates/paimon/src/arrow/format/blob.rs b/crates/paimon/src/arrow/format/blob.rs index 69d93efbe..fc46f7852 100644 --- a/crates/paimon/src/arrow/format/blob.rs +++ b/crates/paimon/src/arrow/format/blob.rs @@ -2151,12 +2151,12 @@ impl FormatFileWriter for BlobFormatWriter { .to_string(), source: None, })?; - let input = file_io.new_input(desc.uri())?; + let input = crate::io::uri_reader::UriInput::new(file_io, desc.uri())?; let offset = range.offset(); let payload_len = match range.length() { Some(length) => length, None => input - .metadata() + .size() .await .map_err(|e| Error::UnexpectedError { message: format!( @@ -2165,7 +2165,6 @@ impl FormatFileWriter for BlobFormatWriter { ), source: Some(Box::new(e)), })? - .size .saturating_sub(offset), }; let end = offset @@ -2191,7 +2190,7 @@ impl FormatFileWriter for BlobFormatWriter { let reader = if payload_len == 0 { None } else { - Some(input.reader().await?) + Some(input.reader_for_range(offset..end).await?) }; let mut hasher = crc32fast::Hasher::new(); diff --git a/crates/paimon/src/io/mod.rs b/crates/paimon/src/io/mod.rs index 5a075bcd2..c7eea7c12 100644 --- a/crates/paimon/src/io/mod.rs +++ b/crates/paimon/src/io/mod.rs @@ -17,6 +17,7 @@ pub(crate) mod cache; mod file_io; +pub(crate) mod uri_reader; pub use file_io::*; mod storage; diff --git a/crates/paimon/src/io/uri_reader.rs b/crates/paimon/src/io/uri_reader.rs new file mode 100644 index 000000000..2d35dd34b --- /dev/null +++ b/crates/paimon/src/io/uri_reader.rs @@ -0,0 +1,353 @@ +// 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. + +//! URI dispatch for Blob references, matching Java's UriReaderFactory. +//! Ordinary paths retain the table's FileIO, including provider credentials. + +use super::{FileIO, FileRead, InputFile}; +use crate::{Error, Result}; +use bytes::{Buf, Bytes, BytesMut}; +use reqwest::header::CONTENT_ENCODING; +use std::ops::Range; +use std::sync::{Arc, LazyLock}; + +pub(crate) enum UriInput { + File(InputFile), + Http(HttpReader), +} + +impl UriInput { + pub(crate) fn new(file_io: &FileIO, uri: &str) -> Result { + let http = uri.split_once(':').is_some_and(|(scheme, _)| { + scheme.eq_ignore_ascii_case("http") || scheme.eq_ignore_ascii_case("https") + }); + if http { + Ok(Self::Http(HttpReader { uri: uri.into() })) + } else { + Ok(Self::File(file_io.new_input(uri)?)) + } + } + + pub(crate) async fn size(&self) -> Result { + match self { + Self::File(input) => Ok(input.metadata().await?.size), + Self::Http(reader) => reader.size().await, + } + } + + pub(crate) async fn reader(&self) -> Result> { + match self { + Self::File(input) => Ok(Arc::new(input.reader().await?)), + Self::Http(reader) => Ok(Arc::new(reader.clone())), + } + } + + /// Reuse one HTTP response while copying a descriptor into a Blob file. + /// Reopening a Range-ignoring endpoint per chunk would reread every prefix. + pub(crate) async fn reader_for_range(&self, range: Range) -> Result> { + match self { + Self::File(input) => Ok(Arc::new(input.reader().await?)), + Self::Http(reader) => Ok(Arc::new(reader.open(range).await?)), + } + } +} + +static HTTP_CLIENT: LazyLock = LazyLock::new(reqwest::Client::new); + +#[derive(Clone)] +pub(crate) struct HttpReader { + uri: String, +} + +impl HttpReader { + fn request(&self) -> reqwest::RequestBuilder { + HTTP_CLIENT.get(&self.uri) + } + + async fn response(&self) -> Result { + let response = self + .request() + .send() + .await + .map_err(http_error)? + .error_for_status() + .map_err(http_error)?; + if response.status() != reqwest::StatusCode::OK { + return Err(invalid(&format!( + "Unexpected HTTP Blob status: {}", + response.status() + ))); + } + // Reqwest removes Content-Encoding after decoding gzip/deflate. + // Never silently copy an unsupported encoded representation as payload. + if response + .headers() + .get(CONTENT_ENCODING) + .is_some_and(|value| value.as_bytes() != b"identity") + { + return Err(invalid("Unsupported HTTP Blob Content-Encoding")); + } + Ok(response) + } + + async fn size(&self) -> Result { + // GET also works with endpoints that reject HEAD. A chunked or decoded + // response has no known size: count without retaining the payload. + let mut response = self.response().await?; + if let Some(size) = response.content_length() { + return Ok(size); + } + let mut size = 0_u64; + while let Some(chunk) = response.chunk().await.map_err(http_error)? { + size = size + .checked_add(chunk.len() as u64) + .ok_or_else(|| invalid("HTTP Blob size exceeds u64"))?; + } + Ok(size) + } +} + +impl HttpReader { + async fn open(&self, range: Range) -> Result { + if range.start >= range.end { + return Err(invalid("Invalid HTTP Blob range")); + } + // Like Java's HttpUriReader, offsets address the decoded entity. + // A wire Range may instead address compressed bytes, so use a decoded + // GET stream and skip the prefix once for the entire copy. + let response = self.response().await?; + let skip = range.start; + Ok(HttpRangeReader { + end: range.end, + state: tokio::sync::Mutex::new(HttpRangeState { + response, + position: range.start, + skip, + pending: Bytes::new(), + }), + }) + } +} + +#[async_trait::async_trait] +impl FileRead for HttpReader { + async fn read(&self, range: Range) -> Result { + if range.start == range.end { + return Ok(Bytes::new()); + } + self.open(range.clone()).await?.read(range).await + } +} + +struct HttpRangeReader { + end: u64, + state: tokio::sync::Mutex, +} + +struct HttpRangeState { + response: reqwest::Response, + position: u64, + skip: u64, + pending: Bytes, +} + +#[async_trait::async_trait] +impl FileRead for HttpRangeReader { + async fn read(&self, range: Range) -> Result { + let mut state = self.state.lock().await; + if range.start != state.position || range.end < range.start || range.end > self.end { + return Err(invalid("HTTP Blob copy requires consecutive ranges")); + } + let length = range.end - range.start; + let mut remaining = length; + let mut result = BytesMut::new(); + while remaining > 0 { + if state.pending.is_empty() { + let Some(chunk) = state.response.chunk().await.map_err(http_error)? else { + return Err(invalid(&format!( + "HTTP Blob short read: expected {length} bytes, received {}", + length - remaining + ))); + }; + state.pending = chunk; + } + let skip = state.skip.min(state.pending.len() as u64) as usize; + state.skip -= skip as u64; + state.pending.advance(skip); + let count = remaining.min(state.pending.len() as u64) as usize; + result.extend_from_slice(&state.pending[..count]); + state.pending.advance(count); + remaining -= count as u64; + } + state.position = range.end; + Ok(result.freeze()) + } +} + +fn http_error(error: reqwest::Error) -> Error { + Error::UnexpectedError { + message: format!("HTTP Blob request failed: {error}"), + source: Some(Box::new(error)), + } +} + +fn invalid(message: &str) -> Error { + Error::DataInvalid { + message: message.into(), + source: None, + } +} + +#[cfg(test)] +mod tests { + use super::*; + use axum::body::Body; + use axum::http::{Response, StatusCode}; + use axum::response::Redirect; + use axum::routing::get; + use axum::Router; + use std::io::Write; + use std::sync::atomic::{AtomicUsize, Ordering}; + + #[tokio::test] + async fn http_reads_use_decoded_offsets_and_reuse_the_copy_stream() { + let requests = Arc::new(AtomicUsize::new(0)); + let counter = requests.clone(); + let app = Router::new() + .route( + "/copy", + get(move || { + counter.fetch_add(1, Ordering::Relaxed); + async { "0123456789" } + }), + ) + .route("/ignored", get(|| async { "0123456789" })) + .route( + "/gzip", + get(|| async { + let mut encoder = + flate2::write::GzEncoder::new(Vec::new(), flate2::Compression::default()); + encoder.write_all(b"0123456789").unwrap(); + Response::builder() + .header(CONTENT_ENCODING, "gzip") + .body(Body::from(encoder.finish().unwrap())) + .unwrap() + }), + ) + .route( + "/deflate", + get(|| async { + let mut encoder = + flate2::write::ZlibEncoder::new(Vec::new(), flate2::Compression::default()); + encoder.write_all(b"0123456789").unwrap(); + Response::builder() + .header(CONTENT_ENCODING, "deflate") + .body(Body::from(encoder.finish().unwrap())) + .unwrap() + }), + ) + .route("/no-content", get(|| async { StatusCode::NO_CONTENT })) + .route( + "/redirect", + get(|| async { Redirect::temporary("/ignored") }), + ) + .route( + "/chunked", + get(|| async { + Body::from_stream(futures::stream::iter([ + Ok::<_, std::io::Error>(Bytes::from_static(b"012")), + Ok(Bytes::from_static(b"34567")), + Ok(Bytes::from_static(b"89")), + ])) + }), + ) + .route( + "/encoded", + get(|| async { + Response::builder() + .status(StatusCode::OK) + .header(CONTENT_ENCODING, "unknown") + .body(Body::from("0123")) + .unwrap() + }), + ); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let server = tokio::spawn(async move { axum::serve(listener, app).await.unwrap() }); + let io = crate::io::FileIOBuilder::new("memory").build().unwrap(); + for route in ["ignored", "redirect", "chunked", "gzip", "deflate"] { + let input = UriInput::new(&io, &format!("http://{address}/{route}")).unwrap(); + let reader = input.reader().await.unwrap(); + assert_eq!( + reader.read(3..7).await.unwrap().as_ref(), + b"3456", + "{route}" + ); + assert!(reader.read(4..4).await.unwrap().is_empty()); + assert_eq!(input.size().await.unwrap(), 10, "{route}"); + } + let ignored = UriInput::new(&io, &format!("http://{address}/ignored")).unwrap(); + assert!(ignored + .reader() + .await + .unwrap() + .read(8..12) + .await + .unwrap_err() + .to_string() + .contains("short read")); + let copy = UriInput::new(&io, &format!("http://{address}/copy")).unwrap(); + let reader = copy.reader_for_range(1..9).await.unwrap(); + assert_eq!(reader.read(1..4).await.unwrap().as_ref(), b"123"); + assert_eq!(reader.read(4..9).await.unwrap().as_ref(), b"45678"); + assert_eq!(requests.load(Ordering::Relaxed), 1); + let empty_status = UriInput::new(&io, &format!("http://{address}/no-content")).unwrap(); + assert!(empty_status + .size() + .await + .unwrap_err() + .to_string() + .contains("204")); + let encoded = UriInput::new(&io, &format!("http://{address}/encoded")).unwrap(); + assert!(encoded + .reader() + .await + .unwrap() + .read(0..2) + .await + .unwrap_err() + .to_string() + .contains("Content-Encoding")); + let missing = UriInput::new(&io, &format!("http://{address}/missing")).unwrap(); + assert!(missing + .size() + .await + .unwrap_err() + .to_string() + .contains("404")); + assert!(missing + .reader() + .await + .unwrap() + .read(0..1) + .await + .unwrap_err() + .to_string() + .contains("404")); + server.abort(); + } +} diff --git a/crates/paimon/src/spec/core_options.rs b/crates/paimon/src/spec/core_options.rs index d474479b1..59fe532f6 100644 --- a/crates/paimon/src/spec/core_options.rs +++ b/crates/paimon/src/spec/core_options.rs @@ -730,6 +730,13 @@ impl<'a> CoreOptions<'a> { .unwrap_or(false) } + pub fn data_evolution_write_cols_optimization_enabled(&self) -> bool { + self.options + .get("data-evolution.write-cols-optimization.enabled") + .map(|value| value.eq_ignore_ascii_case("true")) + .unwrap_or(false) + } + pub fn data_evolution_nested_field_enabled(&self) -> bool { self.options .get(DATA_EVOLUTION_NESTED_FIELD_ENABLED_OPTION) diff --git a/crates/paimon/src/table/blob_resolver.rs b/crates/paimon/src/table/blob_resolver.rs index cdf2fc7d4..a6bed8830 100644 --- a/crates/paimon/src/table/blob_resolver.rs +++ b/crates/paimon/src/table/blob_resolver.rs @@ -409,10 +409,30 @@ pub(crate) async fn resolve_blob_column( col: &LargeBinaryArray, file_io: &FileIO, limiter: BlobReadLimiter, +) -> Result { + resolve_column(col, file_io, limiter, false).await +} + +/// Resolve a schema-declared descriptor column, including version 1 descriptors +/// which have no magic number and cannot be detected in arbitrary payload bytes. +pub(crate) async fn resolve_descriptor_column( + col: &LargeBinaryArray, + file_io: &FileIO, + limiter: BlobReadLimiter, +) -> Result { + resolve_column(col, file_io, limiter, true).await +} + +async fn resolve_column( + col: &LargeBinaryArray, + file_io: &FileIO, + limiter: BlobReadLimiter, + descriptor_only: bool, ) -> Result { let mut needs_resolve = false; for i in 0..col.len() { - if !col.is_null(i) && BlobDescriptor::is_blob_descriptor(col.value(i)) { + if !col.is_null(i) && (descriptor_only || BlobDescriptor::is_blob_descriptor(col.value(i))) + { needs_resolve = true; break; } @@ -433,7 +453,7 @@ pub(crate) async fn resolve_blob_column( } let value = col.value(row); - if BlobDescriptor::is_blob_descriptor(value) { + if descriptor_only || BlobDescriptor::is_blob_descriptor(value) { let desc = BlobDescriptor::deserialize(value)?; let range = desc.range_spec()?; requests_by_uri @@ -453,17 +473,16 @@ pub(crate) async fn resolve_blob_column( let mut read_groups = Vec::with_capacity(requests_by_uri.len()); for (uri, requests) in requests_by_uri { - let input = file_io.new_input(&uri)?; + let input = crate::io::uri_reader::UriInput::new(file_io, &uri)?; let file_size = if requests.iter().any(|request| request.length.is_none()) { let _metadata_permit = limiter.acquire_request(&uri, "metadata").await?; input - .metadata() + .size() .await .map_err(|e| crate::Error::UnexpectedError { message: format!("Failed to read metadata for BlobDescriptor URI '{uri}': {e}"), source: Some(Box::new(e)), })? - .size } else { 0 }; @@ -498,7 +517,7 @@ pub(crate) async fn resolve_blob_column( continue; } - let reader: Arc = Arc::new(input.reader().await?); + let reader = input.reader().await?; read_groups.push(BlobReadGroup { uri, reader, @@ -715,6 +734,30 @@ fn merge_blob_read_requests(mut requests: Vec) -> Vec>, } @@ -57,6 +60,8 @@ impl DataFilePathFactory { let relative = relative_bucket_path(partition, bucket, directory); Ok(Self { bucket_path, + file_uuid: uuid::Uuid::new_v4(), + file_counter: AtomicU64::new(0), external: ExternalPathProvider::new(options, &relative)?.map(Mutex::new), }) } @@ -65,6 +70,12 @@ impl DataFilePathFactory { &self.bucket_path } + /// Java shares the UUID and counter across all physical column files. + pub fn new_file_name(&self, prefix: &str, format: &str) -> String { + let counter = self.file_counter.fetch_add(1, Ordering::Relaxed); + format!("{prefix}{}-{counter}.{format}", self.file_uuid) + } + pub fn new_path(&self, name: &str) -> Result { let external_path = self .external diff --git a/crates/paimon/src/table/data_file_writer.rs b/crates/paimon/src/table/data_file_writer.rs index 62b7da25e..a28c7078e 100644 --- a/crates/paimon/src/table/data_file_writer.rs +++ b/crates/paimon/src/table/data_file_writer.rs @@ -175,6 +175,7 @@ impl DataFileWriter { return Ok(()); } + super::inline_blob::validate_inline_blob_columns(batch, &self.format_options)?; if self.current_writer.is_none() { self.open_new_file(batch.schema()).await?; } @@ -202,19 +203,19 @@ impl DataFileWriter { Ok(()) } + pub(super) fn has_open_file(&self) -> bool { + self.current_writer.is_some() + } + async fn open_new_file(&mut self, schema: arrow_schema::SchemaRef) -> Result<()> { let index = self .index_options .as_ref() .map(|options| options.create_writer()) .transpose()?; - let file_name = format!( - "{}{}-{}.{}", - self.data_file_prefix, - uuid::Uuid::new_v4(), - self.next_file_ordinal, - self.file_format, - ); + let file_name = self + .paths + .new_file_name(&self.data_file_prefix, &self.file_format); let location = self.paths.new_path(&file_name)?; let bucket_dir = location.parent(); self.file_io.mkdirs(&format!("{bucket_dir}/")).await?; @@ -285,6 +286,13 @@ impl DataFileWriter { let file_source = self.file_source; let first_row_id = self.first_row_id; let write_cols = self.write_cols.clone(); + // Java creates a new row sequence counter for each physical DE file. + let max_sequence_number = if CoreOptions::new(&self.format_options).data_evolution_enabled() + { + row_count - 1 + } else { + 0 + }; Some(async move { let write_result = writer.close().await?; @@ -298,6 +306,7 @@ impl DataFileWriter { write_cols, write_result.value_stats, ); + meta.max_sequence_number = max_sequence_number; meta.external_path = location.external_path; if let Some(index) = index { let bytes = index.serialize()?; @@ -355,11 +364,7 @@ impl DataFileWriter { for (writer, result) in results { match result { Ok(files) => { - for file in files { - for path in file.collect_files(writer.bucket_dir()) { - let _ = writer.file_io.delete_file(&path).await; - } - } + writer.delete_files(&files).await; } Err(error) => { first_error.get_or_insert(error); @@ -411,6 +416,14 @@ impl DataFileWriter { self.written_files.clear(); } + pub(super) async fn delete_files(&mut self, files: &[DataFileMeta]) { + for file in files { + for path in file.collect_files(self.bucket_dir()) { + let _ = self.file_io.delete_file(&path).await; + } + } + } + fn bucket_dir(&self) -> &str { self.paths.bucket_path() } diff --git a/crates/paimon/src/table/dedicated_format_file_writer.rs b/crates/paimon/src/table/dedicated_format_file_writer.rs index 7fcc2d2b5..77b24e23a 100644 --- a/crates/paimon/src/table/dedicated_format_file_writer.rs +++ b/crates/paimon/src/table/dedicated_format_file_writer.rs @@ -17,10 +17,10 @@ use crate::io::FileIO; use crate::resource::ResourceContext; -use crate::spec::{BlobViewStruct, DataField, DataFileMeta, DataType}; +use crate::spec::{CoreOptions, DataField, DataFileMeta, DataType}; use crate::table::data_file_writer::DataFileWriter; use crate::Result; -use arrow_array::{Array, RecordBatch}; +use arrow_array::RecordBatch; use std::collections::{HashMap, HashSet}; use std::sync::Arc; @@ -57,7 +57,8 @@ pub(crate) struct AppendDedicatedFormatFileWriter { vector_writer: Option, normal_column_indices: Vec, normal_schema: Arc, - blob_view_column_indices: Vec<(usize, String)>, + // Completed normal + dedicated file groups, still owned until prepare_commit. + written_files: Vec, } impl AppendDedicatedFormatFileWriter { @@ -81,7 +82,6 @@ impl AppendDedicatedFormatFileWriter { table_fields: &[DataField], format_options: &HashMap, blob_inline_fields: &HashSet, - blob_view_fields: &HashSet, ) -> Result { let paths = Arc::new(super::data_file_path_factory::DataFilePathFactory::new( &table_location, @@ -179,6 +179,20 @@ impl AppendDedicatedFormatFileWriter { None }; + let core_options = CoreOptions::new(format_options); + // Full append writes include all normal fields in schema order. Partial + // column writers must retain their explicit write_cols instead. + let normal_write_cols = if core_options.data_evolution_enabled() + && core_options.data_evolution_write_cols_optimization_enabled() + { + None + } else { + Some(normal_field_names) + }; + let normal_index = + super::data_file_index_writer::FileIndexOptions::parse(format_options, table_fields)? + .and_then(|options| options.project_to_fields(&normal_table_fields)) + .map(Arc::new); let normal_writer = DataFileWriter::new( file_io.clone(), table_location, @@ -194,23 +208,19 @@ impl AppendDedicatedFormatFileWriter { format_options.clone(), Some(0), None, - Some(normal_field_names), + normal_write_cols, )? .with_path_factory(paths.clone()) + .with_file_index(normal_index) .with_target_file_row_num(target_file_row_num); Ok(Self { + written_files: Vec::new(), normal_writer, blob_writers, vector_writer, normal_column_indices, normal_schema, - blob_view_column_indices: table_fields - .iter() - .enumerate() - .filter(|(_, field)| blob_view_fields.contains(field.name())) - .map(|(idx, field)| (idx, field.name().to_string())) - .collect(), }) } @@ -230,8 +240,6 @@ impl AppendDedicatedFormatFileWriter { return Ok(()); } - self.validate_blob_view_columns(batch)?; - // Write normal columns let normal_columns: Vec> = self .normal_column_indices @@ -260,7 +268,11 @@ impl AppendDedicatedFormatFileWriter { ), source: None, })?; - blob_writer.writer.write(&blob_batch).await?; + // Serialized descriptor size says nothing about the payload size. + // Java checks the physical Blob bytes after every row. + for row in 0..blob_batch.num_rows() { + blob_writer.writer.write(&blob_batch.slice(row, 1)).await?; + } } if let Some(vector_writer) = &mut self.vector_writer { @@ -285,51 +297,15 @@ impl AppendDedicatedFormatFileWriter { vector_writer.writer.write(&vector_batch).await?; } - Ok(()) - } - - fn validate_blob_view_columns(&self, batch: &RecordBatch) -> Result<()> { - for (column_index, field_name) in &self.blob_view_column_indices { - let Some(col) = batch - .column(*column_index) - .as_any() - .downcast_ref::() - else { - return Err(crate::Error::DataInvalid { - message: format!( - "blob-view-field '{field_name}' requires a LargeBinaryArray value column" - ), - source: None, - }); - }; - for row in 0..col.len() { - if col.is_null(row) { - continue; - } - let value = col.value(row); - if !BlobViewStruct::is_blob_view_struct(value) { - return Err(crate::Error::DataInvalid { - message: format!( - "blob-view-field '{field_name}' requires blob field value to be a serialized BlobViewStruct" - ), - source: None, - }); - } - let view = BlobViewStruct::deserialize(value)?; - if view.serialize()?.as_slice() != value { - return Err(crate::Error::DataInvalid { - message: format!( - "blob-view-field '{field_name}' contains a non-canonical BlobViewStruct payload" - ), - source: None, - }); - } - } + if !self.normal_writer.has_open_file() { + self.close_group().await?; } Ok(()) } pub(crate) async fn abort(&mut self) { + self.normal_writer.delete_files(&self.written_files).await; + self.written_files.clear(); self.normal_writer.abort().await; for writer in &mut self.blob_writers { writer.writer.abort().await; @@ -340,6 +316,11 @@ impl AppendDedicatedFormatFileWriter { } pub(crate) async fn prepare_commit(&mut self) -> Result> { + self.close_group().await?; + Ok(std::mem::take(&mut self.written_files)) + } + + async fn close_group(&mut self) -> Result<()> { let writers = std::iter::once(&mut self.normal_writer) .chain( self.blob_writers @@ -351,11 +332,14 @@ impl AppendDedicatedFormatFileWriter { .iter_mut() .map(|writer| &mut writer.writer), ); - Ok(DataFileWriter::prepare_group(writers) - .await? - .into_iter() - .flatten() - .collect()) + match DataFileWriter::prepare_group(writers).await { + Ok(files) => self.written_files.extend(files.into_iter().flatten()), + Err(error) => { + self.abort().await; + return Err(error); + } + } + Ok(()) } } @@ -368,7 +352,7 @@ mod tests { #[tokio::test] async fn failed_blob_close_cleans_every_physical_column() { - for external in [false, true] { + for (external, rolled) in [(false, false), (true, false), (false, true), (true, true)] { let io = FileIOBuilder::new("memory").build().unwrap(); let batch = RecordBatch::try_from_iter([ ("id", Arc::new(Int32Array::from(vec![1])) as ArrayRef), @@ -387,7 +371,7 @@ mod tests { DataField::new(1, "a".into(), DataType::Blob(BlobType::new())), DataField::new(2, "b".into(), DataType::Blob(BlobType::new())), ]; - let options = if external { + let mut options = if external { HashMap::from([ ( "data-file.external-paths".into(), @@ -401,6 +385,11 @@ mod tests { } else { HashMap::new() }; + options.extend([ + ("file-index.bloom-filter.columns".into(), "id".into()), + ("file-index.bloom-filter.id.items".into(), "10".into()), + ("file-index.in-manifest-threshold".into(), "0 B".into()), + ]); let mut writer = AppendDedicatedFormatFileWriter::new( io.clone(), "memory:/table".into(), @@ -408,7 +397,7 @@ mod tests { 0, 0, i64::MAX, - i64::MAX, + if rolled { 1 } else { i64::MAX }, i64::MAX, "none".into(), 0, @@ -420,17 +409,19 @@ mod tests { &fields, &options, &HashSet::new(), - &HashSet::new(), ) .unwrap(); writer.write(&batch).await.unwrap(); writer.blob_writers[1].writer.inject_close_failure(); - assert!(writer - .prepare_commit() - .await - .unwrap_err() - .to_string() - .contains("injected close failure")); + let error = if rolled { + assert!(!writer.written_files.is_empty()); + // A later group failure must also remove earlier successful + // groups, which have not been handed to a committer yet. + writer.write(&batch).await.unwrap_err() + } else { + writer.prepare_commit().await.unwrap_err() + }; + assert!(error.to_string().contains("injected close failure")); assert!(io .list_status_recursive("memory:/") .await diff --git a/crates/paimon/src/table/inline_blob.rs b/crates/paimon/src/table/inline_blob.rs new file mode 100644 index 000000000..1cd68c978 --- /dev/null +++ b/crates/paimon/src/table/inline_blob.rs @@ -0,0 +1,71 @@ +// 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. + +//! Validate inline Blob values before any physical file is created. + +use crate::spec::{BlobDescriptor, BlobViewStruct, CoreOptions}; +use crate::{Error, Result}; +use arrow_array::{LargeBinaryArray, RecordBatch}; +use std::collections::HashMap; + +pub(super) fn validate_inline_blob_columns( + batch: &RecordBatch, + options: &HashMap, +) -> Result<()> { + if !options.contains_key("blob-descriptor-field") && !options.contains_key("blob-view-field") { + return Ok(()); + } + let options = CoreOptions::new(options); + for (option, fields) in [ + ("blob-descriptor-field", options.blob_descriptor_fields()), + ("blob-view-field", options.blob_view_fields()), + ] { + for field in fields { + // A partial-column update may not contain this field. + let Some(column) = batch.column_by_name(&field) else { + continue; + }; + let invalid = |message: String| Error::DataInvalid { + message: format!("{option} '{field}': {message}"), + source: None, + }; + let column = column + .as_any() + .downcast_ref::() + .ok_or_else(|| invalid("requires a LargeBinaryArray".into()))?; + for value in column.iter().flatten() { + if option == "blob-descriptor-field" { + // Schema-declared descriptors use Java's prefix parser; + // unlike detection in raw payloads, trailing bytes are valid. + BlobDescriptor::deserialize(value).map_err(|error| { + invalid(format!("requires a serialized BlobDescriptor: {error}")) + })?; + } else { + if !BlobViewStruct::is_blob_view_struct(value) { + return Err(invalid("requires a serialized BlobViewStruct".into())); + } + let view = BlobViewStruct::deserialize(value) + .map_err(|error| invalid(error.to_string()))?; + if view.serialize()?.as_slice() != value { + return Err(invalid("non-canonical BlobViewStruct payload".into())); + } + } + } + } + } + Ok(()) +} diff --git a/crates/paimon/src/table/managed_blob_reader.rs b/crates/paimon/src/table/managed_blob_reader.rs index 6cdaf8083..0d311f959 100644 --- a/crates/paimon/src/table/managed_blob_reader.rs +++ b/crates/paimon/src/table/managed_blob_reader.rs @@ -17,7 +17,7 @@ //! Resolve BLOB descriptors after primary-key merge and row selection. -use super::blob_resolver::{resolve_blob_column, BlobReadLimiter}; +use super::blob_resolver::{resolve_descriptor_column, BlobReadLimiter}; use super::managed_blob_writer::{managed_blob_kind, ManagedBlobKind}; use super::{ArrowRecordBatchStream, Table}; use crate::arrow::format::FilePredicates; @@ -195,7 +195,8 @@ async fn resolve_batch( let column = match kind { ManagedBlobKind::Scalar => { let values = blob_values(columns[index].as_ref())?; - Arc::new(resolve_blob_column(values, file_io, limiter.clone()).await?) as ArrayRef + Arc::new(resolve_descriptor_column(values, file_io, limiter.clone()).await?) + as ArrayRef } ManagedBlobKind::Array => { let array = columns[index] @@ -205,7 +206,7 @@ async fn resolve_batch( let values = blob_values(array.values().as_ref())?; let visible = visible_child_values(values, array.value_offsets(), array)?; let resolved = - Arc::new(resolve_blob_column(&visible, file_io, limiter.clone()).await?); + Arc::new(resolve_descriptor_column(&visible, file_io, limiter.clone()).await?); let ArrowDataType::List(element) = array.data_type() else { unreachable!() }; @@ -227,7 +228,7 @@ async fn resolve_batch( let values = blob_values(map.entries().column(1).as_ref())?; let visible = visible_child_values(values, map.value_offsets(), map)?; let resolved = - Arc::new(resolve_blob_column(&visible, file_io, limiter.clone()).await?); + Arc::new(resolve_descriptor_column(&visible, file_io, limiter.clone()).await?); let ArrowDataType::Map(entries_field, ordered) = map.data_type() else { unreachable!() }; diff --git a/crates/paimon/src/table/mod.rs b/crates/paimon/src/table/mod.rs index e779858ec..5bd178357 100644 --- a/crates/paimon/src/table/mod.rs +++ b/crates/paimon/src/table/mod.rs @@ -73,6 +73,7 @@ mod global_index_types; mod hybrid_search_builder; mod incremental_scan; pub(crate) mod index_file_path; +mod inline_blob; mod kv_file_reader; mod kv_file_writer; mod lumina_index_build_builder; diff --git a/crates/paimon/src/table/table_write.rs b/crates/paimon/src/table/table_write.rs index 80f84d4a0..b2d8085b0 100644 --- a/crates/paimon/src/table/table_write.rs +++ b/crates/paimon/src/table/table_write.rs @@ -460,11 +460,9 @@ impl TableWrite { .any(|f| matches!(f.data_type(), DataType::Vector(_))); let file_index_options = FileIndexOptions::parse(schema.options(), schema.fields())?; - if file_index_options.is_some() - && (has_blob_fields || has_dedicated_vector_fields || !blob_view_fields.is_empty()) - { + if file_index_options.is_some() && has_dedicated_vector_fields { return Err(crate::Error::Unsupported { - message: "FileIndex generation does not support dedicated Blob/Vector writes" + message: "FileIndex generation does not support dedicated Vector writes" .to_string(), }); } @@ -620,6 +618,12 @@ impl TableWrite { if batch.num_rows() == 0 { return Ok(None); } + if !self.primary_key_indices.is_empty() { + super::inline_blob::validate_inline_blob_columns( + &batch, + self.table.schema().options(), + )?; + } let batch = self.enrich_rowkind_batch(&batch)?; Ok((batch.num_rows() != 0).then_some(batch)) } @@ -1171,7 +1175,6 @@ impl TableWrite { fields, self.table.schema().options(), &self.blob_inline_fields, - &self.blob_view_fields, )? .with_resources(self.resources.clone()), ))) diff --git a/crates/paimon/tests/dedicated_blob_write_test.rs b/crates/paimon/tests/dedicated_blob_write_test.rs new file mode 100644 index 000000000..3a7548ab3 --- /dev/null +++ b/crates/paimon/tests/dedicated_blob_write_test.rs @@ -0,0 +1,421 @@ +// 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. + +mod common; + +use arrow_array::{ArrayRef, Int32Array, LargeBinaryArray, RecordBatch}; +use common::incremental_helpers::{memory_table, persist_table_schema, setup_dirs}; +use futures::TryStreamExt; +use paimon::spec::{BlobDescriptor, BlobType, DataType, IntType, Schema, TableSchema}; +use paimon::table::Table; +use std::sync::Arc; + +async fn table(options: &[(&str, &str)]) -> Table { + let mut schema = Schema::builder() + .column("id", DataType::Int(IntType::new())) + .column("payload", DataType::Blob(BlobType::new())) + .option("row-tracking.enabled", "true") + .option("data-evolution.enabled", "true"); + for (key, value) in options { + schema = schema.option(*key, *value); + } + let path = "memory:/blob_write"; + let (io, table) = memory_table(path, TableSchema::new(0, &schema.build().unwrap())); + setup_dirs(&io, path).await; + persist_table_schema(&io, path, table.schema()).await; + table +} + +fn batch(ids: Vec, blobs: Vec>) -> RecordBatch { + RecordBatch::try_from_iter([ + ("id", Arc::new(Int32Array::from(ids)) as ArrayRef), + ( + "payload", + Arc::new(LargeBinaryArray::from(blobs)) as ArrayRef, + ), + ]) + .unwrap() +} + +async fn rows(table: &Table) -> Vec<(i32, Option>)> { + let read = table.new_read_builder(); + let plan = read.new_scan().plan().await.unwrap(); + let batches: Vec = read + .new_read() + .unwrap() + .to_arrow(plan.splits()) + .unwrap() + .try_collect() + .await + .unwrap(); + let mut rows = Vec::new(); + for batch in batches { + let ids = batch + .column(0) + .as_any() + .downcast_ref::() + .unwrap(); + let values = batch + .column(1) + .as_any() + .downcast_ref::() + .unwrap(); + rows.extend( + ids.values() + .iter() + .zip(values.iter()) + .map(|(id, value)| (*id, value.map(Vec::from))), + ); + } + rows.sort_by_key(|row| row.0); + rows +} + +#[tokio::test] +async fn prepared_sequences_and_optimized_write_columns_match_java() { + for optimize in [false, true] { + let table = table(&[( + "data-evolution.write-cols-optimization.enabled", + if optimize { "true" } else { "false" }, + )]) + .await; + let builder = table.new_write_builder(); + let mut writer = builder.new_write().unwrap(); + writer + .write_arrow_batch(&batch(vec![1, 2, 3], vec![Some(b"one"), None, Some(b"")])) + .await + .unwrap(); + let messages = writer.prepare_commit().await.unwrap(); + for file in &messages[0].new_files { + assert_eq!((file.min_sequence_number, file.max_sequence_number), (0, 2)); + let expected = if file.file_name.ends_with(".blob") { + Some(vec!["payload".to_string()]) + } else if optimize { + None + } else { + Some(vec!["id".to_string()]) + }; + assert_eq!(file.write_cols, expected); + } + builder.new_commit().commit(messages).await.unwrap(); + assert_eq!( + rows(&table).await, + vec![(1, Some(b"one".to_vec())), (2, None), (3, Some(vec![]))] + ); + } +} + +#[tokio::test] +async fn normal_rolls_keep_blob_groups_and_row_ids_aligned() { + let table = table(&[("target-file-row-num", "2")]).await; + let builder = table.new_write_builder(); + let mut writer = builder.new_write().unwrap(); + for (ids, blobs) in [ + (vec![1, 2], vec![Some(b"a".as_slice()), None]), + ( + vec![3, 4], + vec![Some(b"b".as_slice()), Some(b"".as_slice())], + ), + ] { + writer.write_arrow_batch(&batch(ids, blobs)).await.unwrap(); + } + let messages = writer.prepare_commit().await.unwrap(); + let files = &messages[0].new_files; + assert_eq!( + files + .iter() + .map(|file| file.file_name.ends_with(".blob")) + .collect::>(), + vec![false, true, false, true] + ); + builder.new_commit().commit(messages).await.unwrap(); + assert_eq!( + rows(&table).await, + vec![ + (1, Some(b"a".to_vec())), + (2, None), + (3, Some(b"b".to_vec())), + (4, Some(vec![])) + ] + ); +} + +#[tokio::test] +async fn blob_size_rolls_inside_a_batch_using_resolved_payload_bytes() { + for descriptor in [false, true] { + let table = table(&[("blob.target-file-size", "32 B")]).await; + let payload = vec![b'x'; 40]; + let source = "memory:/source/payload"; + table + .file_io() + .new_output(source) + .unwrap() + .write(payload.clone().into()) + .await + .unwrap(); + let input = if descriptor { + BlobDescriptor::new(source.into(), 0, 40).serialize() + } else { + payload.clone() + }; + let builder = table.new_write_builder(); + let mut writer = builder.new_write().unwrap(); + writer + .write_arrow_batch(&batch( + vec![1, 2, 3], + vec![Some(&input), None, Some(&input)], + )) + .await + .unwrap(); + let messages = writer.prepare_commit().await.unwrap(); + let blobs: Vec<_> = messages[0] + .new_files + .iter() + .filter(|file| file.file_name.ends_with(".blob")) + .collect(); + assert_eq!( + blobs.iter().map(|file| file.row_count).collect::>(), + vec![1, 2] + ); + let prefix = blobs[0].file_name.rsplit_once('-').unwrap().0; + for file in &blobs { + assert_eq!(file.file_name.rsplit_once('-').unwrap().0, prefix); + } + assert_ne!(blobs[0].file_name, blobs[1].file_name); + builder.new_commit().commit(messages).await.unwrap(); + assert_eq!( + rows(&table).await, + vec![(1, Some(payload.clone())), (2, None), (3, Some(payload))] + ); + } +} + +#[tokio::test] +async fn inline_descriptors_reject_payload_bytes_before_creating_files() { + let table = table(&[("blob-descriptor-field", "payload")]).await; + let mut writer = table.new_write_builder().new_write().unwrap(); + let err = writer + .write_arrow_batch(&batch(vec![1], vec![Some(b"raw payload")])) + .await + .unwrap_err(); + assert!(err.to_string().contains("blob-descriptor-field"), "{err}"); + assert!(table + .file_io() + .list_status_recursive("memory:/blob_write/data") + .await + .unwrap() + .iter() + .all(|entry| !entry.path.ends_with(".parquet"))); +} + +#[tokio::test] +async fn multiple_blob_columns_keep_values_across_independent_rolls() { + let schema = Schema::builder() + .column("id", DataType::Int(IntType::new())) + .column("large", DataType::Blob(BlobType::new())) + .column("small", DataType::Blob(BlobType::new())) + .option("row-tracking.enabled", "true") + .option("data-evolution.enabled", "true") + .option("target-file-row-num", "2") + .option("blob.target-file-size", "32 B") + .option("file-index.bloom-filter.columns", "id") + .option("file-index.bloom-filter.id.items", "10") + .option("file-index.in-manifest-threshold", "0 B") + .build() + .unwrap(); + let path = "memory:/two_blobs"; + let (io, table) = memory_table(path, TableSchema::new(0, &schema)); + setup_dirs(&io, path).await; + persist_table_schema(&io, path, table.schema()).await; + let builder = table.new_write_builder(); + let mut writer = builder.new_write().unwrap(); + let mut expected = Vec::new(); + for start in [0, 2, 4] { + let payload = vec![start as u8; 40]; + let input = RecordBatch::try_from_iter([ + ( + "id", + Arc::new(Int32Array::from(vec![start, start + 1])) as ArrayRef, + ), + ( + "large", + Arc::new(LargeBinaryArray::from(vec![Some(payload.as_slice()), None])) as ArrayRef, + ), + ( + "small", + Arc::new(LargeBinaryArray::from(vec![ + Some(b"".as_slice()), + Some(b"abc".as_slice()), + ])) as ArrayRef, + ), + ]) + .unwrap(); + writer.write_arrow_batch(&input).await.unwrap(); + expected.extend([ + (start, Some(payload), Some(Vec::new())), + (start + 1, None, Some(b"abc".to_vec())), + ]); + } + let messages = writer.prepare_commit().await.unwrap(); + let files = &messages[0].new_files; + let prefix = files[0].file_name.rsplit_once('-').unwrap().0; + let mut suffixes = std::collections::HashSet::new(); + for file in files { + let (file_prefix, suffix) = file.file_name.rsplit_once('-').unwrap(); + assert_eq!(file_prefix, prefix); + assert!(suffixes.insert(suffix.split('.').next().unwrap())); + if file.file_name.ends_with(".parquet") { + assert_eq!(file.extra_files.len(), 1); + } + } + builder.new_commit().commit(messages).await.unwrap(); + let read = table.new_read_builder(); + let plan = read.new_scan().plan().await.unwrap(); + let batches: Vec = read + .new_read() + .unwrap() + .to_arrow(plan.splits()) + .unwrap() + .try_collect() + .await + .unwrap(); + let mut actual = Vec::new(); + for batch in batches { + let ids = batch + .column(0) + .as_any() + .downcast_ref::() + .unwrap(); + let large = batch + .column(1) + .as_any() + .downcast_ref::() + .unwrap(); + let small = batch + .column(2) + .as_any() + .downcast_ref::() + .unwrap(); + actual.extend( + ids.values() + .iter() + .zip(large.iter()) + .zip(small.iter()) + .map(|((id, large), small)| (*id, large.map(Vec::from), small.map(Vec::from))), + ); + } + actual.sort_by_key(|row| row.0); + assert_eq!(actual, expected); +} + +#[tokio::test] +async fn inline_descriptors_accept_v1_v2_and_null_but_reject_malformed_values() { + let table = table(&[("blob-descriptor-field", "payload")]).await; + let source = "memory:/source/payload"; + table + .file_io() + .new_output(source) + .unwrap() + .write(b"hello".to_vec().into()) + .await + .unwrap(); + let v2 = BlobDescriptor::new(source.into(), 0, 5).serialize(); + let mut v1 = vec![1]; + v1.extend_from_slice(&v2[9..]); + v1.extend_from_slice(b"padding"); + let mut v2 = v2; + v2.extend_from_slice(b"padding"); + let builder = table.new_write_builder(); + let mut writer = builder.new_write().unwrap(); + writer + .write_arrow_batch(&batch(vec![1, 2, 3], vec![Some(&v1), Some(&v2), None])) + .await + .unwrap(); + let messages = writer.prepare_commit().await.unwrap(); + assert_eq!(messages[0].new_files.len(), 1); + builder.new_commit().commit(messages).await.unwrap(); + assert_eq!( + rows(&table).await, + vec![ + (1, Some(b"hello".to_vec())), + (2, Some(b"hello".to_vec())), + (3, None) + ] + ); + let mut bad_version = v1; + bad_version[0] = 3; + for invalid in [ + vec![], + b"payload".to_vec(), + v2[..v2.len() - 8].to_vec(), + bad_version, + ] { + let mut writer = builder.new_write().unwrap(); + assert!(writer + .write_arrow_batch(&batch(vec![4], vec![Some(&invalid)])) + .await + .is_err()); + } + assert_eq!(rows(&table).await.len(), 3); +} + +#[tokio::test] +async fn primary_key_inline_descriptors_validate_and_resolve_legacy_bytes() { + let schema = Schema::builder() + .column("id", DataType::Int(IntType::with_nullable(false))) + .column("payload", DataType::Blob(BlobType::new())) + .primary_key(["id"]) + .option("bucket", "1") + .option("blob-descriptor-field", "payload") + .build() + .unwrap(); + let path = "memory:/pk_inline_blob"; + let (io, table) = memory_table(path, TableSchema::new(0, &schema)); + setup_dirs(&io, path).await; + persist_table_schema(&io, path, table.schema()).await; + let source = "memory:/source/payload"; + io.new_output(source) + .unwrap() + .write(b"hello".to_vec().into()) + .await + .unwrap(); + let v2 = BlobDescriptor::new(source.into(), 0, 5).serialize(); + let mut v1 = vec![1]; + v1.extend_from_slice(&v2[9..]); + v1.extend_from_slice(b"padding"); + let builder = table.new_write_builder(); + let mut writer = builder.new_write().unwrap(); + let error = writer + .write_arrow_batch(&batch(vec![1], vec![Some(b"raw")])) + .await + .unwrap_err(); + assert!(error.to_string().contains("blob-descriptor-field")); + writer + .write_arrow_batch(&batch(vec![1, 2, 3], vec![Some(&v1), Some(&v2), None])) + .await + .unwrap(); + let messages = writer.prepare_commit().await.unwrap(); + builder.new_commit().commit(messages).await.unwrap(); + assert_eq!( + rows(&table).await, + vec![ + (1, Some(b"hello".to_vec())), + (2, Some(b"hello".to_vec())), + (3, None) + ] + ); +} From ae58da0f07184266b3abc6f3f3517105b1f96adf Mon Sep 17 00:00:00 2001 From: JingsongLi Date: Tue, 29 Sep 2026 10:16:59 +0800 Subject: [PATCH 2/2] [core] Fix Blob review feedback and descriptor update regressions --- bindings/c/DEPENDENCIES.rust.tsv | 3 + bindings/go/DEPENDENCIES.rust.tsv | 3 + .../integration_tests/DEPENDENCIES.rust.tsv | 3 + .../datafusion/tests/blob_tests.rs | 88 +++++++++--- .../paimon-rest-server/DEPENDENCIES.rust.tsv | 3 + crates/paimon/DEPENDENCIES.rust.tsv | 3 + crates/paimon/src/arrow/format/blob.rs | 6 +- crates/paimon/src/io/uri_reader.rs | 57 +++++++- crates/paimon/src/table/blob_resolver.rs | 25 ++-- .../paimon/src/table/data_evolution_writer.rs | 8 +- crates/paimon/src/table/table_write.rs | 6 +- crates/paimon/tests/blob_http_error_test.rs | 133 +++++++++++++++++ .../paimon/tests/dedicated_blob_write_test.rs | 135 ++++++++++++++++++ 13 files changed, 428 insertions(+), 45 deletions(-) create mode 100644 crates/paimon/tests/blob_http_error_test.rs diff --git a/bindings/c/DEPENDENCIES.rust.tsv b/bindings/c/DEPENDENCIES.rust.tsv index a60e13faa..149d0a46f 100644 --- a/bindings/c/DEPENDENCIES.rust.tsv +++ b/bindings/c/DEPENDENCIES.rust.tsv @@ -24,6 +24,7 @@ arrow-row@58.4.0 X arrow-schema@58.4.0 X arrow-select@58.4.0 X arrow-string@58.4.0 X +async-compression@0.4.42 X X async-stream@0.3.6 X async-stream-impl@0.3.6 X async-trait@0.1.91 X X @@ -61,6 +62,8 @@ cmake@0.1.58 X X cmov@0.5.4 X X combine@4.6.7 X comfy-table@7.2.2 X +compression-codecs@0.4.38 X X +compression-core@0.4.32 X X const-oid@0.10.2 X X const-oid@0.9.6 X X const-random@0.1.18 X X diff --git a/bindings/go/DEPENDENCIES.rust.tsv b/bindings/go/DEPENDENCIES.rust.tsv index a60e13faa..149d0a46f 100644 --- a/bindings/go/DEPENDENCIES.rust.tsv +++ b/bindings/go/DEPENDENCIES.rust.tsv @@ -24,6 +24,7 @@ arrow-row@58.4.0 X arrow-schema@58.4.0 X arrow-select@58.4.0 X arrow-string@58.4.0 X +async-compression@0.4.42 X X async-stream@0.3.6 X async-stream-impl@0.3.6 X async-trait@0.1.91 X X @@ -61,6 +62,8 @@ cmake@0.1.58 X X cmov@0.5.4 X X combine@4.6.7 X comfy-table@7.2.2 X +compression-codecs@0.4.38 X X +compression-core@0.4.32 X X const-oid@0.10.2 X X const-oid@0.9.6 X X const-random@0.1.18 X X diff --git a/crates/integration_tests/DEPENDENCIES.rust.tsv b/crates/integration_tests/DEPENDENCIES.rust.tsv index ede3938da..abe736638 100644 --- a/crates/integration_tests/DEPENDENCIES.rust.tsv +++ b/crates/integration_tests/DEPENDENCIES.rust.tsv @@ -23,6 +23,7 @@ arrow-row@58.4.0 X arrow-schema@58.4.0 X arrow-select@58.4.0 X arrow-string@58.4.0 X +async-compression@0.4.42 X X async-stream@0.3.6 X async-stream-impl@0.3.6 X async-trait@0.1.91 X X @@ -56,6 +57,8 @@ cmake@0.1.58 X X cmov@0.5.4 X X combine@4.6.7 X comfy-table@7.2.2 X +compression-codecs@0.4.38 X X +compression-core@0.4.32 X X const-oid@0.10.2 X X const-random@0.1.18 X X const-random-macro@0.1.16 X X diff --git a/crates/integrations/datafusion/tests/blob_tests.rs b/crates/integrations/datafusion/tests/blob_tests.rs index 47364ead7..8b2a9aed4 100644 --- a/crates/integrations/datafusion/tests/blob_tests.rs +++ b/crates/integrations/datafusion/tests/blob_tests.rs @@ -35,6 +35,14 @@ fn to_hex(bytes: &[u8]) -> String { bytes.iter().map(|b| format!("{b:02X}")).collect() } +// Inline descriptor columns contain references, not arbitrary payload bytes. +fn descriptor_hex(directory: &std::path::Path, name: &str, payload: &[u8]) -> String { + let path = directory.join(name); + std::fs::write(&path, payload).unwrap(); + let uri = path.to_str().unwrap().to_string(); + to_hex(&BlobDescriptor::new(uri, 0, payload.len() as i64).serialize()) +} + async fn setup(table_ddl: &str) -> (tempfile::TempDir, SQLContext) { let (tmp, catalog) = create_test_env(); let sql_context = create_sql_context(catalog).await; @@ -249,7 +257,7 @@ async fn test_blob_multiple_inserts() { /// blob-descriptor-field: listed fields are stored inline in parquet (no .blob files). #[tokio::test] async fn test_blob_descriptor_field_inline() { - let (_tmp, sql_context) = setup( + let (tmp, sql_context) = setup( "CREATE TABLE paimon.test_db.t (\ id INT, \ name STRING, \ @@ -262,9 +270,20 @@ async fn test_blob_descriptor_field_inline() { ) .await; + let hello_hex = descriptor_hex(tmp.path(), "hello.bin", b"Hello"); + + assert_sql_error( + &sql_context, + "INSERT INTO paimon.test_db.t (id, name, picture) VALUES (0, 'Invalid', X'48656C6C6F')", + "blob-descriptor-field", + ) + .await; + exec( &sql_context, - "INSERT INTO paimon.test_db.t (id, name, picture) VALUES (1, 'Alice', X'48656C6C6F')", + &format!( + "INSERT INTO paimon.test_db.t (id, name, picture) VALUES (1, 'Alice', X'{hello_hex}')" + ), ) .await; @@ -357,7 +376,7 @@ async fn test_merge_into_rejects_raw_blob_update() { /// Reference: BlobTestBase "Blob: merge-into updates non-blob column on descriptor blob table" #[tokio::test] async fn test_merge_into_updates_non_blob_on_descriptor_table() { - let (_tmp, sql_context) = setup( + let (tmp, sql_context) = setup( "CREATE TABLE paimon.test_db.t (\ id INT, \ name STRING, \ @@ -370,11 +389,16 @@ async fn test_merge_into_updates_non_blob_on_descriptor_table() { ) .await; + let first_hex = descriptor_hex(tmp.path(), "first.bin", b"AA"); + let second_hex = descriptor_hex(tmp.path(), "second.bin", b"BB"); + exec( &sql_context, - "INSERT INTO paimon.test_db.t (id, name, picture) VALUES \ - (1, 'Alice', X'4141'), \ - (2, 'Bob', X'4242')", + &format!( + "INSERT INTO paimon.test_db.t (id, name, picture) VALUES \ + (1, 'Alice', X'{first_hex}'), \ + (2, 'Bob', X'{second_hex}')" + ), ) .await; @@ -406,7 +430,7 @@ async fn test_merge_into_updates_non_blob_on_descriptor_table() { /// Merge-into on a descriptor blob table: updating the blob column should succeed. #[tokio::test] async fn test_merge_into_updates_blob_on_descriptor_table() { - let (_tmp, sql_context) = setup( + let (tmp, sql_context) = setup( "CREATE TABLE paimon.test_db.t (\ id INT, \ name STRING, \ @@ -419,17 +443,23 @@ async fn test_merge_into_updates_blob_on_descriptor_table() { ) .await; + let first_hex = descriptor_hex(tmp.path(), "first.bin", b"AA"); + let second_hex = descriptor_hex(tmp.path(), "second.bin", b"BB"); + let updated_hex = descriptor_hex(tmp.path(), "updated.bin", b"CC"); + exec( &sql_context, - "INSERT INTO paimon.test_db.t (id, name, picture) VALUES \ - (1, 'Alice', X'4141'), \ - (2, 'Bob', X'4242')", + &format!( + "INSERT INTO paimon.test_db.t (id, name, picture) VALUES \ + (1, 'Alice', X'{first_hex}'), \ + (2, 'Bob', X'{second_hex}')" + ), ) .await; exec( &sql_context, - "CREATE TEMPORARY TABLE paimon.test_db.src AS SELECT * FROM (VALUES (1, X'4343')) AS t(id, picture)", + &format!("CREATE TEMPORARY TABLE paimon.test_db.src AS SELECT * FROM (VALUES (1, X'{updated_hex}')) AS t(id, picture)"), ) .await; @@ -673,7 +703,7 @@ async fn test_blob_rolling() { /// blob-descriptor-field with multiple inserts: descriptor values are resolved to actual data. #[tokio::test] async fn test_blob_descriptor_field_resolve_on_read() { - let (_tmp, sql_context) = setup( + let (tmp, sql_context) = setup( "CREATE TABLE paimon.test_db.t (\ id INT, \ name STRING, \ @@ -686,20 +716,28 @@ async fn test_blob_descriptor_field_resolve_on_read() { ) .await; + let hello_hex = descriptor_hex(tmp.path(), "hello.bin", b"Hello"); + let world_hex = descriptor_hex(tmp.path(), "world.bin", b"World"); + let paimon_hex = descriptor_hex(tmp.path(), "paimon.bin", b"Paimon"); + exec( &sql_context, - "INSERT INTO paimon.test_db.t (id, name, picture) VALUES \ - (1, 'Alice', X'48656C6C6F'), \ + &format!( + "INSERT INTO paimon.test_db.t (id, name, picture) VALUES \ + (1, 'Alice', X'{hello_hex}'), \ (2, 'Bob', NULL), \ - (3, 'Carol', X'576F726C64')", + (3, 'Carol', X'{world_hex}')" + ), ) .await; // Second insert to exercise multi-file merge path. exec( &sql_context, - "INSERT INTO paimon.test_db.t (id, name, picture) VALUES \ - (4, 'Dave', X'5061696D6F6E')", + &format!( + "INSERT INTO paimon.test_db.t (id, name, picture) VALUES \ + (4, 'Dave', X'{paimon_hex}')" + ), ) .await; @@ -735,6 +773,8 @@ async fn test_blob_descriptor_field_resolve_descriptor_value() { ) .await; + let hello_hex = descriptor_hex(tmp.path(), "hello.bin", b"Hello"); + // Write a source file that the BlobDescriptor will reference. let source_data = b"DescriptorResolved"; let source_path = tmp.path().join("desc_source.bin"); @@ -747,7 +787,7 @@ async fn test_blob_descriptor_field_resolve_descriptor_value() { let sql = format!( "INSERT INTO paimon.test_db.t (id, name, picture) VALUES \ (1, 'FromDesc', X'{desc_hex}'), \ - (2, 'Raw', X'48656C6C6F')" + (2, 'FromSecondDesc', X'{hello_hex}')" ); exec(&sql_context, &sql).await; @@ -760,7 +800,7 @@ async fn test_blob_descriptor_field_resolve_descriptor_value() { rows, vec![ (1, "FromDesc".into(), Some(b"DescriptorResolved".to_vec())), - (2, "Raw".into(), Some(b"Hello".to_vec())), + (2, "FromSecondDesc".into(), Some(b"Hello".to_vec())), ] ); } @@ -780,6 +820,8 @@ async fn test_blob_descriptor_field_resolve_unknown_length_descriptor() { ) .await; + let known_hex = descriptor_hex(tmp.path(), "known.bin", b"RAW"); + let source_data = b"HEADER_PAYLOAD_TRAILER"; let source_path = tmp.path().join("descriptor_unknown_length.bin"); std::fs::write(&source_path, source_data).unwrap(); @@ -797,7 +839,7 @@ async fn test_blob_descriptor_field_resolve_unknown_length_descriptor() { (2, 'Suffix', X'{suffix_hex}'), \ (3, 'Eof', X'{eof_hex}'), \ (4, 'PastEof', X'{past_eof_hex}'), \ - (5, 'Raw', X'524157'), \ + (5, 'Known', X'{known_hex}'), \ (6, 'Null', NULL)" ); exec(&sql_context, &sql).await; @@ -814,7 +856,7 @@ async fn test_blob_descriptor_field_resolve_unknown_length_descriptor() { (2, "Suffix".into(), Some(b"PAYLOAD_TRAILER".to_vec())), (3, "Eof".into(), Some(Vec::new())), (4, "PastEof".into(), Some(Vec::new())), - (5, "Raw".into(), Some(b"RAW".to_vec())), + (5, "Known".into(), Some(b"RAW".to_vec())), (6, "Null".into(), None), ] ); @@ -897,13 +939,15 @@ async fn test_blob_descriptor_filter_before_resolve_by_default() { ) .await; + let kept_hex = descriptor_hex(tmp.path(), "kept.bin", b"OK"); + let missing_uri = format!("file://{}", tmp.path().join("missing_blob.bin").display()); let bad_desc = BlobDescriptor::new(missing_uri, 0, 1); let bad_desc_hex = to_hex(&bad_desc.serialize()); let sql = format!( "INSERT INTO paimon.test_db.t (id, name, picture) VALUES \ (1, 'Filtered', X'{bad_desc_hex}'), \ - (2, 'Kept', X'4F4B')" + (2, 'Kept', X'{kept_hex}')" ); exec(&sql_context, &sql).await; let rows = query_id_name_picture( diff --git a/crates/paimon-rest-server/DEPENDENCIES.rust.tsv b/crates/paimon-rest-server/DEPENDENCIES.rust.tsv index 3b4730a4b..2e6184706 100644 --- a/crates/paimon-rest-server/DEPENDENCIES.rust.tsv +++ b/crates/paimon-rest-server/DEPENDENCIES.rust.tsv @@ -23,6 +23,7 @@ arrow-row@58.4.0 X arrow-schema@58.4.0 X arrow-select@58.4.0 X arrow-string@58.4.0 X +async-compression@0.4.42 X X async-stream@0.3.6 X async-stream-impl@0.3.6 X async-trait@0.1.91 X X @@ -59,6 +60,8 @@ cmake@0.1.58 X X cmov@0.5.4 X X combine@4.6.7 X comfy-table@7.2.2 X +compression-codecs@0.4.38 X X +compression-core@0.4.32 X X const-oid@0.10.2 X X const-random@0.1.18 X X const-random-macro@0.1.16 X X diff --git a/crates/paimon/DEPENDENCIES.rust.tsv b/crates/paimon/DEPENDENCIES.rust.tsv index af93908d6..b05ba1276 100644 --- a/crates/paimon/DEPENDENCIES.rust.tsv +++ b/crates/paimon/DEPENDENCIES.rust.tsv @@ -29,6 +29,7 @@ arrow-schema@58.4.0 X arrow-select@58.4.0 X arrow-string@58.4.0 X async-channel@2.5.0 X X +async-compression@0.4.42 X X async-executor@1.14.0 X X async-fs@2.2.0 X X async-io@2.6.0 X X @@ -89,6 +90,8 @@ cmake@0.1.58 X X cmov@0.5.4 X X combine@4.6.7 X comfy-table@7.2.2 X +compression-codecs@0.4.38 X X +compression-core@0.4.32 X X concurrent-queue@2.5.0 X X const-oid@0.10.2 X X const-oid@0.9.6 X X diff --git a/crates/paimon/src/arrow/format/blob.rs b/crates/paimon/src/arrow/format/blob.rs index fc46f7852..2280b9c12 100644 --- a/crates/paimon/src/arrow/format/blob.rs +++ b/crates/paimon/src/arrow/format/blob.rs @@ -2161,7 +2161,7 @@ impl FormatFileWriter for BlobFormatWriter { .map_err(|e| Error::UnexpectedError { message: format!( "Failed to read metadata for BlobDescriptor '{}': {e}", - desc.uri() + crate::io::uri_reader::sanitize_blob_uri(desc.uri()) ), source: Some(Box::new(e)), })? @@ -2209,7 +2209,7 @@ impl FormatFileWriter for BlobFormatWriter { Error::UnexpectedError { message: format!( "Failed to read BlobDescriptor '{}' range {pos}..{chunk_end}: {e}", - desc.uri() + crate::io::uri_reader::sanitize_blob_uri(desc.uri()) ), source: Some(Box::new(e)), } @@ -2220,7 +2220,7 @@ impl FormatFileWriter for BlobFormatWriter { return Err(Error::DataInvalid { message: format!( "Failed to read BlobDescriptor '{}': short read for range {pos}..{chunk_end}, expected={expected_len} bytes, actual={actual_len} bytes", - desc.uri() + crate::io::uri_reader::sanitize_blob_uri(desc.uri()) ), source: None, }); diff --git a/crates/paimon/src/io/uri_reader.rs b/crates/paimon/src/io/uri_reader.rs index 2d35dd34b..6324671d2 100644 --- a/crates/paimon/src/io/uri_reader.rs +++ b/crates/paimon/src/io/uri_reader.rs @@ -199,10 +199,51 @@ impl FileRead for HttpRangeReader { } fn http_error(error: reqwest::Error) -> Error { + // Both Display and Debug/source chains can expose the final URL after a + // redirect. Like Java HttpClientUtils, retain safe diagnostics only; even + // removing the request URL would not sanitize arbitrary nested causes. + let message = if let Some(status) = error.status() { + format!("HTTP Blob request failed with status {status}") + } else { + let operation = if error.is_timeout() { + "timeout" + } else if error.is_redirect() { + "redirect" + } else if error.is_connect() { + "connection" + } else if error.is_decode() { + "response decoding" + } else if error.is_builder() { + "request construction" + } else if error.is_body() { + "response body" + } else { + "request" + }; + format!("HTTP Blob {operation} failed") + }; Error::UnexpectedError { - message: format!("HTTP Blob request failed: {error}"), - source: Some(Box::new(error)), + message, + source: None, + } +} + +/// Safe URI context for Blob I/O errors, including signed URLs and userinfo. +pub(crate) fn sanitize_blob_uri(uri: &str) -> String { + if let Ok(mut url) = url::Url::parse(uri) { + let _ = url.set_username(""); + let _ = url.set_password(None); + url.set_query(None); + url.set_fragment(None); + return url.to_string(); + } + if uri.split_once(':').is_some_and(|(scheme, _)| { + scheme.eq_ignore_ascii_case("http") || scheme.eq_ignore_ascii_case("https") + }) { + // A malformed authority cannot be safely split into host/userinfo. + return "".into(); } + uri.split(['?', '#']).next().unwrap_or(uri).to_string() } fn invalid(message: &str) -> Error { @@ -223,6 +264,18 @@ mod tests { use std::io::Write; use std::sync::atomic::{AtomicUsize, Ordering}; + #[test] + fn sanitize_http_uri_context() { + assert_eq!( + sanitize_blob_uri("https://user:password@example.com/path?signature=secret#fragment"), + "https://example.com/path" + ); + assert_eq!( + sanitize_blob_uri("http://user:password@example.com:bad/path?signature=secret"), + "" + ); + } + #[tokio::test] async fn http_reads_use_decoded_offsets_and_reuse_the_copy_stream() { let requests = Arc::new(AtomicUsize::new(0)); diff --git a/crates/paimon/src/table/blob_resolver.rs b/crates/paimon/src/table/blob_resolver.rs index a6bed8830..863c65533 100644 --- a/crates/paimon/src/table/blob_resolver.rs +++ b/crates/paimon/src/table/blob_resolver.rs @@ -16,6 +16,7 @@ // under the License. use crate::arrow::format::blob::DEFAULT_BLOB_READ_PARALLELISM; +use crate::io::uri_reader::sanitize_blob_uri; use crate::io::{FileIO, FileRead}; use crate::spec::BlobDescriptor; use crate::Result; @@ -303,17 +304,6 @@ fn blob_error_with_context( } } -fn sanitize_blob_uri(uri: &str) -> String { - if let Ok(mut url) = url::Url::parse(uri) { - let _ = url.set_username(""); - let _ = url.set_password(None); - url.set_query(None); - url.set_fragment(None); - return url.to_string(); - } - uri.split(['?', '#']).next().unwrap_or(uri).to_string() -} - /// Shared admission control for external descriptor metadata and range reads. /// /// The byte semaphore budgets active range I/O only. A single range larger than @@ -381,7 +371,8 @@ impl BlobReadLimiter { .await .map_err(|e| crate::Error::UnexpectedError { message: format!( - "Failed to acquire BlobDescriptor byte permits for URI '{uri}': {e}" + "Failed to acquire BlobDescriptor byte permits for URI '{}': {e}", + sanitize_blob_uri(uri) ), source: Some(Box::new(e)), })?; @@ -395,7 +386,8 @@ impl BlobReadLimiter { .await .map_err(|e| crate::Error::UnexpectedError { message: format!( - "Failed to acquire BlobDescriptor {operation} permit for URI '{uri}': {e}" + "Failed to acquire BlobDescriptor {operation} permit for URI '{}': {e}", + sanitize_blob_uri(uri) ), source: Some(Box::new(e)), }) @@ -480,7 +472,10 @@ async fn resolve_column( .size() .await .map_err(|e| crate::Error::UnexpectedError { - message: format!("Failed to read metadata for BlobDescriptor URI '{uri}': {e}"), + message: format!( + "Failed to read metadata for BlobDescriptor URI '{}': {e}", + sanitize_blob_uri(&uri) + ), source: Some(Box::new(e)), })? } else { @@ -606,7 +601,7 @@ async fn read_merged_blob_ranges( let blob_parallelism = limiter.parallelism(); stream::iter(reads) .map(|merged| { - let uri = uri.to_string(); + let uri = sanitize_blob_uri(uri); let reader = reader.clone(); let limiter = limiter.clone(); async move { diff --git a/crates/paimon/src/table/data_evolution_writer.rs b/crates/paimon/src/table/data_evolution_writer.rs index cbe0b142d..7976bbf34 100644 --- a/crates/paimon/src/table/data_evolution_writer.rs +++ b/crates/paimon/src/table/data_evolution_writer.rs @@ -319,7 +319,13 @@ impl DataEvolutionWriter { if self.matched_batches.is_empty() { return Ok(Vec::new()); } - let read_table = &index.read_table; + // Rewriting an inline Blob column must preserve references for rows + // without an update. Resolving them would replace descriptors with + // payload bytes and unnecessarily require the external files to exist. + let read_table = index.read_table.copy_with_options(HashMap::from([( + "blob-as-descriptor".to_string(), + "true".to_string(), + )])); let file_index = &index.files; if file_index.is_empty() { return Err(crate::Error::DataInvalid { diff --git a/crates/paimon/src/table/table_write.rs b/crates/paimon/src/table/table_write.rs index b2d8085b0..22a423566 100644 --- a/crates/paimon/src/table/table_write.rs +++ b/crates/paimon/src/table/table_write.rs @@ -615,6 +615,9 @@ impl TableWrite { pub(super) fn normalize_write_batch(&self, batch: &RecordBatch) -> Result> { let batch = self.validate_write_batch_schema(batch)?; + // Java filters row kinds before extracting Blob values. Ignored rows + // need neither valid descriptor bytes nor physical files. + let batch = self.enrich_rowkind_batch(&batch)?; if batch.num_rows() == 0 { return Ok(None); } @@ -624,8 +627,7 @@ impl TableWrite { self.table.schema().options(), )?; } - let batch = self.enrich_rowkind_batch(&batch)?; - Ok((batch.num_rows() != 0).then_some(batch)) + Ok(Some(batch)) } pub(super) async fn write_partition_bucket_batch( diff --git a/crates/paimon/tests/blob_http_error_test.rs b/crates/paimon/tests/blob_http_error_test.rs new file mode 100644 index 000000000..0fced0004 --- /dev/null +++ b/crates/paimon/tests/blob_http_error_test.rs @@ -0,0 +1,133 @@ +// 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. + +mod common; + +use arrow_array::{ArrayRef, Int32Array, LargeBinaryArray, RecordBatch}; +use axum::{http::StatusCode, response::Redirect, routing::get, Router}; +use common::incremental_helpers::{memory_table, persist_table_schema, setup_dirs}; +use futures::TryStreamExt; +use paimon::spec::{BlobDescriptor, BlobType, DataType, IntType, Schema, TableSchema}; +use paimon::table::BlobReader; +use std::sync::Arc; + +fn assert_redacted(error: &paimon::Error) { + let mut rendered = format!("{error}\n{error:?}"); + let mut source = std::error::Error::source(error); + while let Some(cause) = source { + rendered.push_str(&format!("\n{cause}\n{cause:?}")); + source = cause.source(); + } + for secret in [ + "original-user", + "original-password", + "original-token", + "original-fragment", + "redirect-secret", + ] { + assert!( + !rendered.contains(secret), + "HTTP credential leaked: {rendered}" + ); + } +} + +#[tokio::test] +async fn http_errors_redact_direct_and_redirected_credentials_in_public_blob_apis() { + let app = Router::new() + .route( + "/redirect", + get(|| async { Redirect::temporary("/failure?signature=redirect-secret") }), + ) + .route("/failure", get(|| async { StatusCode::FORBIDDEN })) + .route("/short", get(|| async { "x" })); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let server = tokio::spawn(async move { axum::serve(listener, app).await.unwrap() }); + for route in ["failure", "redirect", "short"] { + for length in [-1, 2] { + // A one-byte resource only fails a bounded read, not length=-1. + if route == "short" && length == -1 { + continue; + } + for inline in [false, true] { + let mut schema = Schema::builder() + .column("id", DataType::Int(IntType::new())) + .column("payload", DataType::Blob(BlobType::new())) + .option("row-tracking.enabled", "true") + .option("data-evolution.enabled", "true"); + if inline { + schema = schema.option("blob-descriptor-field", "payload"); + } + let path = "memory:/http_failure"; + let (io, table) = memory_table(path, TableSchema::new(0, &schema.build().unwrap())); + setup_dirs(&io, path).await; + persist_table_schema(&io, path, table.schema()).await; + let uri = if route == "redirect" { + // Credentials exist only on the redirected URL. + format!("http://{address}/redirect") + } else { + format!("http://original-user:original-password@{address}/{route}?signature=original-token#original-fragment") + }; + let descriptor = BlobDescriptor::new(uri, 0, length).serialize(); + let error = BlobReader::from_file_io(io.clone()) + .read_blobs(std::slice::from_ref(&descriptor)) + .await + .unwrap_err(); + assert_redacted(&error); + if route != "short" { + assert!(error.to_string().contains("403")); + } + let input = RecordBatch::try_from_iter([ + ("id", Arc::new(Int32Array::from(vec![1])) as ArrayRef), + ( + "payload", + Arc::new(LargeBinaryArray::from(vec![Some(descriptor.as_slice())])) + as ArrayRef, + ), + ]) + .unwrap(); + let builder = table.new_write_builder(); + let mut writer = builder.new_write().unwrap(); + let error = if inline { + writer.write_arrow_batch(&input).await.unwrap(); + builder + .new_commit() + .commit(writer.prepare_commit().await.unwrap()) + .await + .unwrap(); + let read = table.new_read_builder(); + let plan = read.new_scan().plan().await.unwrap(); + read.new_read() + .unwrap() + .to_arrow(plan.splits()) + .unwrap() + .try_collect::>() + .await + .unwrap_err() + } else { + writer.write_arrow_batch(&input).await.unwrap_err() + }; + assert_redacted(&error); + if route != "short" { + assert!(error.to_string().contains("403")); + } + } + } + } + server.abort(); +} diff --git a/crates/paimon/tests/dedicated_blob_write_test.rs b/crates/paimon/tests/dedicated_blob_write_test.rs index 3a7548ab3..cdcbf5d18 100644 --- a/crates/paimon/tests/dedicated_blob_write_test.rs +++ b/crates/paimon/tests/dedicated_blob_write_test.rs @@ -85,6 +85,64 @@ async fn rows(table: &Table) -> Vec<(i32, Option>)> { rows } +#[tokio::test] +async fn partial_blob_updates_preserve_unmatched_descriptors_without_resolving() { + let table = table(&[("blob-descriptor-field", "payload")]).await; + // None of these references exist. Updating an inline reference must not + // read the old or new payloads, including rows retained in the same file. + let old = BlobDescriptor::new("memory:/missing/old".into(), 0, 3).serialize(); + let retained = BlobDescriptor::new("memory:/missing/retained".into(), 7, -1).serialize(); + let updated = BlobDescriptor::new("memory:/missing/updated".into(), 2, 5).serialize(); + let builder = table.new_write_builder(); + let mut writer = builder.new_write().unwrap(); + writer + .write_arrow_batch(&batch( + vec![1, 2, 3], + vec![Some(&old), Some(&retained), None], + )) + .await + .unwrap(); + builder + .new_commit() + .commit(writer.prepare_commit().await.unwrap()) + .await + .unwrap(); + + let mut update = builder + .new_data_evolution_writer(vec!["payload".into()]) + .unwrap(); + update + .add_matched_batch( + RecordBatch::try_from_iter([ + ( + "_ROW_ID", + Arc::new(arrow_array::Int64Array::from(vec![0])) as ArrayRef, + ), + ( + "payload", + Arc::new(LargeBinaryArray::from(vec![Some(updated.as_slice())])) as ArrayRef, + ), + ]) + .unwrap(), + ) + .unwrap(); + builder + .new_commit() + .commit(update.prepare_commit().await.unwrap()) + .await + .unwrap(); + + let descriptor_table = table.copy_with_options( + [("blob-as-descriptor".to_string(), "true".to_string())] + .into_iter() + .collect(), + ); + assert_eq!( + rows(&descriptor_table).await, + vec![(1, Some(updated)), (2, Some(retained)), (3, None)] + ); +} + #[tokio::test] async fn prepared_sequences_and_optimized_write_columns_match_java() { for optimize in [false, true] { @@ -419,3 +477,80 @@ async fn primary_key_inline_descriptors_validate_and_resolve_legacy_bytes() { ] ); } + +#[tokio::test] +async fn ignored_row_kinds_skip_inline_blob_validation() { + for generated in [false, true] { + let mut schema = Schema::builder() + .column("id", DataType::Int(IntType::with_nullable(false))) + .column("payload", DataType::Blob(BlobType::new())) + .primary_key(["id"]) + .option("bucket", "1") + .option("ignore-delete", "true") + .option("ignore-update-before", "true") + .option("blob-descriptor-field", "payload"); + if generated { + schema = schema + .column( + "kind", + DataType::VarChar(paimon::spec::VarCharType::string_type()), + ) + .option("rowkind.field", "kind"); + } + let path = "memory:/ignored_inline_blob"; + let (io, table) = memory_table(path, TableSchema::new(0, &schema.build().unwrap())); + setup_dirs(&io, path).await; + persist_table_schema(&io, path, table.schema()).await; + io.new_output("memory:/payload") + .unwrap() + .write(b"hello".to_vec().into()) + .await + .unwrap(); + let descriptor = BlobDescriptor::new("memory:/payload".into(), 0, 5).serialize(); + let input = batch( + vec![1, 2, 3], + vec![Some(&descriptor), Some(b""), Some(b"raw")], + ); + let mut columns: Vec<(String, ArrayRef)> = input + .schema() + .fields() + .iter() + .zip(input.columns()) + .map(|(field, column)| (field.name().clone(), column.clone())) + .collect(); + if generated { + columns.push(( + "kind".into(), + Arc::new(arrow_array::StringArray::from(vec!["+I", "-D", "-U"])), + )); + } else { + columns.push(( + "_VALUE_KIND".into(), + Arc::new(arrow_array::Int8Array::from(vec![0, 3, 1])), + )); + } + let input = RecordBatch::try_from_iter(columns).unwrap(); + let builder = table.new_write_builder(); + let mut writer = builder.new_write().unwrap(); + // An all-ignored batch creates no files, even with invalid payloads. + writer.write_arrow_batch(&input.slice(1, 2)).await.unwrap(); + assert!(writer.prepare_commit().await.unwrap().is_empty()); + writer.write_arrow_batch(&input).await.unwrap(); + let messages = writer.prepare_commit().await.unwrap(); + builder.new_commit().commit(messages).await.unwrap(); + assert_eq!(rows(&table).await, vec![(1, Some(b"hello".to_vec()))]); + // Filtering must not bypass validation of a retained row. + let invalid = input.slice(0, 1); + let mut columns = invalid.columns().to_vec(); + columns[1] = Arc::new(LargeBinaryArray::from(vec![Some(b"raw".as_slice())])); + let invalid = RecordBatch::try_new(invalid.schema(), columns).unwrap(); + assert!(builder + .new_write() + .unwrap() + .write_arrow_batch(&invalid) + .await + .unwrap_err() + .to_string() + .contains("blob-descriptor-field")); + } +}