diff --git a/bindings/go/postpone_fixed_bucket_write.go b/bindings/go/postpone_fixed_bucket_write.go new file mode 100644 index 000000000..c037ad8a1 --- /dev/null +++ b/bindings/go/postpone_fixed_bucket_write.go @@ -0,0 +1,329 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package paimon + +import ( + "context" + "sync" + "unsafe" + + "github.com/apache/arrow-go/v18/arrow" +) + +// PostponeFixedBucketWriteBuilder creates fixed-bucket writers for bucket=-2 +// tables. A resolved bucket plan must be supplied before NewWrite. +type PostponeFixedBucketWriteBuilder struct { + ctx context.Context + lib *libRef + inner *paimonPostponeFixedBucketWriteBuilder + closeOnce sync.Once +} + +// NewPostponeFixedBucketWriteBuilder creates an explicitly selected +// fixed-bucket builder for a postpone table. +func (t *Table) NewPostponeFixedBucketWriteBuilder() (*PostponeFixedBucketWriteBuilder, error) { + if t.inner == nil { + return nil, ErrClosed + } + inner, err := ffiTableNewPostponeFixedBucketWriteBuilder.symbol(t.ctx)(t.inner) + if err != nil { + return nil, err + } + t.lib.acquire() + return &PostponeFixedBucketWriteBuilder{ctx: t.ctx, lib: t.lib, inner: inner}, nil +} + +// NewPostponeFixedBucketWriteBuilderWithCommitUser creates a fixed-bucket +// builder with a stable commit identity. +func (t *Table) NewPostponeFixedBucketWriteBuilderWithCommitUser( + commitUser string, +) (*PostponeFixedBucketWriteBuilder, error) { + if t.inner == nil { + return nil, ErrClosed + } + inner, err := ffiTableNewPostponeFixedBucketWriteBuilderWithCommitUser.symbol(t.ctx)( + t.inner, + commitUser, + ) + if err != nil { + return nil, err + } + t.lib.acquire() + return &PostponeFixedBucketWriteBuilder{ctx: t.ctx, lib: t.lib, inner: inner}, nil +} + +// Close releases the builder resources. Safe to call multiple times. +func (wb *PostponeFixedBucketWriteBuilder) Close() { + wb.closeOnce.Do(func() { + ffiPostponeFixedBucketWriteBuilderFree.symbol(wb.ctx)(wb.inner) + wb.inner = nil + wb.lib.release() + }) +} + +// WithOverwrite enables overwrite mode for both writers and committers created +// by this builder. +func (wb *PostponeFixedBucketWriteBuilder) WithOverwrite() error { + if wb.inner == nil { + return ErrClosed + } + return ffiPostponeFixedBucketWriteBuilderWithOverwrite.symbol(wb.ctx)(wb.inner) +} + +// WithBucketPlan sets a resolved partition-to-bucket-count plan. The plan must +// contain the table partition columns followed by a non-null Int32 +// total_buckets column. The caller retains ownership of plan. +func (wb *PostponeFixedBucketWriteBuilder) WithBucketPlan(plan arrow.Record) error { + if wb.inner == nil { + return ErrClosed + } + return withOwnedArrowRecord( + plan, + "paimon: bucket plan must not be nil", + func(array, schema unsafe.Pointer) error { + return ffiPostponeFixedBucketWriteBuilderWithBucketPlan.symbol(wb.ctx)( + wb.inner, + array, + schema, + ) + }, + ) +} + +// NewWrite creates a fixed-bucket writer. WithBucketPlan must be called first. +func (wb *PostponeFixedBucketWriteBuilder) NewWrite() (*PostponeFixedBucketTableWrite, error) { + if wb.inner == nil { + return nil, ErrClosed + } + inner, err := ffiPostponeFixedBucketWriteBuilderNewWrite.symbol(wb.ctx)(wb.inner) + if err != nil { + return nil, err + } + wb.lib.acquire() + return &PostponeFixedBucketTableWrite{ctx: wb.ctx, lib: wb.lib, inner: inner}, nil +} + +// NewCommit creates a committer using this builder's commit identity and mode. +func (wb *PostponeFixedBucketWriteBuilder) NewCommit() (*PostponeFixedBucketTableCommit, error) { + if wb.inner == nil { + return nil, ErrClosed + } + inner, err := ffiPostponeFixedBucketWriteBuilderNewCommit.symbol(wb.ctx)(wb.inner) + if err != nil { + return nil, err + } + wb.lib.acquire() + return &PostponeFixedBucketTableCommit{ctx: wb.ctx, lib: wb.lib, inner: inner}, nil +} + +// PostponeFixedBucketTableWrite writes rows according to a resolved bucket plan. +type PostponeFixedBucketTableWrite struct { + ctx context.Context + lib *libRef + inner *paimonPostponeFixedBucketTableWrite + closeOnce sync.Once +} + +// Close releases the writer resources. Safe to call multiple times. +func (tw *PostponeFixedBucketTableWrite) Close() { + tw.closeOnce.Do(func() { + ffiPostponeFixedBucketTableWriteFree.symbol(tw.ctx)(tw.inner) + tw.inner = nil + tw.lib.release() + }) +} + +// WriteArrowBatch writes one Arrow record batch. The caller retains ownership. +func (tw *PostponeFixedBucketTableWrite) WriteArrowBatch(record arrow.Record) error { + if tw.inner == nil { + return ErrClosed + } + return withOwnedArrowRecord( + record, + "paimon: record batch must not be nil", + func(array, schema unsafe.Pointer) error { + return ffiPostponeFixedBucketTableWriteWriteArrowBatch.symbol(tw.ctx)( + tw.inner, + array, + schema, + ) + }, + ) +} + +// PrepareCommit closes current writers and returns fixed-bucket messages. +// The writer is single-use; create a new writer for the next batch. +func (tw *PostponeFixedBucketTableWrite) PrepareCommit() (*PostponeFixedBucketCommitMessages, error) { + if tw.inner == nil { + return nil, ErrClosed + } + inner, err := ffiPostponeFixedBucketTableWritePrepareCommit.symbol(tw.ctx)(tw.inner) + if err != nil { + return nil, err + } + tw.lib.acquire() + return &PostponeFixedBucketCommitMessages{ctx: tw.ctx, lib: tw.lib, inner: inner}, nil +} + +// PostponeFixedBucketCommitMessages contains files produced by fixed-bucket +// writers. It is a process-local native handle and cannot be transferred +// between processes or passed to a standard TableCommit. +type PostponeFixedBucketCommitMessages struct { + ctx context.Context + lib *libRef + inner *paimonPostponeFixedBucketCommitMessages + closeOnce sync.Once +} + +// Close releases the messages. Safe to call multiple times. +func (m *PostponeFixedBucketCommitMessages) Close() { + m.closeOnce.Do(func() { + ffiPostponeFixedBucketCommitMessagesFree.symbol(m.ctx)(m.inner) + m.inner = nil + m.lib.release() + }) +} + +// Merge appends a copy of source's messages. Both handles must belong to the +// same process, and both builders must use the same table, commit user, and +// overwrite mode. +func (m *PostponeFixedBucketCommitMessages) Merge( + source *PostponeFixedBucketCommitMessages, +) error { + if m.inner == nil { + return ErrClosed + } + if source == nil || source.inner == nil { + return ErrClosed + } + return ffiPostponeFixedBucketCommitMessagesMerge.symbol(m.ctx)(m.inner, source.inner) +} + +// PostponeFixedBucketTableCommit commits fixed-bucket messages using the mode +// selected on its builder. +type PostponeFixedBucketTableCommit struct { + ctx context.Context + lib *libRef + inner *paimonPostponeFixedBucketTableCommit + closeOnce sync.Once +} + +// Close releases the committer resources. Safe to call multiple times. +func (tc *PostponeFixedBucketTableCommit) Close() { + tc.closeOnce.Do(func() { + ffiPostponeFixedBucketTableCommitFree.symbol(tc.ctx)(tc.inner) + tc.inner = nil + tc.lib.release() + }) +} + +func (tc *PostponeFixedBucketTableCommit) withMessages( + messages *PostponeFixedBucketCommitMessages, + operation func( + *paimonPostponeFixedBucketTableCommit, + *paimonPostponeFixedBucketCommitMessages, + ) error, +) error { + if tc.inner == nil { + return ErrClosed + } + if messages == nil || messages.inner == nil { + return ErrClosed + } + return operation(tc.inner, messages.inner) +} + +func (tc *PostponeFixedBucketTableCommit) withMessagesAndIdentifier( + messages *PostponeFixedBucketCommitMessages, + commitIdentifier int64, + operation func( + *paimonPostponeFixedBucketTableCommit, + *paimonPostponeFixedBucketCommitMessages, + int64, + ) error, +) error { + if tc.inner == nil { + return ErrClosed + } + if messages == nil || messages.inner == nil { + return ErrClosed + } + return operation(tc.inner, messages.inner, commitIdentifier) +} + +// Commit persists fixed-bucket messages using the builder's append or overwrite +// mode. +func (tc *PostponeFixedBucketTableCommit) Commit( + messages *PostponeFixedBucketCommitMessages, +) error { + return tc.withMessages(messages, ffiPostponeFixedBucketTableCommitCommit.symbol(tc.ctx)) +} + +// CommitWithIdentifier commits with a stable identifier. +func (tc *PostponeFixedBucketTableCommit) CommitWithIdentifier( + messages *PostponeFixedBucketCommitMessages, + commitIdentifier int64, +) error { + return tc.withMessagesAndIdentifier( + messages, + commitIdentifier, + ffiPostponeFixedBucketTableCommitCommitWithIdentifier.symbol(tc.ctx), + ) +} + +// FilterAndCommitWithIdentifier makes a retry idempotent. +func (tc *PostponeFixedBucketTableCommit) FilterAndCommitWithIdentifier( + messages *PostponeFixedBucketCommitMessages, + commitIdentifier int64, +) error { + return tc.withMessagesAndIdentifier( + messages, + commitIdentifier, + ffiPostponeFixedBucketTableCommitFilterAndCommitWithIdentifier.symbol(tc.ctx), + ) +} + +// TruncateTable removes all table data. +func (tc *PostponeFixedBucketTableCommit) TruncateTable() error { + if tc.inner == nil { + return ErrClosed + } + return ffiPostponeFixedBucketTableCommitTruncateTable.symbol(tc.ctx)(tc.inner) +} + +// TruncateTableWithIdentifier removes all table data with a stable identifier. +func (tc *PostponeFixedBucketTableCommit) TruncateTableWithIdentifier( + commitIdentifier int64, +) error { + if tc.inner == nil { + return ErrClosed + } + return ffiPostponeFixedBucketTableCommitTruncateTableWithIdentifier.symbol(tc.ctx)( + tc.inner, + commitIdentifier, + ) +} + +// Abort performs best-effort cleanup of files in prepared messages. +func (tc *PostponeFixedBucketTableCommit) Abort( + messages *PostponeFixedBucketCommitMessages, +) error { + return tc.withMessages(messages, ffiPostponeFixedBucketTableCommitAbort.symbol(tc.ctx)) +} diff --git a/bindings/go/postpone_fixed_bucket_write_ffi.go b/bindings/go/postpone_fixed_bucket_write_ffi.go new file mode 100644 index 000000000..d9154e587 --- /dev/null +++ b/bindings/go/postpone_fixed_bucket_write_ffi.go @@ -0,0 +1,394 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package paimon + +import ( + "context" + "unsafe" + + "github.com/jupiterrider/ffi" +) + +var ffiTableNewPostponeFixedBucketWriteBuilder = newFFI(ffiOpts{ + sym: "paimon_table_new_postpone_fixed_bucket_write_builder", + rType: &typeResultWriteBuilder, + aTypes: []*ffi.Type{&ffi.TypePointer}, +}, func( + ctx context.Context, + ffiCall ffiCall, +) func(*paimonTable) (*paimonPostponeFixedBucketWriteBuilder, error) { + return func(table *paimonTable) (*paimonPostponeFixedBucketWriteBuilder, error) { + var result resultPostponeFixedBucketWriteBuilder + ffiCall(unsafe.Pointer(&result), unsafe.Pointer(&table)) + if result.error != nil { + return nil, parseError(ctx, result.error) + } + return result.writeBuilder, nil + } +}) + +var ffiTableNewPostponeFixedBucketWriteBuilderWithCommitUser = newFFI(ffiOpts{ + sym: "paimon_table_new_postpone_fixed_bucket_write_builder_with_commit_user", + rType: &typeResultWriteBuilder, + aTypes: []*ffi.Type{&ffi.TypePointer, &ffi.TypePointer}, +}, func( + ctx context.Context, + ffiCall ffiCall, +) func(*paimonTable, string) (*paimonPostponeFixedBucketWriteBuilder, error) { + return func( + table *paimonTable, + commitUser string, + ) (*paimonPostponeFixedBucketWriteBuilder, error) { + commitUserPtr, err := bytePtrFromString(commitUser) + if err != nil { + return nil, err + } + var result resultPostponeFixedBucketWriteBuilder + ffiCall( + unsafe.Pointer(&result), + unsafe.Pointer(&table), + unsafe.Pointer(&commitUserPtr), + ) + if result.error != nil { + return nil, parseError(ctx, result.error) + } + return result.writeBuilder, nil + } +}) + +var ffiPostponeFixedBucketWriteBuilderFree = newFFI(ffiOpts{ + sym: "paimon_postpone_fixed_bucket_write_builder_free", + rType: &ffi.TypeVoid, + aTypes: []*ffi.Type{&ffi.TypePointer}, +}, func( + _ context.Context, + ffiCall ffiCall, +) func(*paimonPostponeFixedBucketWriteBuilder) { + return func(builder *paimonPostponeFixedBucketWriteBuilder) { + ffiCall(nil, unsafe.Pointer(&builder)) + } +}) + +var ffiPostponeFixedBucketWriteBuilderWithOverwrite = newFFI(ffiOpts{ + sym: "paimon_postpone_fixed_bucket_write_builder_with_overwrite", + rType: &ffi.TypePointer, + aTypes: []*ffi.Type{&ffi.TypePointer}, +}, func( + ctx context.Context, + ffiCall ffiCall, +) func(*paimonPostponeFixedBucketWriteBuilder) error { + return func(builder *paimonPostponeFixedBucketWriteBuilder) error { + var ffiError *paimonError + ffiCall(unsafe.Pointer(&ffiError), unsafe.Pointer(&builder)) + return parseError(ctx, ffiError) + } +}) + +var ffiPostponeFixedBucketWriteBuilderWithBucketPlan = newFFI(ffiOpts{ + sym: "paimon_postpone_fixed_bucket_write_builder_with_bucket_plan", + rType: &ffi.TypePointer, + aTypes: []*ffi.Type{ + &ffi.TypePointer, + &ffi.TypePointer, + &ffi.TypePointer, + }, +}, func( + ctx context.Context, + ffiCall ffiCall, +) func(*paimonPostponeFixedBucketWriteBuilder, unsafe.Pointer, unsafe.Pointer) error { + return func( + builder *paimonPostponeFixedBucketWriteBuilder, + array unsafe.Pointer, + schema unsafe.Pointer, + ) error { + var ffiError *paimonError + ffiCall( + unsafe.Pointer(&ffiError), + unsafe.Pointer(&builder), + unsafe.Pointer(&array), + unsafe.Pointer(&schema), + ) + return parseError(ctx, ffiError) + } +}) + +var ffiPostponeFixedBucketWriteBuilderNewWrite = newFFI(ffiOpts{ + sym: "paimon_postpone_fixed_bucket_write_builder_new_write", + rType: &typeResultTableWrite, + aTypes: []*ffi.Type{&ffi.TypePointer}, +}, func( + ctx context.Context, + ffiCall ffiCall, +) func(*paimonPostponeFixedBucketWriteBuilder) (*paimonPostponeFixedBucketTableWrite, error) { + return func( + builder *paimonPostponeFixedBucketWriteBuilder, + ) (*paimonPostponeFixedBucketTableWrite, error) { + var result resultPostponeFixedBucketTableWrite + ffiCall(unsafe.Pointer(&result), unsafe.Pointer(&builder)) + if result.error != nil { + return nil, parseError(ctx, result.error) + } + return result.write, nil + } +}) + +var ffiPostponeFixedBucketWriteBuilderNewCommit = newFFI(ffiOpts{ + sym: "paimon_postpone_fixed_bucket_write_builder_new_commit", + rType: &typeResultTableCommit, + aTypes: []*ffi.Type{&ffi.TypePointer}, +}, func( + ctx context.Context, + ffiCall ffiCall, +) func(*paimonPostponeFixedBucketWriteBuilder) (*paimonPostponeFixedBucketTableCommit, error) { + return func( + builder *paimonPostponeFixedBucketWriteBuilder, + ) (*paimonPostponeFixedBucketTableCommit, error) { + var result resultPostponeFixedBucketTableCommit + ffiCall(unsafe.Pointer(&result), unsafe.Pointer(&builder)) + if result.error != nil { + return nil, parseError(ctx, result.error) + } + return result.commit, nil + } +}) + +var ffiPostponeFixedBucketTableWriteFree = newFFI(ffiOpts{ + sym: "paimon_postpone_fixed_bucket_table_write_free", + rType: &ffi.TypeVoid, + aTypes: []*ffi.Type{&ffi.TypePointer}, +}, func( + _ context.Context, + ffiCall ffiCall, +) func(*paimonPostponeFixedBucketTableWrite) { + return func(write *paimonPostponeFixedBucketTableWrite) { + ffiCall(nil, unsafe.Pointer(&write)) + } +}) + +var ffiPostponeFixedBucketTableWriteWriteArrowBatch = newFFI(ffiOpts{ + sym: "paimon_postpone_fixed_bucket_table_write_write_arrow_batch", + rType: &ffi.TypePointer, + aTypes: []*ffi.Type{ + &ffi.TypePointer, + &ffi.TypePointer, + &ffi.TypePointer, + }, +}, func( + ctx context.Context, + ffiCall ffiCall, +) func(*paimonPostponeFixedBucketTableWrite, unsafe.Pointer, unsafe.Pointer) error { + return func( + write *paimonPostponeFixedBucketTableWrite, + array unsafe.Pointer, + schema unsafe.Pointer, + ) error { + var ffiError *paimonError + ffiCall( + unsafe.Pointer(&ffiError), + unsafe.Pointer(&write), + unsafe.Pointer(&array), + unsafe.Pointer(&schema), + ) + return parseError(ctx, ffiError) + } +}) + +var ffiPostponeFixedBucketTableWritePrepareCommit = newFFI(ffiOpts{ + sym: "paimon_postpone_fixed_bucket_table_write_prepare_commit", + rType: &typeResultPrepareCommit, + aTypes: []*ffi.Type{&ffi.TypePointer}, +}, func( + ctx context.Context, + ffiCall ffiCall, +) func(*paimonPostponeFixedBucketTableWrite) (*paimonPostponeFixedBucketCommitMessages, error) { + return func( + write *paimonPostponeFixedBucketTableWrite, + ) (*paimonPostponeFixedBucketCommitMessages, error) { + var result resultPostponeFixedBucketPrepareCommit + ffiCall(unsafe.Pointer(&result), unsafe.Pointer(&write)) + if result.error != nil { + return nil, parseError(ctx, result.error) + } + return result.messages, nil + } +}) + +var ffiPostponeFixedBucketTableCommitFree = newFFI(ffiOpts{ + sym: "paimon_postpone_fixed_bucket_table_commit_free", + rType: &ffi.TypeVoid, + aTypes: []*ffi.Type{&ffi.TypePointer}, +}, func( + _ context.Context, + ffiCall ffiCall, +) func(*paimonPostponeFixedBucketTableCommit) { + return func(commit *paimonPostponeFixedBucketTableCommit) { + ffiCall(nil, unsafe.Pointer(&commit)) + } +}) + +var ffiPostponeFixedBucketCommitMessagesFree = newFFI(ffiOpts{ + sym: "paimon_postpone_fixed_bucket_commit_messages_free", + rType: &ffi.TypeVoid, + aTypes: []*ffi.Type{&ffi.TypePointer}, +}, func( + _ context.Context, + ffiCall ffiCall, +) func(*paimonPostponeFixedBucketCommitMessages) { + return func(messages *paimonPostponeFixedBucketCommitMessages) { + ffiCall(nil, unsafe.Pointer(&messages)) + } +}) + +var ffiPostponeFixedBucketCommitMessagesMerge = newFFI(ffiOpts{ + sym: "paimon_postpone_fixed_bucket_commit_messages_merge", + rType: &ffi.TypePointer, + aTypes: []*ffi.Type{&ffi.TypePointer, &ffi.TypePointer}, +}, func( + ctx context.Context, + ffiCall ffiCall, +) func(*paimonPostponeFixedBucketCommitMessages, *paimonPostponeFixedBucketCommitMessages) error { + return func( + target *paimonPostponeFixedBucketCommitMessages, + source *paimonPostponeFixedBucketCommitMessages, + ) error { + var ffiError *paimonError + ffiCall( + unsafe.Pointer(&ffiError), + unsafe.Pointer(&target), + unsafe.Pointer(&source), + ) + return parseError(ctx, ffiError) + } +}) + +var ffiPostponeFixedBucketTableCommitCommit = newFixedCommitMessagesFFI( + "paimon_postpone_fixed_bucket_table_commit_commit", +) +var ffiPostponeFixedBucketTableCommitCommitWithIdentifier = newFixedCommitMessagesIdentifierFFI( + "paimon_postpone_fixed_bucket_table_commit_commit_with_identifier", +) +var ffiPostponeFixedBucketTableCommitFilterAndCommitWithIdentifier = newFixedCommitMessagesIdentifierFFI( + "paimon_postpone_fixed_bucket_table_commit_filter_and_commit_with_identifier", +) +var ffiPostponeFixedBucketTableCommitAbort = newFixedCommitMessagesFFI( + "paimon_postpone_fixed_bucket_table_commit_abort", +) + +func newFixedCommitMessagesFFI( + symbol contextKey, +) *FFI[func( + *paimonPostponeFixedBucketTableCommit, + *paimonPostponeFixedBucketCommitMessages, +) error] { + return newFFI(ffiOpts{ + sym: symbol, + rType: &ffi.TypePointer, + aTypes: []*ffi.Type{&ffi.TypePointer, &ffi.TypePointer}, + }, func( + ctx context.Context, + ffiCall ffiCall, + ) func( + *paimonPostponeFixedBucketTableCommit, + *paimonPostponeFixedBucketCommitMessages, + ) error { + return func( + commit *paimonPostponeFixedBucketTableCommit, + messages *paimonPostponeFixedBucketCommitMessages, + ) error { + var ffiError *paimonError + ffiCall( + unsafe.Pointer(&ffiError), + unsafe.Pointer(&commit), + unsafe.Pointer(&messages), + ) + return parseError(ctx, ffiError) + } + }) +} + +func newFixedCommitMessagesIdentifierFFI( + symbol contextKey, +) *FFI[func( + *paimonPostponeFixedBucketTableCommit, + *paimonPostponeFixedBucketCommitMessages, + int64, +) error] { + return newFFI(ffiOpts{ + sym: symbol, + rType: &ffi.TypePointer, + aTypes: []*ffi.Type{&ffi.TypePointer, &ffi.TypePointer, &ffi.TypeSint64}, + }, func( + ctx context.Context, + ffiCall ffiCall, + ) func( + *paimonPostponeFixedBucketTableCommit, + *paimonPostponeFixedBucketCommitMessages, + int64, + ) error { + return func( + commit *paimonPostponeFixedBucketTableCommit, + messages *paimonPostponeFixedBucketCommitMessages, + identifier int64, + ) error { + var ffiError *paimonError + ffiCall( + unsafe.Pointer(&ffiError), + unsafe.Pointer(&commit), + unsafe.Pointer(&messages), + unsafe.Pointer(&identifier), + ) + return parseError(ctx, ffiError) + } + }) +} + +var ffiPostponeFixedBucketTableCommitTruncateTable = newFFI(ffiOpts{ + sym: "paimon_postpone_fixed_bucket_table_commit_truncate_table", + rType: &ffi.TypePointer, + aTypes: []*ffi.Type{&ffi.TypePointer}, +}, func( + ctx context.Context, + ffiCall ffiCall, +) func(*paimonPostponeFixedBucketTableCommit) error { + return func(commit *paimonPostponeFixedBucketTableCommit) error { + var ffiError *paimonError + ffiCall(unsafe.Pointer(&ffiError), unsafe.Pointer(&commit)) + return parseError(ctx, ffiError) + } +}) + +var ffiPostponeFixedBucketTableCommitTruncateTableWithIdentifier = newFFI(ffiOpts{ + sym: "paimon_postpone_fixed_bucket_table_commit_truncate_table_with_identifier", + rType: &ffi.TypePointer, + aTypes: []*ffi.Type{&ffi.TypePointer, &ffi.TypeSint64}, +}, func( + ctx context.Context, + ffiCall ffiCall, +) func(*paimonPostponeFixedBucketTableCommit, int64) error { + return func(commit *paimonPostponeFixedBucketTableCommit, identifier int64) error { + var ffiError *paimonError + ffiCall( + unsafe.Pointer(&ffiError), + unsafe.Pointer(&commit), + unsafe.Pointer(&identifier), + ) + return parseError(ctx, ffiError) + } +}) diff --git a/bindings/go/tests/paimon_test.go b/bindings/go/tests/paimon_test.go index ed5054ce0..47d0b244f 100644 --- a/bindings/go/tests/paimon_test.go +++ b/bindings/go/tests/paimon_test.go @@ -24,6 +24,7 @@ import ( "io" "os" "path/filepath" + "reflect" "sort" "strings" "testing" @@ -39,6 +40,12 @@ type row struct { name string } +type partitionedRow struct { + id int32 + name string + dt string +} + func testWarehouse() string { warehouse := os.Getenv("PAIMON_TEST_WAREHOUSE") if warehouse == "" { @@ -157,6 +164,40 @@ func makeRecord(t *testing.T, rows []row) arrow.Record { return builder.NewRecord() } +func makePartitionedRecord(t *testing.T, value partitionedRow) arrow.Record { + t.Helper() + + schema := arrow.NewSchema([]arrow.Field{ + {Name: "id", Type: arrow.PrimitiveTypes.Int32, Nullable: false}, + {Name: "name", Type: arrow.BinaryTypes.String, Nullable: true}, + {Name: "dt", Type: arrow.BinaryTypes.String, Nullable: true}, + }, nil) + builder := array.NewRecordBuilder(memory.DefaultAllocator, schema) + defer builder.Release() + builder.Field(0).(*array.Int32Builder).Append(value.id) + builder.Field(1).(*array.StringBuilder).Append(value.name) + builder.Field(2).(*array.StringBuilder).Append(value.dt) + return builder.NewRecord() +} + +func makePartitionedBucketPlan(t *testing.T, partitions []string, totalBuckets int32) arrow.Record { + t.Helper() + + schema := arrow.NewSchema([]arrow.Field{ + {Name: "dt", Type: arrow.BinaryTypes.String, Nullable: true}, + {Name: "total_buckets", Type: arrow.PrimitiveTypes.Int32, Nullable: false}, + }, nil) + builder := array.NewRecordBuilder(memory.DefaultAllocator, schema) + defer builder.Release() + partitionBuilder := builder.Field(0).(*array.StringBuilder) + countBuilder := builder.Field(1).(*array.Int32Builder) + for _, partition := range partitions { + partitionBuilder.Append(partition) + countBuilder.Append(totalBuckets) + } + return builder.NewRecord() +} + func readTableRows(t *testing.T, table *paimon.Table) []row { t.Helper() rb, err := table.NewReadBuilder() @@ -549,6 +590,111 @@ func TestAppendOnlyWriteMergeAndIdempotentCommit(t *testing.T) { } } +// TestPostponeFixedBucketTypesAreIsolated pins the compile-time separation +// between the standard and postpone fixed-bucket handles: neither commit +// method accepts the other's messages, so the two paths cannot be mixed. +func TestPostponeFixedBucketTypesAreIsolated(t *testing.T) { + standard, ok := reflect.TypeOf(&paimon.TableCommit{}).MethodByName("Commit") + if !ok { + t.Fatal("TableCommit.Commit not found") + } + fixed, ok := reflect.TypeOf(&paimon.PostponeFixedBucketTableCommit{}).MethodByName("Commit") + if !ok { + t.Fatal("PostponeFixedBucketTableCommit.Commit not found") + } + + standardMessages := standard.Type.In(1) + fixedMessages := fixed.Type.In(1) + if standardMessages == fixedMessages { + t.Fatalf("Commit message types must stay distinct, both are %s", standardMessages) + } + if standardMessages != reflect.TypeOf(&paimon.CommitMessages{}) { + t.Errorf("TableCommit.Commit takes %s", standardMessages) + } + if fixedMessages != reflect.TypeOf(&paimon.PostponeFixedBucketCommitMessages{}) { + t.Errorf("PostponeFixedBucketTableCommit.Commit takes %s", fixedMessages) + } +} + +func TestMultiplePostponeFixedBucketWritersSharePlan(t *testing.T) { + table := openCopiedTable(t, "postpone_fixed_bucket_pk_table") + const commitUser = "go-postpone-fixed-bucket-write" + + builders := make([]*paimon.PostponeFixedBucketWriteBuilder, 2) + for index := range builders { + builder, err := table.NewPostponeFixedBucketWriteBuilderWithCommitUser(commitUser) + if err != nil { + t.Fatalf("Failed to create fixed-bucket builder %d: %v", index, err) + } + builders[index] = builder + defer builder.Close() + } + + if _, err := builders[0].NewWrite(); err == nil || !strings.Contains(err.Error(), "bucket plan is required") { + t.Fatalf("Expected missing bucket plan error, got: %v", err) + } + plan := makePartitionedBucketPlan(t, []string{"2026-08-14", "2026-08-15"}, 1) + for index, builder := range builders { + if err := builder.WithBucketPlan(plan); err != nil { + plan.Release() + t.Fatalf("Failed to set shared bucket plan on builder %d: %v", index, err) + } + } + plan.Release() + + writeAndPrepare := func( + builder *paimon.PostponeFixedBucketWriteBuilder, + value partitionedRow, + ) *paimon.PostponeFixedBucketCommitMessages { + write, err := builder.NewWrite() + if err != nil { + t.Fatalf("Failed to create fixed-bucket writer: %v", err) + } + defer write.Close() + + record := makePartitionedRecord(t, value) + if err := write.WriteArrowBatch(record); err != nil { + record.Release() + t.Fatalf("Failed to write Arrow record batch: %v", err) + } + record.Release() + messages, err := write.PrepareCommit() + if err != nil { + t.Fatalf("Failed to prepare fixed-bucket commit: %v", err) + } + return messages + } + + messages1 := writeAndPrepare(builders[0], partitionedRow{4, "dave", "2026-08-14"}) + defer messages1.Close() + messages2 := writeAndPrepare(builders[1], partitionedRow{5, "eve", "2026-08-15"}) + defer messages2.Close() + if err := messages1.Merge(messages2); err != nil { + t.Fatalf("Failed to merge fixed-bucket commit messages: %v", err) + } + + commit, err := builders[0].NewCommit() + if err != nil { + t.Fatalf("Failed to create fixed-bucket table commit: %v", err) + } + defer commit.Close() + if err := commit.Commit(messages1); err != nil { + t.Fatalf("Failed to commit fixed-bucket write: %v", err) + } + + rows := readTableRows(t, table) + sort.Slice(rows, func(i, j int) bool { return rows[i].id < rows[j].id }) + expected := []row{{4, "dave"}, {5, "eve"}} + if len(rows) != len(expected) { + t.Fatalf("Expected %d rows, got %d: %v", len(expected), len(rows), rows) + } + for index := range expected { + if rows[index] != expected[index] { + t.Errorf("Row %d: expected %v, got %v", index, expected[index], rows[index]) + } + } +} + // TestReadLogTable reads the test table and verifies the data matches expected values. // // The table was populated by Docker provisioning with: diff --git a/bindings/go/types.go b/bindings/go/types.go index a8cec54e4..95b10cef5 100644 --- a/bindings/go/types.go +++ b/bindings/go/types.go @@ -223,6 +223,10 @@ type paimonWriteBuilder struct{} type paimonTableWrite struct{} type paimonTableCommit struct{} type paimonCommitMessages struct{} +type paimonPostponeFixedBucketWriteBuilder struct{} +type paimonPostponeFixedBucketTableWrite struct{} +type paimonPostponeFixedBucketTableCommit struct{} +type paimonPostponeFixedBucketCommitMessages struct{} // Result types matching the C repr structs type resultCatalogNew struct { @@ -290,6 +294,26 @@ type resultPrepareCommit struct { error *paimonError } +type resultPostponeFixedBucketWriteBuilder struct { + writeBuilder *paimonPostponeFixedBucketWriteBuilder + error *paimonError +} + +type resultPostponeFixedBucketTableWrite struct { + write *paimonPostponeFixedBucketTableWrite + error *paimonError +} + +type resultPostponeFixedBucketTableCommit struct { + commit *paimonPostponeFixedBucketTableCommit + error *paimonError +} + +type resultPostponeFixedBucketPrepareCommit struct { + messages *paimonPostponeFixedBucketCommitMessages + error *paimonError +} + // paimonDatumC mirrors the C paimon_datum struct. type paimonDatumC struct { tag int32 diff --git a/dev/spark/provision.py b/dev/spark/provision.py index 2ce1c724c..0b209d0ed 100644 --- a/dev/spark/provision.py +++ b/dev/spark/provision.py @@ -989,6 +989,21 @@ def main(): """ ) + # Empty postpone table for Go fixed-bucket write tests. + spark.sql( + """ + CREATE TABLE IF NOT EXISTS postpone_fixed_bucket_pk_table ( + id INT, + name STRING, + dt STRING + ) USING paimon + PARTITIONED BY (dt) + TBLPROPERTIES ( + 'primary-key' = 'id,dt', + 'bucket' = '-2' + ) + """ + ) # ===== Dynamic bucket PK table (bucket=-1) ===== # Two commits with overlapping keys to exercise dynamic bucket assignment diff --git a/docs/src/go-binding.md b/docs/src/go-binding.md index 3524e8446..0bbc4025b 100644 --- a/docs/src/go-binding.md +++ b/docs/src/go-binding.md @@ -125,6 +125,63 @@ process-local and cannot be sent to another process. For primary-key fixed-bucket tables, assign each `(partition, bucket)` to one writer before writing; merging messages does not establish ownership. +### Postpone Fixed-Bucket Writes + +For a `bucket = -2` table, the ordinary builder writes postpone files. To write +real buckets directly, use the dedicated builder with a plan mapping each +partition to its bucket count: + +```go +// Plan schema: partition keys in order, then a non-null Int32 total_buckets. +schema := arrow.NewSchema([]arrow.Field{ + {Name: "dt", Type: arrow.BinaryTypes.String, Nullable: true}, + {Name: "total_buckets", Type: arrow.PrimitiveTypes.Int32, Nullable: false}, +}, nil) +rb := array.NewRecordBuilder(memory.DefaultAllocator, schema) +defer rb.Release() +rb.Field(0).(*array.StringBuilder).AppendValues([]string{"2026-08-14", "2026-08-15"}, nil) +rb.Field(1).(*array.Int32Builder).AppendValues([]int32{1, 1}, nil) +plan := rb.NewRecord() +defer plan.Release() + +builder, err := table.NewPostponeFixedBucketWriteBuilderWithCommitUser("job-1") +if err != nil { + log.Fatal(err) +} +defer builder.Close() +if err := builder.WithBucketPlan(plan); err != nil { + log.Fatal(err) +} +writer, err := builder.NewWrite() // fails without a plan +if err != nil { + log.Fatal(err) +} +defer writer.Close() + +if err := writer.WriteArrowBatch(record); err != nil { + log.Fatal(err) +} +messages, err := writer.PrepareCommit() +if err != nil { + log.Fatal(err) +} +defer messages.Close() + +commit, err := builder.NewCommit() +if err != nil { + log.Fatal(err) +} +defer commit.Close() +if err := commit.Commit(messages); err != nil { + log.Fatal(err) +} +``` + +An unpartitioned plan holds only `total_buckets`. Multiple writers in one process must share the plan +and commit user and assign each `(partition, bucket)` to one writer. Commit +messages are process-local. A fixed-bucket writer is single-use; create a new +writer after `PrepareCommit`. + ## Reading a Table Paimon Go uses a **scan-then-read** pattern: first scan the table to produce splits, then read data from those splits as Arrow RecordBatches.