Skip to content
Draft
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
252 changes: 252 additions & 0 deletions bindings/c/src/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,42 @@ fn not_null_table_schema() -> TableSchema {
TableSchema::new(0, &schema)
}

fn postpone_table_schema() -> TableSchema {
let schema = Schema::builder()
.column("id", DataType::Int(IntType::new()))
.column("name", DataType::VarChar(VarCharType::string_type()))
.primary_key(["id"])
.option("bucket", "-2")
.build()
.unwrap();
TableSchema::new(0, &schema)
}

fn legacy_postpone_table_schema() -> TableSchema {
let schema = Schema::builder()
.column("id", DataType::Int(IntType::new()))
.column("name", DataType::VarChar(VarCharType::string_type()))
.primary_key(["id"])
.option("bucket", "-2")
.option("postpone.batch-write-fixed-bucket", "false")
.build()
.unwrap();
TableSchema::new(0, &schema)
}

fn partitioned_postpone_table_schema() -> TableSchema {
let schema = Schema::builder()
.column("pt", DataType::VarChar(VarCharType::string_type()))
.column("id", DataType::Int(IntType::new()))
.column("name", DataType::VarChar(VarCharType::string_type()))
.primary_key(["pt", "id"])
.partition_keys(["pt"])
.option("bucket", "-2")
.build()
.unwrap();
TableSchema::new(0, &schema)
}

unsafe fn wrap_table(table: Table) -> *mut paimon_table {
let inner = Box::into_raw(Box::new(table)) as *mut c_void;
Box::into_raw(Box::new(paimon_table { inner }))
Expand Down Expand Up @@ -104,6 +140,38 @@ fn make_batch(ids: Vec<i32>, names: Vec<&str>) -> RecordBatch {
.unwrap()
}

fn make_partitioned_write_batch(pts: Vec<&str>, ids: Vec<i32>, names: Vec<&str>) -> RecordBatch {
let schema = Arc::new(ArrowSchema::new(vec![
ArrowField::new("pt", ArrowDataType::Utf8, false),
ArrowField::new("id", ArrowDataType::Int32, false),
ArrowField::new("name", ArrowDataType::Utf8, true),
]));
RecordBatch::try_new(
schema,
vec![
Arc::new(StringArray::from(pts)),
Arc::new(Int32Array::from(ids)),
Arc::new(StringArray::from(names)),
],
)
.unwrap()
}

fn make_postpone_bucket_plan_batch(partitions: Vec<&str>, counts: Vec<i32>) -> RecordBatch {
let schema = Arc::new(ArrowSchema::new(vec![
ArrowField::new("pt", ArrowDataType::Utf8, false),
ArrowField::new("total_buckets", ArrowDataType::Int32, false),
]));
RecordBatch::try_new(
schema,
vec![
Arc::new(StringArray::from(partitions)),
Arc::new(Int32Array::from(counts)),
],
)
.unwrap()
}

fn make_type_mismatch_batch(ids: Vec<&str>, names: Vec<&str>) -> RecordBatch {
let schema = Arc::new(ArrowSchema::new(vec![
ArrowField::new("id", ArrowDataType::Utf8, false),
Expand Down Expand Up @@ -950,6 +1018,93 @@ fn test_write_new_builder_and_free() {
}
}

#[test]
fn test_postpone_fixed_bucket_builder_respects_option() {
let path = "memory:/test_postpone_fixed_bucket_builder_option";
let file_io = memory_file_io();
setup_table_dirs(&file_io, path);
let table = Table::new(
file_io,
Identifier::new("default", "test"),
path.to_string(),
postpone_table_schema(),
None,
);
let handle = unsafe { wrap_table(table) };

unsafe {
let normal = paimon_table_new_write_builder(handle);
assert!(normal.error.is_null());
let normal_state = &*((*normal.write_builder).inner as *const WriteBuilderState);
assert!(normal_state.postpone_fixed_bucket);

let write = paimon_write_builder_new_write(normal.write_builder);
assert!(write.error.is_null());
let (array, schema) = export_batch_to_ffi(make_batch(vec![1], vec!["a"]));
let error = paimon_table_write_write_arrow_batch(
write.write,
(&**array) as *const FFI_ArrowArray as *mut c_void,
(&**schema) as *const FFI_ArrowSchema as *mut c_void,
);
assert!(error.is_null());
let prepared = paimon_table_write_prepare_commit(write.write);
assert!(prepared.error.is_null());
let messages = &*((*prepared.messages).inner as *const CommitMessagesState);
assert!(messages.messages.iter().all(|message| message.bucket >= 0));
assert!(messages
.messages
.iter()
.all(|message| message.total_buckets == Some(1)));
paimon_commit_messages_free(prepared.messages);
paimon_table_write_free(write.write);
paimon_write_builder_free(normal.write_builder);

let fixed = paimon_table_new_postpone_fixed_bucket_write_builder(handle);
assert!(fixed.error.is_null());
let fixed_state = &*((*fixed.write_builder).inner as *const WriteBuilderState);
assert!(fixed_state.postpone_fixed_bucket);
paimon_write_builder_free(fixed.write_builder);

let commit_user = CString::new("fixed-user").unwrap();
let fixed = paimon_table_new_postpone_fixed_bucket_write_builder_with_commit_user(
handle,
commit_user.as_ptr(),
);
assert!(fixed.error.is_null());
let fixed_state = &*((*fixed.write_builder).inner as *const WriteBuilderState);
assert!(fixed_state.postpone_fixed_bucket);
assert_eq!(fixed_state.commit_user, "fixed-user");
paimon_write_builder_free(fixed.write_builder);
unwrap_table(handle);
}

let legacy_path = "memory:/test_legacy_postpone_builder_option";
let file_io = memory_file_io();
setup_table_dirs(&file_io, legacy_path);
let legacy_table = Table::new(
file_io,
Identifier::new("default", "test"),
legacy_path.to_string(),
legacy_postpone_table_schema(),
None,
);
let legacy_handle = unsafe { wrap_table(legacy_table) };
unsafe {
let normal = paimon_table_new_write_builder(legacy_handle);
assert!(normal.error.is_null());
let normal_state = &*((*normal.write_builder).inner as *const WriteBuilderState);
assert!(!normal_state.postpone_fixed_bucket);
paimon_write_builder_free(normal.write_builder);

let fixed = paimon_table_new_postpone_fixed_bucket_write_builder(legacy_handle);
assert!(fixed.error.is_null());
let fixed_state = &*((*fixed.write_builder).inner as *const WriteBuilderState);
assert!(fixed_state.postpone_fixed_bucket);
paimon_write_builder_free(fixed.write_builder);
unwrap_table(legacy_handle);
}
}

#[test]
fn test_write_commit_read_roundtrip() {
let path = "memory:/test_write_roundtrip";
Expand Down Expand Up @@ -1442,6 +1597,90 @@ fn test_commit_messages_merge_preserves_all_writer_files() {
}
}

#[test]
fn test_distributed_postpone_writers_share_bucket_plan() {
let path = "memory:/test_distributed_postpone_bucket_plan";
let file_io = memory_file_io();
setup_table_dirs(&file_io, path);
let table = Table::new(
file_io,
Identifier::new("default", "test"),
path.to_string(),
partitioned_postpone_table_schema(),
None,
);
let handle = unsafe { wrap_table(table) };
let commit_user = CString::new("distributed-postpone-job").unwrap();

unsafe {
let wb1 = paimon_table_new_postpone_fixed_bucket_write_builder_with_commit_user(
handle,
commit_user.as_ptr(),
)
.write_builder;
let wb2 = paimon_table_new_postpone_fixed_bucket_write_builder_with_commit_user(
handle,
commit_user.as_ptr(),
)
.write_builder;

for wb in [wb1, wb2] {
let (array, schema) =
export_batch_to_ffi(make_postpone_bucket_plan_batch(vec!["p"], vec![3]));
let error = paimon_write_builder_with_postpone_bucket_plan(
wb,
(&**array) as *const FFI_ArrowArray as *mut c_void,
(&**schema) as *const FFI_ArrowSchema as *mut c_void,
);
assert!(error.is_null());
}

let tw1 = paimon_write_builder_new_write(wb1).write;
let tw2 = paimon_write_builder_new_write(wb2).write;
for (tw, ids, names) in [
(tw1, vec![1], vec!["a"]),
(tw2, vec![2, 3, 4, 5], vec!["b", "c", "d", "e"]),
] {
let (array, schema) = export_batch_to_ffi(make_partitioned_write_batch(
vec!["p"; ids.len()],
ids,
names,
));
let error = paimon_table_write_write_arrow_batch(
tw,
(&**array) as *const FFI_ArrowArray as *mut c_void,
(&**schema) as *const FFI_ArrowSchema as *mut c_void,
);
assert!(error.is_null());
}

let messages1 = paimon_table_write_prepare_commit(tw1).messages;
let messages2 = paimon_table_write_prepare_commit(tw2).messages;
for messages in [messages1, messages2] {
let state = &*((*messages).inner as *const CommitMessagesState);
assert!(!state.messages.is_empty());
assert!(state
.messages
.iter()
.all(|message| message.total_buckets == Some(3)));
}
let error = paimon_commit_messages_merge(messages1, messages2);
assert!(error.is_null());
let commit = paimon_write_builder_new_commit(wb1).commit;
let error = paimon_table_commit_commit(commit, messages1);
assert!(error.is_null());

paimon_table_commit_free(commit);
paimon_commit_messages_free(messages2);
paimon_commit_messages_free(messages1);
paimon_table_write_free(tw2);
paimon_table_write_free(tw1);
paimon_write_builder_free(wb2);
paimon_write_builder_free(wb1);
unwrap_table(handle);
}
}

#[test]
fn test_write_multiple_batches() {
let path = "memory:/test_write_multi_batch";
Expand Down Expand Up @@ -1721,6 +1960,11 @@ fn test_null_pointer_handling() {
assert!(result.write_builder.is_null());
paimon_error_free(result.error);

let result = paimon_table_new_postpone_fixed_bucket_write_builder(ptr::null());
assert!(!result.error.is_null());
assert!(result.write_builder.is_null());
paimon_error_free(result.error);

let result = paimon_write_builder_new_write(ptr::null());
assert!(!result.error.is_null());
assert!(result.write.is_null());
Expand All @@ -1731,6 +1975,14 @@ fn test_null_pointer_handling() {
assert!(result.commit.is_null());
paimon_error_free(result.error);

let err = paimon_write_builder_with_postpone_bucket_plan(
ptr::null_mut(),
ptr::null_mut(),
ptr::null_mut(),
);
assert!(!err.is_null());
paimon_error_free(err);

let err =
paimon_table_write_write_arrow_batch(ptr::null_mut(), ptr::null_mut(), ptr::null_mut());
assert!(!err.is_null());
Expand Down
4 changes: 3 additions & 1 deletion bindings/c/src/types.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@ use std::sync::Arc;

use arrow_schema::Schema as ArrowSchema;
use paimon::spec::{DataField, Predicate};
use paimon::table::{CommitMessage, Table, TableCommit, TableWrite};
use paimon::table::{CommitMessage, PostponeBucketPlan, Table, TableCommit, TableWrite};

/// C-compatible key-value pair for options.
#[repr(C)]
Expand Down Expand Up @@ -210,6 +210,8 @@ pub(crate) struct WriteBuilderState {
pub table: Table,
pub commit_user: String,
pub overwrite: bool,
pub postpone_fixed_bucket: bool,
pub postpone_bucket_plan: Option<PostponeBucketPlan>,
}

pub(crate) struct TableWriteState {
Expand Down
Loading
Loading