From ae2121e7f206856b4c83824ba385bbbad1f1faca Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Sat, 15 Aug 2026 02:07:50 -0700 Subject: [PATCH 01/12] feat(go): add table write and postpone fixed-bucket bindings --- bindings/go/arrow_ffi.go | 81 ++++ bindings/go/postpone_fixed_bucket_write.go | 326 +++++++++++++++ .../go/postpone_fixed_bucket_write_ffi.go | 394 ++++++++++++++++++ bindings/go/tests/paimon_test.go | 387 ++++++++++++++++- bindings/go/types.go | 85 ++++ bindings/go/write.go | 299 +++++++++++++ bindings/go/write_ffi.go | 284 +++++++++++++ dev/spark/provision.py | 15 + docs/src/go-binding.md | 79 +++- 9 files changed, 1936 insertions(+), 14 deletions(-) create mode 100644 bindings/go/arrow_ffi.go create mode 100644 bindings/go/postpone_fixed_bucket_write.go create mode 100644 bindings/go/postpone_fixed_bucket_write_ffi.go create mode 100644 bindings/go/write.go create mode 100644 bindings/go/write_ffi.go diff --git a/bindings/go/arrow_ffi.go b/bindings/go/arrow_ffi.go new file mode 100644 index 000000000..0c3542b53 --- /dev/null +++ b/bindings/go/arrow_ffi.go @@ -0,0 +1,81 @@ +/* + * 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 ( + "errors" + "fmt" + "runtime" + "unsafe" + + "github.com/apache/arrow-go/v18/arrow" + "github.com/apache/arrow-go/v18/arrow/array" + "github.com/apache/arrow-go/v18/arrow/cdata" + "github.com/apache/arrow-go/v18/arrow/memory/mallocator" +) + +// cloneRecordToCMemory makes the Arrow buffers safe for the native writer to +// retain after the FFI call returns. Arrow's default Go allocator may move or +// reclaim its buffers once a cgo call ends, while postpone writes deliberately +// hold record batches until PrepareCommit. +func cloneRecordToCMemory(record arrow.Record) (arrow.Record, error) { + allocator := mallocator.NewMallocator() + columns := make([]arrow.Array, record.NumCols()) + for index := range columns { + column, err := array.Concatenate([]arrow.Array{record.Column(index)}, allocator) + if err != nil { + for _, allocated := range columns[:index] { + allocated.Release() + } + return nil, fmt.Errorf("paimon: failed to copy Arrow column %d to C memory: %w", index, err) + } + columns[index] = column + } + owned := array.NewRecord(record.Schema(), columns, record.NumRows()) + for _, column := range columns { + column.Release() + } + return owned, nil +} + +func withOwnedArrowRecord( + record arrow.Record, + nilMessage string, + operation func(unsafe.Pointer, unsafe.Pointer) error, +) error { + if record == nil { + return errors.New(nilMessage) + } + owned, err := cloneRecordToCMemory(record) + if err != nil { + return err + } + defer owned.Release() + + var array cdata.CArrowArray + var schema cdata.CArrowSchema + cdata.ExportArrowRecordBatch(owned, &array, &schema) + defer cdata.ReleaseCArrowArray(&array) + defer cdata.ReleaseCArrowSchema(&schema) + + err = operation(unsafe.Pointer(&array), unsafe.Pointer(&schema)) + runtime.KeepAlive(owned) + return err +} diff --git a/bindings/go/postpone_fixed_bucket_write.go b/bindings/go/postpone_fixed_bucket_write.go new file mode 100644 index 000000000..bf82868b4 --- /dev/null +++ b/bindings/go/postpone_fixed_bucket_write.go @@ -0,0 +1,326 @@ +/* + * 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. +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 cannot be 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 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 e9ee196f5..461cb61bc 100644 --- a/bindings/go/tests/paimon_test.go +++ b/bindings/go/tests/paimon_test.go @@ -23,10 +23,14 @@ import ( "errors" "io" "os" + "path/filepath" "sort" + "strings" "testing" + "github.com/apache/arrow-go/v18/arrow" "github.com/apache/arrow-go/v18/arrow/array" + "github.com/apache/arrow-go/v18/arrow/memory" paimon "github.com/apache/paimon-rust/bindings/go" ) @@ -35,6 +39,174 @@ 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 == "" { + return "/tmp/paimon-warehouse" + } + return warehouse +} + +func copyDirectory(source, target string) error { + info, err := os.Stat(source) + if err != nil { + return err + } + if err := os.MkdirAll(target, info.Mode()); err != nil { + return err + } + + entries, err := os.ReadDir(source) + if err != nil { + return err + } + for _, entry := range entries { + sourcePath := filepath.Join(source, entry.Name()) + targetPath := filepath.Join(target, entry.Name()) + if entry.IsDir() { + if err := copyDirectory(sourcePath, targetPath); err != nil { + return err + } + continue + } + + entryInfo, err := entry.Info() + if err != nil { + return err + } + input, err := os.Open(sourcePath) + if err != nil { + return err + } + output, err := os.OpenFile(targetPath, os.O_CREATE|os.O_WRONLY|os.O_TRUNC, entryInfo.Mode()) + if err != nil { + input.Close() + return err + } + _, copyErr := io.Copy(output, input) + inputCloseErr := input.Close() + outputCloseErr := output.Close() + if copyErr != nil { + return copyErr + } + if inputCloseErr != nil { + return inputCloseErr + } + if outputCloseErr != nil { + return outputCloseErr + } + } + return nil +} + +func openTableAt(t *testing.T, warehouse, tableName string) *paimon.Table { + t.Helper() + + catalog, err := paimon.NewCatalog(map[string]string{ + "warehouse": warehouse, + }) + if err != nil { + t.Fatalf("Failed to create catalog: %v", err) + } + t.Cleanup(func() { catalog.Close() }) + + table, err := catalog.GetTable(paimon.NewIdentifier("default", tableName)) + if err != nil { + t.Fatalf("Failed to get table: %v", err) + } + t.Cleanup(func() { table.Close() }) + return table +} + +func openCopiedTestTable(t *testing.T) *paimon.Table { + return openCopiedTable(t, "simple_pk_table") +} + +func openCopiedTable(t *testing.T, tableName string) *paimon.Table { + t.Helper() + + warehouse := testWarehouse() + source := filepath.Join(warehouse, "default.db", tableName) + if _, err := os.Stat(source); os.IsNotExist(err) { + t.Skipf("Skipping: table %s does not exist (run 'make docker-up' first)", source) + } + + targetWarehouse := t.TempDir() + target := filepath.Join(targetWarehouse, "default.db", tableName) + if err := copyDirectory(source, target); err != nil { + t.Fatalf("Failed to copy test table: %v", err) + } + return openTableAt(t, targetWarehouse, tableName) +} + +func makeRecord(t *testing.T, rows []row) 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}, + }, nil) + builder := array.NewRecordBuilder(memory.DefaultAllocator, schema) + defer builder.Release() + idBuilder := builder.Field(0).(*array.Int32Builder) + nameBuilder := builder.Field(1).(*array.StringBuilder) + for _, value := range rows { + idBuilder.Append(value.id) + nameBuilder.Append(value.name) + } + 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() + if err != nil { + t.Fatalf("Failed to create read builder: %v", err) + } + defer rb.Close() + return readRows(t, rb) +} + // readRows scans and reads all (id, name) rows from a ReadBuilder. func readRows(t *testing.T, rb *paimon.ReadBuilder) []row { t.Helper() @@ -107,29 +279,218 @@ func readRows(t *testing.T, rb *paimon.ReadBuilder) []row { func openTestTable(t *testing.T) *paimon.Table { t.Helper() - warehouse := os.Getenv("PAIMON_TEST_WAREHOUSE") - if warehouse == "" { - warehouse = "/tmp/paimon-warehouse" - } + warehouse := testWarehouse() if _, err := os.Stat(warehouse); os.IsNotExist(err) { t.Skipf("Skipping: warehouse %s does not exist (run 'make docker-up' first)", warehouse) } + return openTableAt(t, warehouse, "simple_log_table") +} - catalog, err := paimon.NewCatalog(map[string]string{ - "warehouse": warehouse, - }) +func TestWriteCommitReadRoundTrip(t *testing.T) { + table := openCopiedTestTable(t) + + builder, err := table.NewWriteBuilder() if err != nil { - t.Fatalf("Failed to create catalog: %v", err) + t.Fatalf("Failed to create write builder: %v", err) } - t.Cleanup(func() { catalog.Close() }) + defer builder.Close() - table, err := catalog.GetTable(paimon.NewIdentifier("default", "simple_log_table")) + write, err := builder.NewWrite() if err != nil { - t.Fatalf("Failed to get table: %v", err) + t.Fatalf("Failed to create table write: %v", err) } - t.Cleanup(func() { table.Close() }) + defer write.Close() - return table + record := makeRecord(t, []row{{4, "dave"}}) + 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 commit: %v", err) + } + defer messages.Close() + + commit, err := builder.NewCommit() + if err != nil { + t.Fatalf("Failed to create table commit: %v", err) + } + defer commit.Close() + if err := commit.Commit(messages); err != nil { + t.Fatalf("Failed to commit: %v", err) + } + + rows := readTableRows(t, table) + sort.Slice(rows, func(i, j int) bool { return rows[i].id < rows[j].id }) + expected := []row{{1, "alice"}, {2, "bob"}, {3, "carol"}, {4, "dave"}} + if len(rows) != len(expected) { + t.Fatalf("Expected %d rows, got %d: %v", len(expected), len(rows), rows) + } + for i := range expected { + if rows[i] != expected[i] { + t.Errorf("Row %d: expected %v, got %v", i, expected[i], rows[i]) + } + } +} + +func TestWriteMergeAndIdempotentCommit(t *testing.T) { + table := openCopiedTestTable(t) + const commitUser = "go-binding-distributed-write" + + builder1, err := table.NewWriteBuilderWithCommitUser(commitUser) + if err != nil { + t.Fatalf("Failed to create first write builder: %v", err) + } + defer builder1.Close() + builder2, err := table.NewWriteBuilderWithCommitUser(commitUser) + if err != nil { + t.Fatalf("Failed to create second write builder: %v", err) + } + defer builder2.Close() + + writeAndPrepare := func(builder *paimon.WriteBuilder, value row) *paimon.CommitMessages { + write, err := builder.NewWrite() + if err != nil { + t.Fatalf("Failed to create table write: %v", err) + } + defer write.Close() + record := makeRecord(t, []row{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 commit: %v", err) + } + return messages + } + + messages1 := writeAndPrepare(builder1, row{4, "dave"}) + defer messages1.Close() + messages2 := writeAndPrepare(builder2, row{5, "eve"}) + defer messages2.Close() + if err := messages1.Merge(messages2); err != nil { + t.Fatalf("Failed to merge commit messages: %v", err) + } + + commit, err := builder1.NewCommit() + if err != nil { + t.Fatalf("Failed to create table commit: %v", err) + } + defer commit.Close() + if err := commit.CommitWithIdentifier(messages1, 7); err != nil { + t.Fatalf("Failed to commit with identifier: %v", err) + } + + retryBuilder, err := table.NewWriteBuilderWithCommitUser(commitUser) + if err != nil { + t.Fatalf("Failed to create retry write builder: %v", err) + } + defer retryBuilder.Close() + retryCommit, err := retryBuilder.NewCommit() + if err != nil { + t.Fatalf("Failed to create retry table commit: %v", err) + } + defer retryCommit.Close() + if err := retryCommit.FilterAndCommitWithIdentifier(messages1, 7); err != nil { + t.Fatalf("Failed to retry commit idempotently: %v", err) + } + + rows := readTableRows(t, table) + sort.Slice(rows, func(i, j int) bool { return rows[i].id < rows[j].id }) + expected := []row{{1, "alice"}, {2, "bob"}, {3, "carol"}, {4, "dave"}, {5, "eve"}} + if len(rows) != len(expected) { + t.Fatalf("Expected %d rows, got %d: %v", len(expected), len(rows), rows) + } + for i := range expected { + if rows[i] != expected[i] { + t.Errorf("Row %d: expected %v, got %v", i, expected[i], rows[i]) + } + } +} + +func TestDistributedPostponeFixedBucketWritersSharePlan(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. diff --git a/bindings/go/types.go b/bindings/go/types.go index 7d943cb84..95b10cef5 100644 --- a/bindings/go/types.go +++ b/bindings/go/types.go @@ -130,6 +130,43 @@ var ( }[0], } + // Write result types all have the layout { opaque pointer, *paimon_error }. + typeResultWriteBuilder = ffi.Type{ + Type: ffi.Struct, + Elements: &[]*ffi.Type{ + &ffi.TypePointer, + &ffi.TypePointer, + nil, + }[0], + } + + typeResultTableWrite = ffi.Type{ + Type: ffi.Struct, + Elements: &[]*ffi.Type{ + &ffi.TypePointer, + &ffi.TypePointer, + nil, + }[0], + } + + typeResultTableCommit = ffi.Type{ + Type: ffi.Struct, + Elements: &[]*ffi.Type{ + &ffi.TypePointer, + &ffi.TypePointer, + nil, + }[0], + } + + typeResultPrepareCommit = ffi.Type{ + Type: ffi.Struct, + Elements: &[]*ffi.Type{ + &ffi.TypePointer, + &ffi.TypePointer, + nil, + }[0], + } + // paimon_datum { tag: i32, int_val: i64, double_val: f64, str_data: *u8, str_len: usize, // int_val2: i64, uint_val: u32, uint_val2: u32 } typePaimonDatum = ffi.Type{ @@ -182,6 +219,14 @@ type paimonTableRead struct{} type paimonPlan struct{} type paimonRecordBatchReader struct{} type paimonPredicate struct{} +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 { @@ -229,6 +274,46 @@ type resultPredicate struct { error *paimonError } +type resultWriteBuilder struct { + writeBuilder *paimonWriteBuilder + error *paimonError +} + +type resultTableWrite struct { + write *paimonTableWrite + error *paimonError +} + +type resultTableCommit struct { + commit *paimonTableCommit + error *paimonError +} + +type resultPrepareCommit struct { + messages *paimonCommitMessages + 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/bindings/go/write.go b/bindings/go/write.go new file mode 100644 index 000000000..979558f62 --- /dev/null +++ b/bindings/go/write.go @@ -0,0 +1,299 @@ +/* + * 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" +) + +// WriteBuilder creates writers and committers that share one commit identity. +// Writers whose messages are combined into one logical commit must use builders +// created with the same caller-provided commit user. +type WriteBuilder struct { + ctx context.Context + lib *libRef + inner *paimonWriteBuilder + closeOnce sync.Once +} + +// NewWriteBuilder creates a write builder with an automatically generated +// commit identity. +func (t *Table) NewWriteBuilder() (*WriteBuilder, error) { + if t.inner == nil { + return nil, ErrClosed + } + inner, err := ffiTableNewWriteBuilder.symbol(t.ctx)(t.inner) + if err != nil { + return nil, err + } + t.lib.acquire() + return &WriteBuilder{ctx: t.ctx, lib: t.lib, inner: inner}, nil +} + +// NewWriteBuilderWithCommitUser creates a write builder with a stable commit +// identity. Use the same commitUser for distributed writers whose messages will +// be merged into one commit, or when retrying a commit with an identifier. +func (t *Table) NewWriteBuilderWithCommitUser(commitUser string) (*WriteBuilder, error) { + if t.inner == nil { + return nil, ErrClosed + } + inner, err := ffiTableNewWriteBuilderWithCommitUser.symbol(t.ctx)(t.inner, commitUser) + if err != nil { + return nil, err + } + t.lib.acquire() + return &WriteBuilder{ctx: t.ctx, lib: t.lib, inner: inner}, nil +} + +// Close releases the write builder resources. Safe to call multiple times. +func (wb *WriteBuilder) Close() { + wb.closeOnce.Do(func() { + ffiWriteBuilderFree.symbol(wb.ctx)(wb.inner) + wb.inner = nil + wb.lib.release() + }) +} + +// WithOverwrite enables overwrite mode for writers created by this builder. +// Commit their messages with TableCommit.Overwrite rather than Commit. +func (wb *WriteBuilder) WithOverwrite() error { + if wb.inner == nil { + return ErrClosed + } + return ffiWriteBuilderWithOverwrite.symbol(wb.ctx)(wb.inner) +} + +// NewWrite creates a writer that accumulates Arrow record batches. +func (wb *WriteBuilder) NewWrite() (*TableWrite, error) { + if wb.inner == nil { + return nil, ErrClosed + } + inner, err := ffiWriteBuilderNewWrite.symbol(wb.ctx)(wb.inner) + if err != nil { + return nil, err + } + wb.lib.acquire() + return &TableWrite{ctx: wb.ctx, lib: wb.lib, inner: inner}, nil +} + +// NewCommit creates a committer that shares this builder's commit identity. +func (wb *WriteBuilder) NewCommit() (*TableCommit, error) { + if wb.inner == nil { + return nil, ErrClosed + } + inner, err := ffiWriteBuilderNewCommit.symbol(wb.ctx)(wb.inner) + if err != nil { + return nil, err + } + wb.lib.acquire() + return &TableCommit{ctx: wb.ctx, lib: wb.lib, inner: inner}, nil +} + +// TableWrite accumulates Arrow record batches until PrepareCommit is called. +type TableWrite struct { + ctx context.Context + lib *libRef + inner *paimonTableWrite + closeOnce sync.Once +} + +// Close releases the writer resources. Unprepared data is discarded. Safe to +// call multiple times. +func (tw *TableWrite) Close() { + tw.closeOnce.Do(func() { + ffiTableWriteFree.symbol(tw.ctx)(tw.inner) + tw.inner = nil + tw.lib.release() + }) +} + +// WriteArrowBatch writes one Arrow record batch. Its field count, order, names, +// and types must match the table schema. The record remains owned by the caller. +func (tw *TableWrite) 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 ffiTableWriteWriteArrowBatch.symbol(tw.ctx)(tw.inner, array, schema) + }, + ) +} + +// PrepareCommit closes the current file writers and returns opaque commit +// messages. The writer can be reused for another round of writes afterwards. +func (tw *TableWrite) PrepareCommit() (*CommitMessages, error) { + if tw.inner == nil { + return nil, ErrClosed + } + inner, err := ffiTableWritePrepareCommit.symbol(tw.ctx)(tw.inner) + if err != nil { + return nil, err + } + tw.lib.acquire() + return &CommitMessages{ctx: tw.ctx, lib: tw.lib, inner: inner}, nil +} + +// CommitMessages contains the files produced by one or more writers. +type CommitMessages struct { + ctx context.Context + lib *libRef + inner *paimonCommitMessages + closeOnce sync.Once +} + +// Close releases the commit messages. Safe to call multiple times. +func (m *CommitMessages) Close() { + m.closeOnce.Do(func() { + ffiCommitMessagesFree.symbol(m.ctx)(m.inner) + m.inner = nil + m.lib.release() + }) +} + +// Merge appends a copy of source's messages to this handle. Both handles remain +// valid and must be closed separately. They must share a table and commit user. +func (m *CommitMessages) Merge(source *CommitMessages) error { + if m.inner == nil { + return ErrClosed + } + if source == nil || source.inner == nil { + return ErrClosed + } + return ffiCommitMessagesMerge.symbol(m.ctx)(m.inner, source.inner) +} + +// TableCommit persists or aborts prepared commit messages. +type TableCommit struct { + ctx context.Context + lib *libRef + inner *paimonTableCommit + closeOnce sync.Once +} + +// Close releases the committer resources. Safe to call multiple times. +func (tc *TableCommit) Close() { + tc.closeOnce.Do(func() { + ffiTableCommitFree.symbol(tc.ctx)(tc.inner) + tc.inner = nil + tc.lib.release() + }) +} + +func (tc *TableCommit) withMessages( + messages *CommitMessages, + operation func(*paimonTableCommit, *paimonCommitMessages) error, +) error { + if tc.inner == nil { + return ErrClosed + } + if messages == nil || messages.inner == nil { + return ErrClosed + } + return operation(tc.inner, messages.inner) +} + +func (tc *TableCommit) withMessagesAndIdentifier( + messages *CommitMessages, + commitIdentifier int64, + operation func(*paimonTableCommit, *paimonCommitMessages, 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 appends the prepared data to the table. +func (tc *TableCommit) Commit(messages *CommitMessages) error { + return tc.withMessages(messages, ffiTableCommitCommit.symbol(tc.ctx)) +} + +// CommitWithIdentifier appends data with a caller-provided monotonically +// increasing identifier. +func (tc *TableCommit) CommitWithIdentifier(messages *CommitMessages, commitIdentifier int64) error { + return tc.withMessagesAndIdentifier( + messages, + commitIdentifier, + ffiTableCommitCommitWithIdentifier.symbol(tc.ctx), + ) +} + +// FilterAndCommitWithIdentifier makes a retry idempotent by filtering a +// previously committed identifier before committing it if it is new. +func (tc *TableCommit) FilterAndCommitWithIdentifier( + messages *CommitMessages, + commitIdentifier int64, +) error { + return tc.withMessagesAndIdentifier( + messages, + commitIdentifier, + ffiTableCommitFilterAndCommitWithIdentifier.symbol(tc.ctx), + ) +} + +// Overwrite replaces data in the partitions written by an overwrite-enabled +// WriteBuilder. +func (tc *TableCommit) Overwrite(messages *CommitMessages) error { + return tc.withMessages(messages, ffiTableCommitOverwrite.symbol(tc.ctx)) +} + +// OverwriteWithIdentifier overwrites data with a stable commit identifier. +func (tc *TableCommit) OverwriteWithIdentifier( + messages *CommitMessages, + commitIdentifier int64, +) error { + return tc.withMessagesAndIdentifier( + messages, + commitIdentifier, + ffiTableCommitOverwriteWithIdentifier.symbol(tc.ctx), + ) +} + +// TruncateTable removes all table data. +func (tc *TableCommit) TruncateTable() error { + if tc.inner == nil { + return ErrClosed + } + return ffiTableCommitTruncateTable.symbol(tc.ctx)(tc.inner) +} + +// TruncateTableWithIdentifier removes all table data with a stable commit +// identifier. +func (tc *TableCommit) TruncateTableWithIdentifier(commitIdentifier int64) error { + if tc.inner == nil { + return ErrClosed + } + return ffiTableCommitTruncateTableWithIdentifier.symbol(tc.ctx)(tc.inner, commitIdentifier) +} + +// Abort performs best-effort cleanup of files created for a prepared commit. +func (tc *TableCommit) Abort(messages *CommitMessages) error { + return tc.withMessages(messages, ffiTableCommitAbort.symbol(tc.ctx)) +} diff --git a/bindings/go/write_ffi.go b/bindings/go/write_ffi.go new file mode 100644 index 000000000..539fafdcb --- /dev/null +++ b/bindings/go/write_ffi.go @@ -0,0 +1,284 @@ +/* + * 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 ffiTableNewWriteBuilder = newFFI(ffiOpts{ + sym: "paimon_table_new_write_builder", + rType: &typeResultWriteBuilder, + aTypes: []*ffi.Type{&ffi.TypePointer}, +}, func(ctx context.Context, ffiCall ffiCall) func(*paimonTable) (*paimonWriteBuilder, error) { + return func(table *paimonTable) (*paimonWriteBuilder, error) { + var result resultWriteBuilder + ffiCall(unsafe.Pointer(&result), unsafe.Pointer(&table)) + if result.error != nil { + return nil, parseError(ctx, result.error) + } + return result.writeBuilder, nil + } +}) + +var ffiTableNewWriteBuilderWithCommitUser = newFFI(ffiOpts{ + sym: "paimon_table_new_write_builder_with_commit_user", + rType: &typeResultWriteBuilder, + aTypes: []*ffi.Type{&ffi.TypePointer, &ffi.TypePointer}, +}, func(ctx context.Context, ffiCall ffiCall) func(*paimonTable, string) (*paimonWriteBuilder, error) { + return func(table *paimonTable, commitUser string) (*paimonWriteBuilder, error) { + commitUserPtr, err := bytePtrFromString(commitUser) + if err != nil { + return nil, err + } + var result resultWriteBuilder + 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 ffiWriteBuilderFree = newFFI(ffiOpts{ + sym: "paimon_write_builder_free", + rType: &ffi.TypeVoid, + aTypes: []*ffi.Type{&ffi.TypePointer}, +}, func(_ context.Context, ffiCall ffiCall) func(*paimonWriteBuilder) { + return func(builder *paimonWriteBuilder) { + ffiCall(nil, unsafe.Pointer(&builder)) + } +}) + +var ffiWriteBuilderWithOverwrite = newFFI(ffiOpts{ + sym: "paimon_write_builder_with_overwrite", + rType: &ffi.TypePointer, + aTypes: []*ffi.Type{&ffi.TypePointer}, +}, func(ctx context.Context, ffiCall ffiCall) func(*paimonWriteBuilder) error { + return func(builder *paimonWriteBuilder) error { + var ffiError *paimonError + ffiCall(unsafe.Pointer(&ffiError), unsafe.Pointer(&builder)) + return parseError(ctx, ffiError) + } +}) + +var ffiWriteBuilderNewWrite = newFFI(ffiOpts{ + sym: "paimon_write_builder_new_write", + rType: &typeResultTableWrite, + aTypes: []*ffi.Type{&ffi.TypePointer}, +}, func(ctx context.Context, ffiCall ffiCall) func(*paimonWriteBuilder) (*paimonTableWrite, error) { + return func(builder *paimonWriteBuilder) (*paimonTableWrite, error) { + var result resultTableWrite + ffiCall(unsafe.Pointer(&result), unsafe.Pointer(&builder)) + if result.error != nil { + return nil, parseError(ctx, result.error) + } + return result.write, nil + } +}) + +var ffiTableWriteFree = newFFI(ffiOpts{ + sym: "paimon_table_write_free", + rType: &ffi.TypeVoid, + aTypes: []*ffi.Type{&ffi.TypePointer}, +}, func(_ context.Context, ffiCall ffiCall) func(*paimonTableWrite) { + return func(write *paimonTableWrite) { + ffiCall(nil, unsafe.Pointer(&write)) + } +}) + +var ffiTableWriteWriteArrowBatch = newFFI(ffiOpts{ + sym: "paimon_table_write_write_arrow_batch", + rType: &ffi.TypePointer, + aTypes: []*ffi.Type{&ffi.TypePointer, &ffi.TypePointer, &ffi.TypePointer}, +}, func(ctx context.Context, ffiCall ffiCall) func(*paimonTableWrite, unsafe.Pointer, unsafe.Pointer) error { + return func(write *paimonTableWrite, 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 ffiTableWritePrepareCommit = newFFI(ffiOpts{ + sym: "paimon_table_write_prepare_commit", + rType: &typeResultPrepareCommit, + aTypes: []*ffi.Type{&ffi.TypePointer}, +}, func(ctx context.Context, ffiCall ffiCall) func(*paimonTableWrite) (*paimonCommitMessages, error) { + return func(write *paimonTableWrite) (*paimonCommitMessages, error) { + var result resultPrepareCommit + ffiCall(unsafe.Pointer(&result), unsafe.Pointer(&write)) + if result.error != nil { + return nil, parseError(ctx, result.error) + } + return result.messages, nil + } +}) + +var ffiWriteBuilderNewCommit = newFFI(ffiOpts{ + sym: "paimon_write_builder_new_commit", + rType: &typeResultTableCommit, + aTypes: []*ffi.Type{&ffi.TypePointer}, +}, func(ctx context.Context, ffiCall ffiCall) func(*paimonWriteBuilder) (*paimonTableCommit, error) { + return func(builder *paimonWriteBuilder) (*paimonTableCommit, error) { + var result resultTableCommit + ffiCall(unsafe.Pointer(&result), unsafe.Pointer(&builder)) + if result.error != nil { + return nil, parseError(ctx, result.error) + } + return result.commit, nil + } +}) + +var ffiTableCommitFree = newFFI(ffiOpts{ + sym: "paimon_table_commit_free", + rType: &ffi.TypeVoid, + aTypes: []*ffi.Type{&ffi.TypePointer}, +}, func(_ context.Context, ffiCall ffiCall) func(*paimonTableCommit) { + return func(commit *paimonTableCommit) { + ffiCall(nil, unsafe.Pointer(&commit)) + } +}) + +var ffiCommitMessagesFree = newFFI(ffiOpts{ + sym: "paimon_commit_messages_free", + rType: &ffi.TypeVoid, + aTypes: []*ffi.Type{&ffi.TypePointer}, +}, func(_ context.Context, ffiCall ffiCall) func(*paimonCommitMessages) { + return func(messages *paimonCommitMessages) { + ffiCall(nil, unsafe.Pointer(&messages)) + } +}) + +var ffiCommitMessagesMerge = newFFI(ffiOpts{ + sym: "paimon_commit_messages_merge", + rType: &ffi.TypePointer, + aTypes: []*ffi.Type{&ffi.TypePointer, &ffi.TypePointer}, +}, func(ctx context.Context, ffiCall ffiCall) func(*paimonCommitMessages, *paimonCommitMessages) error { + return func(target *paimonCommitMessages, source *paimonCommitMessages) error { + var ffiError *paimonError + ffiCall( + unsafe.Pointer(&ffiError), + unsafe.Pointer(&target), + unsafe.Pointer(&source), + ) + return parseError(ctx, ffiError) + } +}) + +var ffiTableCommitCommit = newCommitMessagesFFI("paimon_table_commit_commit") +var ffiTableCommitCommitWithIdentifier = newCommitMessagesIdentifierFFI( + "paimon_table_commit_commit_with_identifier", +) +var ffiTableCommitFilterAndCommitWithIdentifier = newCommitMessagesIdentifierFFI( + "paimon_table_commit_filter_and_commit_with_identifier", +) +var ffiTableCommitOverwrite = newCommitMessagesFFI("paimon_table_commit_overwrite") +var ffiTableCommitOverwriteWithIdentifier = newCommitMessagesIdentifierFFI( + "paimon_table_commit_overwrite_with_identifier", +) +var ffiTableCommitAbort = newCommitMessagesFFI("paimon_table_commit_abort") + +func newCommitMessagesFFI(symbol contextKey) *FFI[func(*paimonTableCommit, *paimonCommitMessages) error] { + return newFFI(ffiOpts{ + sym: symbol, + rType: &ffi.TypePointer, + aTypes: []*ffi.Type{&ffi.TypePointer, &ffi.TypePointer}, + }, func(ctx context.Context, ffiCall ffiCall) func(*paimonTableCommit, *paimonCommitMessages) error { + return func(commit *paimonTableCommit, messages *paimonCommitMessages) error { + var ffiError *paimonError + ffiCall( + unsafe.Pointer(&ffiError), + unsafe.Pointer(&commit), + unsafe.Pointer(&messages), + ) + return parseError(ctx, ffiError) + } + }) +} + +func newCommitMessagesIdentifierFFI( + symbol contextKey, +) *FFI[func(*paimonTableCommit, *paimonCommitMessages, 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(*paimonTableCommit, *paimonCommitMessages, int64) error { + return func(commit *paimonTableCommit, messages *paimonCommitMessages, identifier int64) error { + var ffiError *paimonError + ffiCall( + unsafe.Pointer(&ffiError), + unsafe.Pointer(&commit), + unsafe.Pointer(&messages), + unsafe.Pointer(&identifier), + ) + return parseError(ctx, ffiError) + } + }) +} + +var ffiTableCommitTruncateTable = newCommitIdentifierFFI( + "paimon_table_commit_truncate_table", + false, +) +var ffiTableCommitTruncateTableWithIdentifier = newCommitIdentifierFFI( + "paimon_table_commit_truncate_table_with_identifier", + true, +) + +// newCommitIdentifierFFI registers truncate operations, whose ABI differs only +// by the optional identifier. The returned function always accepts an int64; +// the no-identifier variant ignores it. +func newCommitIdentifierFFI( + symbol contextKey, + withIdentifier bool, +) *FFI[func(*paimonTableCommit, ...int64) error] { + aTypes := []*ffi.Type{&ffi.TypePointer} + if withIdentifier { + aTypes = append(aTypes, &ffi.TypeSint64) + } + return newFFI(ffiOpts{ + sym: symbol, + rType: &ffi.TypePointer, + aTypes: aTypes, + }, func(ctx context.Context, ffiCall ffiCall) func(*paimonTableCommit, ...int64) error { + return func(commit *paimonTableCommit, identifiers ...int64) error { + var ffiError *paimonError + args := []unsafe.Pointer{unsafe.Pointer(&commit)} + if withIdentifier { + args = append(args, unsafe.Pointer(&identifiers[0])) + } + ffiCall(unsafe.Pointer(&ffiError), args...) + return parseError(ctx, ffiError) + } + }) +} 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 41571112d..8b5c32b75 100644 --- a/docs/src/go-binding.md +++ b/docs/src/go-binding.md @@ -72,6 +72,82 @@ catalog, err := paimon.NewCatalog(map[string]string{ }) ``` +## Writing a Table + +Use `NewWriteBuilder` for ordinary tables, including fixed-bucket tables. For +`bucket = -2` tables, it retains postpone semantics and writes bucket `-2` +files. Input `arrow.Record` batches must match the table schema. Create the +writer and committer from the same builder. + +```go +builder, err := table.NewWriteBuilder() +if err != nil { + log.Fatal(err) +} +defer builder.Close() + +writer, err := builder.NewWrite() +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) +} +``` + +`TableWrite` accepts multiple batches before `PrepareCommit`. For distributed +writes, create every builder with `NewWriteBuilderWithCommitUser` and the same +commit user, then merge their `CommitMessages`. Identifier-based methods +support idempotent retries. Overwrite, truncate, and abort are also available. + +### Postpone Fixed-Bucket Writes + +`NewPostponeFixedBucketWriteBuilder` writes real buckets for a `bucket = -2` +table. Before `NewWrite`, provide a resolved Arrow plan containing the +partition columns in table-schema order followed by a non-null Int32 +`total_buckets` column. An unpartitioned plan contains only `total_buckets`. + +```go +builder, err := table.NewPostponeFixedBucketWriteBuilderWithCommitUser(commitUser) +if err != nil { + log.Fatal(err) +} +defer builder.Close() + +if err := builder.WithBucketPlan(bucketPlan); err != nil { + log.Fatal(err) +} + +writer, err := builder.NewWrite() +if err != nil { + log.Fatal(err) +} +defer writer.Close() + +``` + +Writing, preparing, and committing use the same flow as above, but return +dedicated fixed-bucket types. The Go binding does not create the plan or shuffle +rows. Distributed integrations must share one plan and commit user, and assign +each `(partition, bucket)` to one writer before writing. + ## 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. @@ -284,7 +360,8 @@ pred, _ := pb.Eq("ts", paimon.Timestamp{Millis: 1700000000000, Nanos: 0}) ## Resource Management -All Paimon objects (`Catalog`, `Table`, `ReadBuilder`, `TableScan`, `Plan`, `TableRead`, `RecordBatchReader`) hold native resources and must be closed when no longer needed. Use `defer` to ensure cleanup: +Paimon objects with a `Close` method hold native resources and must be closed. +Use `defer` immediately after creation: ```go catalog, err := paimon.NewCatalog(opts) From 7775678ff9747c9fbf8f42dd7661a90a26a4ef1a Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Sat, 15 Aug 2026 07:30:23 -0700 Subject: [PATCH 02/12] refactor(go): split truncate FFI helpers --- bindings/go/write_ffi.go | 50 +++++++++++++++++++++------------------- 1 file changed, 26 insertions(+), 24 deletions(-) diff --git a/bindings/go/write_ffi.go b/bindings/go/write_ffi.go index 539fafdcb..3916b4c66 100644 --- a/bindings/go/write_ffi.go +++ b/bindings/go/write_ffi.go @@ -246,38 +246,40 @@ func newCommitMessagesIdentifierFFI( }) } -var ffiTableCommitTruncateTable = newCommitIdentifierFFI( - "paimon_table_commit_truncate_table", - false, -) -var ffiTableCommitTruncateTableWithIdentifier = newCommitIdentifierFFI( +var ffiTableCommitTruncateTable = newCommitFFI("paimon_table_commit_truncate_table") +var ffiTableCommitTruncateTableWithIdentifier = newCommitWithIdentifierFFI( "paimon_table_commit_truncate_table_with_identifier", - true, ) -// newCommitIdentifierFFI registers truncate operations, whose ABI differs only -// by the optional identifier. The returned function always accepts an int64; -// the no-identifier variant ignores it. -func newCommitIdentifierFFI( +func newCommitFFI(symbol contextKey) *FFI[func(*paimonTableCommit) error] { + return newFFI(ffiOpts{ + sym: symbol, + rType: &ffi.TypePointer, + aTypes: []*ffi.Type{&ffi.TypePointer}, + }, func(ctx context.Context, ffiCall ffiCall) func(*paimonTableCommit) error { + return func(commit *paimonTableCommit) error { + var ffiError *paimonError + ffiCall(unsafe.Pointer(&ffiError), unsafe.Pointer(&commit)) + return parseError(ctx, ffiError) + } + }) +} + +func newCommitWithIdentifierFFI( symbol contextKey, - withIdentifier bool, -) *FFI[func(*paimonTableCommit, ...int64) error] { - aTypes := []*ffi.Type{&ffi.TypePointer} - if withIdentifier { - aTypes = append(aTypes, &ffi.TypeSint64) - } +) *FFI[func(*paimonTableCommit, int64) error] { return newFFI(ffiOpts{ sym: symbol, rType: &ffi.TypePointer, - aTypes: aTypes, - }, func(ctx context.Context, ffiCall ffiCall) func(*paimonTableCommit, ...int64) error { - return func(commit *paimonTableCommit, identifiers ...int64) error { + aTypes: []*ffi.Type{&ffi.TypePointer, &ffi.TypeSint64}, + }, func(ctx context.Context, ffiCall ffiCall) func(*paimonTableCommit, int64) error { + return func(commit *paimonTableCommit, identifier int64) error { var ffiError *paimonError - args := []unsafe.Pointer{unsafe.Pointer(&commit)} - if withIdentifier { - args = append(args, unsafe.Pointer(&identifiers[0])) - } - ffiCall(unsafe.Pointer(&ffiError), args...) + ffiCall( + unsafe.Pointer(&ffiError), + unsafe.Pointer(&commit), + unsafe.Pointer(&identifier), + ) return parseError(ctx, ffiError) } }) From 4ebf0c78e60fda4893b3a74d351daedffe79ca1b Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Sat, 15 Aug 2026 07:33:08 -0700 Subject: [PATCH 03/12] docs(go): clarify write semantics --- bindings/go/postpone_fixed_bucket_write.go | 1 + docs/src/go-binding.md | 14 +++++++++++--- 2 files changed, 12 insertions(+), 3 deletions(-) diff --git a/bindings/go/postpone_fixed_bucket_write.go b/bindings/go/postpone_fixed_bucket_write.go index bf82868b4..5381b85a9 100644 --- a/bindings/go/postpone_fixed_bucket_write.go +++ b/bindings/go/postpone_fixed_bucket_write.go @@ -169,6 +169,7 @@ func (tw *PostponeFixedBucketTableWrite) WriteArrowBatch(record arrow.Record) er } // 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 diff --git a/docs/src/go-binding.md b/docs/src/go-binding.md index 8b5c32b75..60cce67cf 100644 --- a/docs/src/go-binding.md +++ b/docs/src/go-binding.md @@ -19,7 +19,11 @@ under the License. # Go Integration -The Go integration is a binding built on top of Apache Paimon Rust, allowing you to access Paimon tables from Go programs. It uses the [Arrow C Data Interface](https://arrow.apache.org/docs/format/CDataInterface.html) for zero-copy data transfer. +The Go integration is a binding built on top of Apache Paimon Rust, allowing you +to access Paimon tables from Go programs. It uses the +[Arrow C Data Interface](https://arrow.apache.org/docs/format/CDataInterface.html) +for data transfer. Writes copy Arrow buffers into C-owned memory because native +writers may retain batches after the Go call returns. ## Prerequisites @@ -121,8 +125,9 @@ support idempotent retries. Overwrite, truncate, and abort are also available. `NewPostponeFixedBucketWriteBuilder` writes real buckets for a `bucket = -2` table. Before `NewWrite`, provide a resolved Arrow plan containing the -partition columns in table-schema order followed by a non-null Int32 -`total_buckets` column. An unpartitioned plan contains only `total_buckets`. +partition columns in the table's partition-key order followed by a non-null +Int32 `total_buckets` column. An unpartitioned plan contains only +`total_buckets`. ```go builder, err := table.NewPostponeFixedBucketWriteBuilderWithCommitUser(commitUser) @@ -148,6 +153,9 @@ dedicated fixed-bucket types. The Go binding does not create the plan or shuffle rows. Distributed integrations must share one plan and commit user, and assign each `(partition, bucket)` to one writer before writing. +`PostponeFixedBucketTableWrite` is single-use. After calling `PrepareCommit`, +create a new writer with `builder.NewWrite()` for the next batch. + ## 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. From 91ffba7d803a05e89b368532b1e956c78e1787a0 Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Sat, 15 Aug 2026 08:06:21 -0700 Subject: [PATCH 04/12] fix(go): enforce write mode and ownership contract --- bindings/go/tests/paimon_test.go | 47 ++++++++++++++++++++++- bindings/go/write.go | 65 ++++++++++++++++---------------- docs/src/go-binding.md | 10 +++-- 3 files changed, 86 insertions(+), 36 deletions(-) diff --git a/bindings/go/tests/paimon_test.go b/bindings/go/tests/paimon_test.go index 461cb61bc..418d5acfe 100644 --- a/bindings/go/tests/paimon_test.go +++ b/bindings/go/tests/paimon_test.go @@ -336,8 +336,53 @@ func TestWriteCommitReadRoundTrip(t *testing.T) { } } -func TestWriteMergeAndIdempotentCommit(t *testing.T) { +func TestWriteOverwriteUsesBuilderMode(t *testing.T) { table := openCopiedTestTable(t) + + builder, err := table.NewWriteBuilder() + if err != nil { + t.Fatalf("Failed to create write builder: %v", err) + } + defer builder.Close() + if err := builder.WithOverwrite(); err != nil { + t.Fatalf("Failed to enable overwrite: %v", err) + } + + write, err := builder.NewWrite() + if err != nil { + t.Fatalf("Failed to create table write: %v", err) + } + defer write.Close() + record := makeRecord(t, []row{{4, "dave"}}) + 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 commit: %v", err) + } + defer messages.Close() + commit, err := builder.NewCommit() + if err != nil { + t.Fatalf("Failed to create table commit: %v", err) + } + defer commit.Close() + if err := commit.Commit(messages); err != nil { + t.Fatalf("Failed to overwrite: %v", err) + } + + rows := readTableRows(t, table) + expected := []row{{4, "dave"}} + if len(rows) != len(expected) || rows[0] != expected[0] { + t.Fatalf("Expected %v after overwrite, got %v", expected, rows) + } +} + +func TestAppendOnlyWriteMergeAndIdempotentCommit(t *testing.T) { + table := openCopiedTable(t, "simple_log_table") const commitUser = "go-binding-distributed-write" builder1, err := table.NewWriteBuilderWithCommitUser(commitUser) diff --git a/bindings/go/write.go b/bindings/go/write.go index 979558f62..5d09b4156 100644 --- a/bindings/go/write.go +++ b/bindings/go/write.go @@ -34,6 +34,7 @@ type WriteBuilder struct { ctx context.Context lib *libRef inner *paimonWriteBuilder + overwrite bool closeOnce sync.Once } @@ -53,7 +54,8 @@ func (t *Table) NewWriteBuilder() (*WriteBuilder, error) { // NewWriteBuilderWithCommitUser creates a write builder with a stable commit // identity. Use the same commitUser for distributed writers whose messages will -// be merged into one commit, or when retrying a commit with an identifier. +// be merged into one commit, or when retrying a commit with an identifier. For +// fixed-bucket primary-key tables, assign each partition and bucket to one writer. func (t *Table) NewWriteBuilderWithCommitUser(commitUser string) (*WriteBuilder, error) { if t.inner == nil { return nil, ErrClosed @@ -75,13 +77,17 @@ func (wb *WriteBuilder) Close() { }) } -// WithOverwrite enables overwrite mode for writers created by this builder. -// Commit their messages with TableCommit.Overwrite rather than Commit. +// WithOverwrite enables overwrite mode for writers and committers created by +// this builder. func (wb *WriteBuilder) WithOverwrite() error { if wb.inner == nil { return ErrClosed } - return ffiWriteBuilderWithOverwrite.symbol(wb.ctx)(wb.inner) + if err := ffiWriteBuilderWithOverwrite.symbol(wb.ctx)(wb.inner); err != nil { + return err + } + wb.overwrite = true + return nil } // NewWrite creates a writer that accumulates Arrow record batches. @@ -97,7 +103,7 @@ func (wb *WriteBuilder) NewWrite() (*TableWrite, error) { return &TableWrite{ctx: wb.ctx, lib: wb.lib, inner: inner}, nil } -// NewCommit creates a committer that shares this builder's commit identity. +// NewCommit creates a committer that shares this builder's identity and mode. func (wb *WriteBuilder) NewCommit() (*TableCommit, error) { if wb.inner == nil { return nil, ErrClosed @@ -107,7 +113,7 @@ func (wb *WriteBuilder) NewCommit() (*TableCommit, error) { return nil, err } wb.lib.acquire() - return &TableCommit{ctx: wb.ctx, lib: wb.lib, inner: inner}, nil + return &TableCommit{ctx: wb.ctx, lib: wb.lib, inner: inner, overwrite: wb.overwrite}, nil } // TableWrite accumulates Arrow record batches until PrepareCommit is called. @@ -176,6 +182,7 @@ func (m *CommitMessages) Close() { // Merge appends a copy of source's messages to this handle. Both handles remain // valid and must be closed separately. They must share a table and commit user. +// Merging does not establish fixed-bucket primary-key ownership. func (m *CommitMessages) Merge(source *CommitMessages) error { if m.inner == nil { return ErrClosed @@ -191,6 +198,7 @@ type TableCommit struct { ctx context.Context lib *libRef inner *paimonTableCommit + overwrite bool closeOnce sync.Once } @@ -230,49 +238,42 @@ func (tc *TableCommit) withMessagesAndIdentifier( return operation(tc.inner, messages.inner, commitIdentifier) } -// Commit appends the prepared data to the table. +// Commit persists messages using the builder's append or overwrite mode. func (tc *TableCommit) Commit(messages *CommitMessages) error { - return tc.withMessages(messages, ffiTableCommitCommit.symbol(tc.ctx)) + operation := ffiTableCommitCommit.symbol(tc.ctx) + if tc.overwrite { + operation = ffiTableCommitOverwrite.symbol(tc.ctx) + } + return tc.withMessages(messages, operation) } -// CommitWithIdentifier appends data with a caller-provided monotonically -// increasing identifier. +// CommitWithIdentifier commits with a caller-provided monotonically increasing +// identifier. func (tc *TableCommit) CommitWithIdentifier(messages *CommitMessages, commitIdentifier int64) error { + operation := ffiTableCommitCommitWithIdentifier.symbol(tc.ctx) + if tc.overwrite { + operation = ffiTableCommitOverwriteWithIdentifier.symbol(tc.ctx) + } return tc.withMessagesAndIdentifier( messages, commitIdentifier, - ffiTableCommitCommitWithIdentifier.symbol(tc.ctx), + operation, ) } -// FilterAndCommitWithIdentifier makes a retry idempotent by filtering a -// previously committed identifier before committing it if it is new. +// FilterAndCommitWithIdentifier makes a retry idempotent. func (tc *TableCommit) FilterAndCommitWithIdentifier( messages *CommitMessages, commitIdentifier int64, ) error { + operation := ffiTableCommitFilterAndCommitWithIdentifier.symbol(tc.ctx) + if tc.overwrite { + operation = ffiTableCommitOverwriteWithIdentifier.symbol(tc.ctx) + } return tc.withMessagesAndIdentifier( messages, commitIdentifier, - ffiTableCommitFilterAndCommitWithIdentifier.symbol(tc.ctx), - ) -} - -// Overwrite replaces data in the partitions written by an overwrite-enabled -// WriteBuilder. -func (tc *TableCommit) Overwrite(messages *CommitMessages) error { - return tc.withMessages(messages, ffiTableCommitOverwrite.symbol(tc.ctx)) -} - -// OverwriteWithIdentifier overwrites data with a stable commit identifier. -func (tc *TableCommit) OverwriteWithIdentifier( - messages *CommitMessages, - commitIdentifier int64, -) error { - return tc.withMessagesAndIdentifier( - messages, - commitIdentifier, - ffiTableCommitOverwriteWithIdentifier.symbol(tc.ctx), + operation, ) } diff --git a/docs/src/go-binding.md b/docs/src/go-binding.md index 60cce67cf..a363c24b7 100644 --- a/docs/src/go-binding.md +++ b/docs/src/go-binding.md @@ -117,9 +117,13 @@ if err := commit.Commit(messages); err != nil { ``` `TableWrite` accepts multiple batches before `PrepareCommit`. For distributed -writes, create every builder with `NewWriteBuilderWithCommitUser` and the same -commit user, then merge their `CommitMessages`. Identifier-based methods -support idempotent retries. Overwrite, truncate, and abort are also available. +append-only writes, create every builder with `NewWriteBuilderWithCommitUser` +and the same commit user, then merge their `CommitMessages`. Distributed +fixed-bucket primary-key writes must also assign each `(partition, bucket)` to +one writer before writing; sharing a commit user and merging messages does not +establish bucket ownership. Identifier-based methods support idempotent retries. +`WithOverwrite` makes `Commit` use overwrite mode. Truncate and abort are also +available. ### Postpone Fixed-Bucket Writes From 22b310acde31556631728e36fd7e238673f10aab Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Sun, 16 Aug 2026 02:14:27 -0700 Subject: [PATCH 05/12] docs(go): add bucket plan example --- docs/src/go-binding.md | 126 +++++++++++++++++++++++++++++------------ 1 file changed, 89 insertions(+), 37 deletions(-) diff --git a/docs/src/go-binding.md b/docs/src/go-binding.md index a363c24b7..581f463e6 100644 --- a/docs/src/go-binding.md +++ b/docs/src/go-binding.md @@ -28,6 +28,7 @@ writers may retain batches after the Go call returns. ## Prerequisites - Go 1.22.4 or later +- CGO enabled with a C toolchain - Supported platforms: Linux (amd64, arm64), macOS (amd64, arm64) ## Installation @@ -36,7 +37,17 @@ writers may retain batches after the Go call returns. go get github.com/apache/paimon-rust/bindings/go ``` -The pre-built native library is embedded in the package and automatically loaded at runtime — no manual build step is needed. +The native library is embedded and loaded automatically. Build with +`CGO_ENABLED=1`. + +For an unpacked offline package, keep the whole directory inside the project +build context and point `go.mod` to it: + +```mod +require github.com/apache/paimon-rust/bindings/go v0.0.0 + +replace github.com/apache/paimon-rust/bindings/go => ./third_party/paimon-go +``` ## Creating a Catalog @@ -78,10 +89,8 @@ catalog, err := paimon.NewCatalog(map[string]string{ ## Writing a Table -Use `NewWriteBuilder` for ordinary tables, including fixed-bucket tables. For -`bucket = -2` tables, it retains postpone semantics and writes bucket `-2` -files. Input `arrow.Record` batches must match the table schema. Create the -writer and committer from the same builder. +Use `NewWriteBuilder` for ordinary and fixed-bucket tables. The `arrow.Record` +schema must match the table schema. ```go builder, err := table.NewWriteBuilder() @@ -116,49 +125,92 @@ if err := commit.Commit(messages); err != nil { } ``` -`TableWrite` accepts multiple batches before `PrepareCommit`. For distributed -append-only writes, create every builder with `NewWriteBuilderWithCommitUser` -and the same commit user, then merge their `CommitMessages`. Distributed -fixed-bucket primary-key writes must also assign each `(partition, bucket)` to -one writer before writing; sharing a commit user and merging messages does not -establish bucket ownership. Identifier-based methods support idempotent retries. -`WithOverwrite` makes `Commit` use overwrite mode. Truncate and abort are also -available. +Call `WriteArrowBatch` multiple times before `PrepareCommit` to write multiple +batches. Use `WithOverwrite` before `NewWrite` for overwrite mode. Distributed +writers must use the same commit user and merge their messages. For +primary-key fixed-bucket tables, route each `(partition, bucket)` to one writer; +merging messages does not establish ownership. ### Postpone Fixed-Bucket Writes -`NewPostponeFixedBucketWriteBuilder` writes real buckets for a `bucket = -2` -table. Before `NewWrite`, provide a resolved Arrow plan containing the -partition columns in the table's partition-key order followed by a non-null -Int32 `total_buckets` column. An unpartitioned plan contains only -`total_buckets`. +For a `bucket = -2` table, the ordinary builder writes postpone files. To write +real buckets directly, use the dedicated builder and provide a plan mapping +each partition to its bucket count. This example plans two `dt` partitions with +one bucket each: ```go -builder, err := table.NewPostponeFixedBucketWriteBuilderWithCommitUser(commitUser) -if err != nil { - log.Fatal(err) -} -defer builder.Close() +import ( + "log" -if err := builder.WithBucketPlan(bucketPlan); err != nil { - log.Fatal(err) -} + "github.com/apache/arrow-go/v18/arrow" + "github.com/apache/arrow-go/v18/arrow/array" + "github.com/apache/arrow-go/v18/arrow/memory" + paimon "github.com/apache/paimon-rust/bindings/go" +) -writer, err := builder.NewWrite() -if err != nil { - log.Fatal(err) +func newBucketPlan(partitions []string, totalBuckets int32) arrow.Record { + 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) + bucketBuilder := builder.Field(1).(*array.Int32Builder) + for _, partition := range partitions { + partitionBuilder.Append(partition) + bucketBuilder.Append(totalBuckets) + } + return builder.NewRecord() } -defer writer.Close() -``` +func writePostponeFixedBuckets(table *paimon.Table, record arrow.Record) { + const commitUser = "my-fixed-bucket-job" + builder, err := table.NewPostponeFixedBucketWriteBuilderWithCommitUser(commitUser) + if err != nil { + log.Fatal(err) + } + defer builder.Close() + + plan := newBucketPlan([]string{"2026-08-14", "2026-08-15"}, 1) + defer plan.Release() + if err := builder.WithBucketPlan(plan); err != nil { + log.Fatal(err) + } + + writer, err := builder.NewWrite() + 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() -Writing, preparing, and committing use the same flow as above, but return -dedicated fixed-bucket types. The Go binding does not create the plan or shuffle -rows. Distributed integrations must share one plan and commit user, and assign -each `(partition, bucket)` to one writer before writing. + commit, err := builder.NewCommit() + if err != nil { + log.Fatal(err) + } + defer commit.Close() + if err := commit.Commit(messages); err != nil { + log.Fatal(err) + } +} +``` -`PostponeFixedBucketTableWrite` is single-use. After calling `PrepareCommit`, -create a new writer with `builder.NewWrite()` for the next batch. +Plan columns are the table partition keys in partition-key order, followed by a +non-null Int32 `total_buckets`; an unpartitioned plan contains only +`total_buckets`. The Go binding does not calculate this plan. Distributed +writers must share the plan and commit user, and route each `(partition, +bucket)` to one writer. After `PrepareCommit`, call `builder.NewWrite()` for the +next batch. ## Reading a Table From f0205906e7600157c9e92b61578b67fbdc652c09 Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Sun, 16 Aug 2026 02:43:35 -0700 Subject: [PATCH 06/12] docs(go): clarify process-local commit messages --- bindings/go/postpone_fixed_bucket_write.go | 8 +++++--- bindings/go/tests/paimon_test.go | 4 ++-- bindings/go/write.go | 14 ++++++++------ docs/src/go-binding.md | 18 ++++++++++-------- 4 files changed, 25 insertions(+), 19 deletions(-) diff --git a/bindings/go/postpone_fixed_bucket_write.go b/bindings/go/postpone_fixed_bucket_write.go index 5381b85a9..c037ad8a1 100644 --- a/bindings/go/postpone_fixed_bucket_write.go +++ b/bindings/go/postpone_fixed_bucket_write.go @@ -183,7 +183,8 @@ func (tw *PostponeFixedBucketTableWrite) PrepareCommit() (*PostponeFixedBucketCo } // PostponeFixedBucketCommitMessages contains files produced by fixed-bucket -// writers. It cannot be passed to a standard TableCommit. +// 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 @@ -200,8 +201,9 @@ func (m *PostponeFixedBucketCommitMessages) Close() { }) } -// Merge appends a copy of source's messages. Both builders must use the same -// table, commit user, and overwrite mode. +// 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 { diff --git a/bindings/go/tests/paimon_test.go b/bindings/go/tests/paimon_test.go index 418d5acfe..61e7331d5 100644 --- a/bindings/go/tests/paimon_test.go +++ b/bindings/go/tests/paimon_test.go @@ -383,7 +383,7 @@ func TestWriteOverwriteUsesBuilderMode(t *testing.T) { func TestAppendOnlyWriteMergeAndIdempotentCommit(t *testing.T) { table := openCopiedTable(t, "simple_log_table") - const commitUser = "go-binding-distributed-write" + const commitUser = "go-binding-multiple-writers" builder1, err := table.NewWriteBuilderWithCommitUser(commitUser) if err != nil { @@ -459,7 +459,7 @@ func TestAppendOnlyWriteMergeAndIdempotentCommit(t *testing.T) { } } -func TestDistributedPostponeFixedBucketWritersSharePlan(t *testing.T) { +func TestMultiplePostponeFixedBucketWritersSharePlan(t *testing.T) { table := openCopiedTable(t, "postpone_fixed_bucket_pk_table") const commitUser = "go-postpone-fixed-bucket-write" diff --git a/bindings/go/write.go b/bindings/go/write.go index 5d09b4156..2257cd861 100644 --- a/bindings/go/write.go +++ b/bindings/go/write.go @@ -53,8 +53,8 @@ func (t *Table) NewWriteBuilder() (*WriteBuilder, error) { } // NewWriteBuilderWithCommitUser creates a write builder with a stable commit -// identity. Use the same commitUser for distributed writers whose messages will -// be merged into one commit, or when retrying a commit with an identifier. For +// identity. Use the same commitUser for multiple writers in one process whose +// messages will be merged, or when retrying a commit with an identifier. For // fixed-bucket primary-key tables, assign each partition and bucket to one writer. func (t *Table) NewWriteBuilderWithCommitUser(commitUser string) (*WriteBuilder, error) { if t.inner == nil { @@ -163,7 +163,8 @@ func (tw *TableWrite) PrepareCommit() (*CommitMessages, error) { return &CommitMessages{ctx: tw.ctx, lib: tw.lib, inner: inner}, nil } -// CommitMessages contains the files produced by one or more writers. +// CommitMessages contains the files produced by one or more writers. It is a +// process-local native handle and cannot be transferred between processes. type CommitMessages struct { ctx context.Context lib *libRef @@ -180,9 +181,10 @@ func (m *CommitMessages) Close() { }) } -// Merge appends a copy of source's messages to this handle. Both handles remain -// valid and must be closed separately. They must share a table and commit user. -// Merging does not establish fixed-bucket primary-key ownership. +// Merge appends a copy of source's messages to this handle. Both handles must +// belong to the same process, remain valid, and be closed separately. They must +// share a table and commit user. Merging does not establish fixed-bucket +// primary-key ownership. func (m *CommitMessages) Merge(source *CommitMessages) error { if m.inner == nil { return ErrClosed diff --git a/docs/src/go-binding.md b/docs/src/go-binding.md index 581f463e6..f3354ab9f 100644 --- a/docs/src/go-binding.md +++ b/docs/src/go-binding.md @@ -126,10 +126,12 @@ if err := commit.Commit(messages); err != nil { ``` Call `WriteArrowBatch` multiple times before `PrepareCommit` to write multiple -batches. Use `WithOverwrite` before `NewWrite` for overwrite mode. Distributed -writers must use the same commit user and merge their messages. For -primary-key fixed-bucket tables, route each `(partition, bucket)` to one writer; -merging messages does not establish ownership. +batches. Use `WithOverwrite` before `NewWrite` for overwrite mode. Multiple +writers in one process must use the same commit user and merge their messages. +`CommitMessages` are process-local native handles and cannot be serialized or +transferred to another process for coordinator aggregation. For primary-key +fixed-bucket tables, route each `(partition, bucket)` to one writer; merging +messages does not establish ownership. ### Postpone Fixed-Bucket Writes @@ -207,10 +209,10 @@ func writePostponeFixedBuckets(table *paimon.Table, record arrow.Record) { Plan columns are the table partition keys in partition-key order, followed by a non-null Int32 `total_buckets`; an unpartitioned plan contains only -`total_buckets`. The Go binding does not calculate this plan. Distributed -writers must share the plan and commit user, and route each `(partition, -bucket)` to one writer. After `PrepareCommit`, call `builder.NewWrite()` for the -next batch. +`total_buckets`. The Go binding does not calculate this plan. Multiple writers +in one process must share the plan and commit user, and route each `(partition, +bucket)` to one writer. Cross-process commit-message aggregation is not +supported. After `PrepareCommit`, call `builder.NewWrite()` for the next batch. ## Reading a Table From 70cbd797391260eb715eda4d30f5206bfc52cca6 Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Sun, 16 Aug 2026 02:50:11 -0700 Subject: [PATCH 07/12] docs(go): streamline write examples --- docs/src/go-binding.md | 102 +++++++++++++---------------------------- 1 file changed, 31 insertions(+), 71 deletions(-) diff --git a/docs/src/go-binding.md b/docs/src/go-binding.md index f3354ab9f..cc118729e 100644 --- a/docs/src/go-binding.md +++ b/docs/src/go-binding.md @@ -19,11 +19,10 @@ under the License. # Go Integration -The Go integration is a binding built on top of Apache Paimon Rust, allowing you -to access Paimon tables from Go programs. It uses the -[Arrow C Data Interface](https://arrow.apache.org/docs/format/CDataInterface.html) -for data transfer. Writes copy Arrow buffers into C-owned memory because native -writers may retain batches after the Go call returns. +The Go binding uses the +[Arrow C Data Interface](https://arrow.apache.org/docs/format/CDataInterface.html). +Writes copy Arrow buffers into C-owned memory because writers may retain them +after the call returns. ## Prerequisites @@ -40,15 +39,6 @@ go get github.com/apache/paimon-rust/bindings/go The native library is embedded and loaded automatically. Build with `CGO_ENABLED=1`. -For an unpacked offline package, keep the whole directory inside the project -build context and point `go.mod` to it: - -```mod -require github.com/apache/paimon-rust/bindings/go v0.0.0 - -replace github.com/apache/paimon-rust/bindings/go => ./third_party/paimon-go -``` - ## Creating a Catalog Use `NewCatalog` with a map of options to create a catalog. The catalog type is determined by the `metastore` option (default: `filesystem`). @@ -125,13 +115,12 @@ if err := commit.Commit(messages); err != nil { } ``` -Call `WriteArrowBatch` multiple times before `PrepareCommit` to write multiple -batches. Use `WithOverwrite` before `NewWrite` for overwrite mode. Multiple -writers in one process must use the same commit user and merge their messages. -`CommitMessages` are process-local native handles and cannot be serialized or -transferred to another process for coordinator aggregation. For primary-key -fixed-bucket tables, route each `(partition, bucket)` to one writer; merging -messages does not establish ownership. +Call `WriteArrowBatch` multiple times before `PrepareCommit`. `WithOverwrite` +replaces the partitions touched by the batch. Multiple writers in one process +must share a commit user and merge their messages. Commit messages are +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 @@ -141,15 +130,6 @@ each partition to its bucket count. This example plans two `dt` partitions with one bucket each: ```go -import ( - "log" - - "github.com/apache/arrow-go/v18/arrow" - "github.com/apache/arrow-go/v18/arrow/array" - "github.com/apache/arrow-go/v18/arrow/memory" - paimon "github.com/apache/paimon-rust/bindings/go" -) - func newBucketPlan(partitions []string, totalBuckets int32) arrow.Record { schema := arrow.NewSchema([]arrow.Field{ {Name: "dt", Type: arrow.BinaryTypes.String, Nullable: true}, @@ -167,52 +147,32 @@ func newBucketPlan(partitions []string, totalBuckets int32) arrow.Record { return builder.NewRecord() } -func writePostponeFixedBuckets(table *paimon.Table, record arrow.Record) { - const commitUser = "my-fixed-bucket-job" - builder, err := table.NewPostponeFixedBucketWriteBuilderWithCommitUser(commitUser) - if err != nil { - log.Fatal(err) - } - defer builder.Close() - - plan := newBucketPlan([]string{"2026-08-14", "2026-08-15"}, 1) - defer plan.Release() - if err := builder.WithBucketPlan(plan); err != nil { - log.Fatal(err) - } - - writer, err := builder.NewWrite() - if err != nil { - log.Fatal(err) - } - defer writer.Close() +plan := newBucketPlan([]string{"2026-08-14", "2026-08-15"}, 1) +defer plan.Release() - if err := writer.WriteArrowBatch(record); err != nil { - log.Fatal(err) - } - messages, err := writer.PrepareCommit() - if err != nil { - log.Fatal(err) - } - defer messages.Close() +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) +} - commit, err := builder.NewCommit() - if err != nil { - log.Fatal(err) - } - defer commit.Close() - if err := commit.Commit(messages); err != nil { - log.Fatal(err) - } +writer, err := builder.NewWrite() +if err != nil { + log.Fatal(err) } +defer writer.Close() ``` -Plan columns are the table partition keys in partition-key order, followed by a -non-null Int32 `total_buckets`; an unpartitioned plan contains only -`total_buckets`. The Go binding does not calculate this plan. Multiple writers -in one process must share the plan and commit user, and route each `(partition, -bucket)` to one writer. Cross-process commit-message aggregation is not -supported. After `PrepareCommit`, call `builder.NewWrite()` for the next batch. +Write, prepare, and commit with the same lifecycle shown above. Plan columns +are the partition keys in partition-key order, followed by a non-null Int32 +`total_buckets`; an unpartitioned plan contains only `total_buckets`. Multiple +writers in one process must share the plan and commit user and assign each +`(partition, bucket)` to one writer. Cross-process message aggregation is not +supported. A fixed-bucket writer is single-use; create a new writer after +`PrepareCommit`. ## Reading a Table From 76a844b4305a4ae00733cff3d35b45a79eaa7f45 Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Sun, 16 Aug 2026 03:27:10 -0700 Subject: [PATCH 08/12] refactor(go): split postpone fixed-bucket bindings Narrow this PR to the standard Go table write binding. Postpone fixed-bucket builder, writer, commit messages, committer, tests, test table, and documentation move to a follow-up PR. --- bindings/go/postpone_fixed_bucket_write.go | 329 --------------- .../go/postpone_fixed_bucket_write_ffi.go | 394 ------------------ bindings/go/tests/paimon_test.go | 120 ------ bindings/go/types.go | 24 -- dev/spark/provision.py | 15 - docs/src/go-binding.md | 52 --- 6 files changed, 934 deletions(-) delete mode 100644 bindings/go/postpone_fixed_bucket_write.go delete mode 100644 bindings/go/postpone_fixed_bucket_write_ffi.go diff --git a/bindings/go/postpone_fixed_bucket_write.go b/bindings/go/postpone_fixed_bucket_write.go deleted file mode 100644 index c037ad8a1..000000000 --- a/bindings/go/postpone_fixed_bucket_write.go +++ /dev/null @@ -1,329 +0,0 @@ -/* - * 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 deleted file mode 100644 index d9154e587..000000000 --- a/bindings/go/postpone_fixed_bucket_write_ffi.go +++ /dev/null @@ -1,394 +0,0 @@ -/* - * 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 61e7331d5..55fba4158 100644 --- a/bindings/go/tests/paimon_test.go +++ b/bindings/go/tests/paimon_test.go @@ -25,7 +25,6 @@ import ( "os" "path/filepath" "sort" - "strings" "testing" "github.com/apache/arrow-go/v18/arrow" @@ -39,12 +38,6 @@ 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 == "" { @@ -163,40 +156,6 @@ 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() @@ -459,85 +418,6 @@ func TestAppendOnlyWriteMergeAndIdempotentCommit(t *testing.T) { } } -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 95b10cef5..a8cec54e4 100644 --- a/bindings/go/types.go +++ b/bindings/go/types.go @@ -223,10 +223,6 @@ 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 { @@ -294,26 +290,6 @@ 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 0b209d0ed..2ce1c724c 100644 --- a/dev/spark/provision.py +++ b/dev/spark/provision.py @@ -989,21 +989,6 @@ 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 cc118729e..e35416e82 100644 --- a/docs/src/go-binding.md +++ b/docs/src/go-binding.md @@ -122,58 +122,6 @@ 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 and provide a plan mapping -each partition to its bucket count. This example plans two `dt` partitions with -one bucket each: - -```go -func newBucketPlan(partitions []string, totalBuckets int32) arrow.Record { - 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) - bucketBuilder := builder.Field(1).(*array.Int32Builder) - for _, partition := range partitions { - partitionBuilder.Append(partition) - bucketBuilder.Append(totalBuckets) - } - return builder.NewRecord() -} - -plan := newBucketPlan([]string{"2026-08-14", "2026-08-15"}, 1) -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() -if err != nil { - log.Fatal(err) -} -defer writer.Close() -``` - -Write, prepare, and commit with the same lifecycle shown above. Plan columns -are the partition keys in partition-key order, followed by a non-null Int32 -`total_buckets`; an unpartitioned plan contains only `total_buckets`. Multiple -writers in one process must share the plan and commit user and assign each -`(partition, bucket)` to one writer. Cross-process message aggregation is not -supported. 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. From d6632e43dad17cfa1a366b5078db177e97dc3a32 Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Sun, 16 Aug 2026 03:35:00 -0700 Subject: [PATCH 09/12] docs(go): drop stale postpone reference from comment --- bindings/go/arrow_ffi.go | 6 ++---- 1 file changed, 2 insertions(+), 4 deletions(-) diff --git a/bindings/go/arrow_ffi.go b/bindings/go/arrow_ffi.go index 0c3542b53..32ac9667e 100644 --- a/bindings/go/arrow_ffi.go +++ b/bindings/go/arrow_ffi.go @@ -31,10 +31,8 @@ import ( "github.com/apache/arrow-go/v18/arrow/memory/mallocator" ) -// cloneRecordToCMemory makes the Arrow buffers safe for the native writer to -// retain after the FFI call returns. Arrow's default Go allocator may move or -// reclaim its buffers once a cgo call ends, while postpone writes deliberately -// hold record batches until PrepareCommit. +// cloneRecordToCMemory copies Arrow buffers into C-owned memory so native +// writers can retain record batches until PrepareCommit. func cloneRecordToCMemory(record arrow.Record) (arrow.Record, error) { allocator := mallocator.NewMallocator() columns := make([]arrow.Array, record.NumCols()) From 2ecff4368970963a01fe5dba2b5ed5b09c543b17 Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Sun, 16 Aug 2026 04:30:40 -0700 Subject: [PATCH 10/12] docs(go): document Arrow FFI move semantics at release sites --- bindings/go/arrow_ffi.go | 2 ++ 1 file changed, 2 insertions(+) diff --git a/bindings/go/arrow_ffi.go b/bindings/go/arrow_ffi.go index 32ac9667e..c755f347a 100644 --- a/bindings/go/arrow_ffi.go +++ b/bindings/go/arrow_ffi.go @@ -70,6 +70,8 @@ func withOwnedArrowRecord( var array cdata.CArrowArray var schema cdata.CArrowSchema cdata.ExportArrowRecordBatch(owned, &array, &schema) + // The native side imports via from_raw, marking these released + // (release = NULL); the defers only fire if an error skips the import. defer cdata.ReleaseCArrowArray(&array) defer cdata.ReleaseCArrowSchema(&schema) From 3ff1ecb2879b0bdad1a8bb3fcf6302ab212fc544 Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Sun, 16 Aug 2026 04:54:48 -0700 Subject: [PATCH 11/12] fix(go): reject filtered retries in overwrite mode The C API has no filtered overwrite, so FilterAndCommitWithIdentifier silently degraded to a non-filtered overwrite; a retry after a successful commit would re-add the same files. Return an explicit error instead, cover it with a test, and tighten doc comments. --- bindings/go/tests/paimon_test.go | 51 ++++++++++++++++++++++++++++++++ bindings/go/write.go | 40 ++++++++++--------------- docs/src/go-binding.md | 3 +- 3 files changed, 68 insertions(+), 26 deletions(-) diff --git a/bindings/go/tests/paimon_test.go b/bindings/go/tests/paimon_test.go index 55fba4158..58336a4c1 100644 --- a/bindings/go/tests/paimon_test.go +++ b/bindings/go/tests/paimon_test.go @@ -25,6 +25,7 @@ import ( "os" "path/filepath" "sort" + "strings" "testing" "github.com/apache/arrow-go/v18/arrow" @@ -340,6 +341,56 @@ func TestWriteOverwriteUsesBuilderMode(t *testing.T) { } } +func TestOverwriteRejectsFilteredRetry(t *testing.T) { + table := openCopiedTestTable(t) + + builder, err := table.NewWriteBuilder() + if err != nil { + t.Fatalf("Failed to create write builder: %v", err) + } + defer builder.Close() + if err := builder.WithOverwrite(); err != nil { + t.Fatalf("Failed to enable overwrite: %v", err) + } + + write, err := builder.NewWrite() + if err != nil { + t.Fatalf("Failed to create table write: %v", err) + } + defer write.Close() + record := makeRecord(t, []row{{4, "dave"}}) + 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 commit: %v", err) + } + defer messages.Close() + commit, err := builder.NewCommit() + if err != nil { + t.Fatalf("Failed to create table commit: %v", err) + } + defer commit.Close() + + err = commit.FilterAndCommitWithIdentifier(messages, 7) + if err == nil || !strings.Contains(err.Error(), "overwrite") { + t.Fatalf("Expected overwrite filtered-retry rejection, got %v", err) + } + if err := commit.CommitWithIdentifier(messages, 7); err != nil { + t.Fatalf("Failed to overwrite with identifier: %v", err) + } + + rows := readTableRows(t, table) + expected := []row{{4, "dave"}} + if len(rows) != len(expected) || rows[0] != expected[0] { + t.Fatalf("Expected %v after overwrite, got %v", expected, rows) + } +} + func TestAppendOnlyWriteMergeAndIdempotentCommit(t *testing.T) { table := openCopiedTable(t, "simple_log_table") const commitUser = "go-binding-multiple-writers" diff --git a/bindings/go/write.go b/bindings/go/write.go index 2257cd861..8a13d18ef 100644 --- a/bindings/go/write.go +++ b/bindings/go/write.go @@ -21,6 +21,7 @@ package paimon import ( "context" + "errors" "sync" "unsafe" @@ -28,8 +29,6 @@ import ( ) // WriteBuilder creates writers and committers that share one commit identity. -// Writers whose messages are combined into one logical commit must use builders -// created with the same caller-provided commit user. type WriteBuilder struct { ctx context.Context lib *libRef @@ -38,8 +37,7 @@ type WriteBuilder struct { closeOnce sync.Once } -// NewWriteBuilder creates a write builder with an automatically generated -// commit identity. +// NewWriteBuilder creates a write builder with a generated commit identity. func (t *Table) NewWriteBuilder() (*WriteBuilder, error) { if t.inner == nil { return nil, ErrClosed @@ -53,9 +51,7 @@ func (t *Table) NewWriteBuilder() (*WriteBuilder, error) { } // NewWriteBuilderWithCommitUser creates a write builder with a stable commit -// identity. Use the same commitUser for multiple writers in one process whose -// messages will be merged, or when retrying a commit with an identifier. For -// fixed-bucket primary-key tables, assign each partition and bucket to one writer. +// identity for merging messages across writers and for identifier retries. func (t *Table) NewWriteBuilderWithCommitUser(commitUser string) (*WriteBuilder, error) { if t.inner == nil { return nil, ErrClosed @@ -77,8 +73,7 @@ func (wb *WriteBuilder) Close() { }) } -// WithOverwrite enables overwrite mode for writers and committers created by -// this builder. +// WithOverwrite enables overwrite mode for this builder's writers and committers. func (wb *WriteBuilder) WithOverwrite() error { if wb.inner == nil { return ErrClosed @@ -134,8 +129,8 @@ func (tw *TableWrite) Close() { }) } -// WriteArrowBatch writes one Arrow record batch. Its field count, order, names, -// and types must match the table schema. The record remains owned by the caller. +// WriteArrowBatch writes one record batch whose schema must match the table +// schema. The record remains owned by the caller. func (tw *TableWrite) WriteArrowBatch(record arrow.Record) error { if tw.inner == nil { return ErrClosed @@ -149,8 +144,7 @@ func (tw *TableWrite) WriteArrowBatch(record arrow.Record) error { ) } -// PrepareCommit closes the current file writers and returns opaque commit -// messages. The writer can be reused for another round of writes afterwards. +// PrepareCommit returns pending writes as commit messages; the writer can be reused. func (tw *TableWrite) PrepareCommit() (*CommitMessages, error) { if tw.inner == nil { return nil, ErrClosed @@ -181,10 +175,8 @@ func (m *CommitMessages) Close() { }) } -// Merge appends a copy of source's messages to this handle. Both handles must -// belong to the same process, remain valid, and be closed separately. They must -// share a table and commit user. Merging does not establish fixed-bucket -// primary-key ownership. +// Merge appends a copy of source's messages; both handles stay valid and must +// share a table and commit user. It does not establish fixed-bucket ownership. func (m *CommitMessages) Merge(source *CommitMessages) error { if m.inner == nil { return ErrClosed @@ -249,8 +241,7 @@ func (tc *TableCommit) Commit(messages *CommitMessages) error { return tc.withMessages(messages, operation) } -// CommitWithIdentifier commits with a caller-provided monotonically increasing -// identifier. +// CommitWithIdentifier commits with a monotonically increasing identifier. func (tc *TableCommit) CommitWithIdentifier(messages *CommitMessages, commitIdentifier int64) error { operation := ffiTableCommitCommitWithIdentifier.symbol(tc.ctx) if tc.overwrite { @@ -263,19 +254,19 @@ func (tc *TableCommit) CommitWithIdentifier(messages *CommitMessages, commitIden ) } -// FilterAndCommitWithIdentifier makes a retry idempotent. +// FilterAndCommitWithIdentifier skips messages already committed under +// commitIdentifier, making retries idempotent. Unsupported in overwrite mode. func (tc *TableCommit) FilterAndCommitWithIdentifier( messages *CommitMessages, commitIdentifier int64, ) error { - operation := ffiTableCommitFilterAndCommitWithIdentifier.symbol(tc.ctx) if tc.overwrite { - operation = ffiTableCommitOverwriteWithIdentifier.symbol(tc.ctx) + return errors.New("paimon: filtered retry is not supported in overwrite mode") } return tc.withMessagesAndIdentifier( messages, commitIdentifier, - operation, + ffiTableCommitFilterAndCommitWithIdentifier.symbol(tc.ctx), ) } @@ -287,8 +278,7 @@ func (tc *TableCommit) TruncateTable() error { return ffiTableCommitTruncateTable.symbol(tc.ctx)(tc.inner) } -// TruncateTableWithIdentifier removes all table data with a stable commit -// identifier. +// TruncateTableWithIdentifier truncates with a stable commit identifier. func (tc *TableCommit) TruncateTableWithIdentifier(commitIdentifier int64) error { if tc.inner == nil { return ErrClosed diff --git a/docs/src/go-binding.md b/docs/src/go-binding.md index e35416e82..2284f032e 100644 --- a/docs/src/go-binding.md +++ b/docs/src/go-binding.md @@ -116,7 +116,8 @@ if err := commit.Commit(messages); err != nil { ``` Call `WriteArrowBatch` multiple times before `PrepareCommit`. `WithOverwrite` -replaces the partitions touched by the batch. Multiple writers in one process +replaces the partitions touched by the batch; filtered identifier retries are +append-only and return an error in overwrite mode. Multiple writers in one process must share a commit user and merge their messages. Commit messages are process-local and cannot be sent to another process. For primary-key fixed-bucket tables, assign each `(partition, bucket)` to one writer before From 5fecaaec95279a113e3d370a55f5e10351101672 Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Sun, 16 Aug 2026 05:09:47 -0700 Subject: [PATCH 12/12] fix(go): restore filtered overwrite retry dispatch Rust overwrite_with_identifier already runs with filter_committed=true, so routing filtered retries to it is the correct idempotent behavior and matches Java StreamTableCommit.filterAndCommit. Replace the wrong rejection with a retry test asserting no new snapshot and preserved intermediate commits, and document one-shot Commit identifier usage. --- bindings/go/tests/paimon_test.go | 136 ++++++++++++++++++++++++------- bindings/go/write.go | 9 +- docs/src/go-binding.md | 6 +- 3 files changed, 117 insertions(+), 34 deletions(-) diff --git a/bindings/go/tests/paimon_test.go b/bindings/go/tests/paimon_test.go index 58336a4c1..ed5054ce0 100644 --- a/bindings/go/tests/paimon_test.go +++ b/bindings/go/tests/paimon_test.go @@ -341,53 +341,133 @@ func TestWriteOverwriteUsesBuilderMode(t *testing.T) { } } -func TestOverwriteRejectsFilteredRetry(t *testing.T) { - table := openCopiedTestTable(t) +func TestOverwriteRetrySameIdentifierIsIdempotent(t *testing.T) { + warehouse := testWarehouse() + source := filepath.Join(warehouse, "default.db", "simple_pk_table") + if _, err := os.Stat(source); os.IsNotExist(err) { + t.Skipf("Skipping: table %s does not exist (run 'make docker-up' first)", source) + } + targetWarehouse := t.TempDir() + target := filepath.Join(targetWarehouse, "default.db", "simple_pk_table") + if err := copyDirectory(source, target); err != nil { + t.Fatalf("Failed to copy test table: %v", err) + } + table := openTableAt(t, targetWarehouse, "simple_pk_table") - builder, err := table.NewWriteBuilder() + countSnapshots := func() int { + entries, err := os.ReadDir(filepath.Join(target, "snapshot")) + if err != nil { + t.Fatalf("Failed to list snapshots: %v", err) + } + count := 0 + for _, entry := range entries { + if strings.HasPrefix(entry.Name(), "snapshot-") { + count++ + } + } + return count + } + + const commitUser = "go-binding-overwrite-retry" + overwriteAndPrepare := func() (*paimon.WriteBuilder, *paimon.CommitMessages) { + builder, err := table.NewWriteBuilderWithCommitUser(commitUser) + if err != nil { + t.Fatalf("Failed to create write builder: %v", err) + } + t.Cleanup(builder.Close) + if err := builder.WithOverwrite(); err != nil { + t.Fatalf("Failed to enable overwrite: %v", err) + } + write, err := builder.NewWrite() + if err != nil { + t.Fatalf("Failed to create table write: %v", err) + } + defer write.Close() + record := makeRecord(t, []row{{4, "dave"}}) + 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 commit: %v", err) + } + t.Cleanup(messages.Close) + return builder, messages + } + + builder, messages := overwriteAndPrepare() + commit, err := builder.NewCommit() if err != nil { - t.Fatalf("Failed to create write builder: %v", err) + t.Fatalf("Failed to create table commit: %v", err) } - defer builder.Close() - if err := builder.WithOverwrite(); err != nil { - t.Fatalf("Failed to enable overwrite: %v", err) + defer commit.Close() + if err := commit.CommitWithIdentifier(messages, 7); err != nil { + t.Fatalf("Failed to overwrite with identifier: %v", err) } - write, err := builder.NewWrite() + appendBuilder, err := table.NewWriteBuilder() if err != nil { - t.Fatalf("Failed to create table write: %v", err) + t.Fatalf("Failed to create append write builder: %v", err) } - defer write.Close() - record := makeRecord(t, []row{{4, "dave"}}) - if err := write.WriteArrowBatch(record); err != nil { + defer appendBuilder.Close() + appendWrite, err := appendBuilder.NewWrite() + if err != nil { + t.Fatalf("Failed to create append table write: %v", err) + } + defer appendWrite.Close() + record := makeRecord(t, []row{{5, "eve"}}) + if err := appendWrite.WriteArrowBatch(record); err != nil { record.Release() - t.Fatalf("Failed to write Arrow record batch: %v", err) + t.Fatalf("Failed to write append batch: %v", err) } record.Release() - - messages, err := write.PrepareCommit() + appendMessages, err := appendWrite.PrepareCommit() if err != nil { - t.Fatalf("Failed to prepare commit: %v", err) + t.Fatalf("Failed to prepare append commit: %v", err) } - defer messages.Close() - commit, err := builder.NewCommit() + defer appendMessages.Close() + appendCommit, err := appendBuilder.NewCommit() if err != nil { - t.Fatalf("Failed to create table commit: %v", err) + t.Fatalf("Failed to create append table commit: %v", err) } - defer commit.Close() + defer appendCommit.Close() + if err := appendCommit.Commit(appendMessages); err != nil { + t.Fatalf("Failed to append: %v", err) + } + snapshotsBeforeRetry := countSnapshots() - err = commit.FilterAndCommitWithIdentifier(messages, 7) - if err == nil || !strings.Contains(err.Error(), "overwrite") { - t.Fatalf("Expected overwrite filtered-retry rejection, got %v", err) + retryBuilder, err := table.NewWriteBuilderWithCommitUser(commitUser) + if err != nil { + t.Fatalf("Failed to create retry write builder: %v", err) } - if err := commit.CommitWithIdentifier(messages, 7); err != nil { - t.Fatalf("Failed to overwrite with identifier: %v", err) + defer retryBuilder.Close() + if err := retryBuilder.WithOverwrite(); err != nil { + t.Fatalf("Failed to enable overwrite on retry builder: %v", err) + } + retryCommit, err := retryBuilder.NewCommit() + if err != nil { + t.Fatalf("Failed to create retry table commit: %v", err) + } + defer retryCommit.Close() + if err := retryCommit.FilterAndCommitWithIdentifier(messages, 7); err != nil { + t.Fatalf("Failed to retry overwrite idempotently: %v", err) } + if got := countSnapshots(); got != snapshotsBeforeRetry { + t.Fatalf("Retry added snapshots: %d != %d", got, snapshotsBeforeRetry) + } rows := readTableRows(t, table) - expected := []row{{4, "dave"}} - if len(rows) != len(expected) || rows[0] != expected[0] { - t.Fatalf("Expected %v after overwrite, got %v", expected, rows) + 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 %v after retry, got %v", expected, rows) + } + for i := range expected { + if rows[i] != expected[i] { + t.Errorf("Row %d: expected %v, got %v", i, expected[i], rows[i]) + } } } diff --git a/bindings/go/write.go b/bindings/go/write.go index 8a13d18ef..2e6549602 100644 --- a/bindings/go/write.go +++ b/bindings/go/write.go @@ -21,7 +21,6 @@ package paimon import ( "context" - "errors" "sync" "unsafe" @@ -255,18 +254,20 @@ func (tc *TableCommit) CommitWithIdentifier(messages *CommitMessages, commitIden } // FilterAndCommitWithIdentifier skips messages already committed under -// commitIdentifier, making retries idempotent. Unsupported in overwrite mode. +// commitIdentifier, making retries idempotent. Identifier commits always +// filter in overwrite mode. func (tc *TableCommit) FilterAndCommitWithIdentifier( messages *CommitMessages, commitIdentifier int64, ) error { + operation := ffiTableCommitFilterAndCommitWithIdentifier.symbol(tc.ctx) if tc.overwrite { - return errors.New("paimon: filtered retry is not supported in overwrite mode") + operation = ffiTableCommitOverwriteWithIdentifier.symbol(tc.ctx) } return tc.withMessagesAndIdentifier( messages, commitIdentifier, - ffiTableCommitFilterAndCommitWithIdentifier.symbol(tc.ctx), + operation, ) } diff --git a/docs/src/go-binding.md b/docs/src/go-binding.md index 2284f032e..3524e8446 100644 --- a/docs/src/go-binding.md +++ b/docs/src/go-binding.md @@ -116,8 +116,10 @@ if err := commit.Commit(messages); err != nil { ``` Call `WriteArrowBatch` multiple times before `PrepareCommit`. `WithOverwrite` -replaces the partitions touched by the batch; filtered identifier retries are -append-only and return an error in overwrite mode. Multiple writers in one process +replaces the partitions touched by the batch. `Commit` is for one-shot batch +jobs: it consumes the maximum commit identifier, so later filtered retries by +the same commit user are treated as already committed. Reuse a writer across +rounds with `CommitWithIdentifier` and increasing identifiers. Multiple writers in one process must share a commit user and merge their messages. Commit messages are process-local and cannot be sent to another process. For primary-key fixed-bucket tables, assign each `(partition, bucket)` to one writer before