Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

3 changes: 3 additions & 0 deletions bindings/c/DEPENDENCIES.rust.tsv

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

3 changes: 3 additions & 0 deletions bindings/go/DEPENDENCIES.rust.tsv

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

3 changes: 3 additions & 0 deletions crates/integration_tests/DEPENDENCIES.rust.tsv

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

88 changes: 66 additions & 22 deletions crates/integrations/datafusion/tests/blob_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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, \
Expand All @@ -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;

Expand Down Expand Up @@ -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, \
Expand All @@ -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;

Expand Down Expand Up @@ -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, \
Expand All @@ -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;

Expand Down Expand Up @@ -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, \
Expand All @@ -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;

Expand Down Expand Up @@ -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");
Expand All @@ -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;

Expand All @@ -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())),
]
);
}
Expand All @@ -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();
Expand All @@ -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;
Expand All @@ -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),
]
);
Expand Down Expand Up @@ -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(
Expand Down
3 changes: 3 additions & 0 deletions crates/paimon-rest-server/DEPENDENCIES.rust.tsv

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion crates/paimon/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
3 changes: 3 additions & 0 deletions crates/paimon/DEPENDENCIES.rust.tsv

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

13 changes: 6 additions & 7 deletions crates/paimon/src/arrow/format/blob.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2151,21 +2151,20 @@ 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!(
"Failed to read metadata for BlobDescriptor '{}': {e}",
desc.uri()
crate::io::uri_reader::sanitize_blob_uri(desc.uri())
),
source: Some(Box::new(e)),
})?
.size
.saturating_sub(offset),
};
let end = offset
Expand All @@ -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();
Expand All @@ -2210,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)),
}
Expand All @@ -2221,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,
});
Expand Down
1 change: 1 addition & 0 deletions crates/paimon/src/io/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@

pub(crate) mod cache;
mod file_io;
pub(crate) mod uri_reader;
pub use file_io::*;

mod storage;
Expand Down
Loading
Loading