This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/paimon-rust.git
The following commit(s) were added to refs/heads/main by this push:
new be5456e9 feat(go): add table write bindings (#658)
be5456e9 is described below
commit be5456e96b326939570b6ef5d63acddda3311c4f
Author: XiaoHongbo <[email protected]>
AuthorDate: Sun Aug 16 22:46:49 2026 +0800
feat(go): add table write bindings (#658)
---
bindings/go/arrow_ffi.go | 81 +++++++
bindings/go/tests/paimon_test.go | 443 +++++++++++++++++++++++++++++++++++++--
bindings/go/types.go | 61 ++++++
bindings/go/write.go | 293 ++++++++++++++++++++++++++
bindings/go/write_ffi.go | 286 +++++++++++++++++++++++++
docs/src/go-binding.md | 60 +++++-
6 files changed, 1208 insertions(+), 16 deletions(-)
diff --git a/bindings/go/arrow_ffi.go b/bindings/go/arrow_ffi.go
new file mode 100644
index 00000000..c755f347
--- /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 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())
+ 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)
+ // 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)
+
+ err = operation(unsafe.Pointer(&array), unsafe.Pointer(&schema))
+ runtime.KeepAlive(owned)
+ return err
+}
diff --git a/bindings/go/tests/paimon_test.go b/bindings/go/tests/paimon_test.go
index e9ee196f..ed5054ce 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,134 @@ type row struct {
name 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 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 +239,314 @@ 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 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 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")
+
+ 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 table commit: %v", err)
+ }
+ defer commit.Close()
+ if err := commit.CommitWithIdentifier(messages, 7); err != nil {
+ t.Fatalf("Failed to overwrite with identifier: %v", err)
+ }
+
+ appendBuilder, err := table.NewWriteBuilder()
+ if err != nil {
+ t.Fatalf("Failed to create append write builder: %v", err)
+ }
+ 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 append batch: %v", err)
+ }
+ record.Release()
+ appendMessages, err := appendWrite.PrepareCommit()
+ if err != nil {
+ t.Fatalf("Failed to prepare append commit: %v", err)
+ }
+ defer appendMessages.Close()
+ appendCommit, err := appendBuilder.NewCommit()
+ if err != nil {
+ t.Fatalf("Failed to create append table commit: %v", err)
+ }
+ defer appendCommit.Close()
+ if err := appendCommit.Commit(appendMessages); err != nil {
+ t.Fatalf("Failed to append: %v", err)
+ }
+ snapshotsBeforeRetry := countSnapshots()
+
+ retryBuilder, err := table.NewWriteBuilderWithCommitUser(commitUser)
+ if err != nil {
+ t.Fatalf("Failed to create retry write builder: %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)
+ 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])
+ }
+ }
+}
+
+func TestAppendOnlyWriteMergeAndIdempotentCommit(t *testing.T) {
+ table := openCopiedTable(t, "simple_log_table")
+ const commitUser = "go-binding-multiple-writers"
+
+ 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])
+ }
+ }
}
// 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 7d943cb8..a8cec54e 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,10 @@ type paimonTableRead struct{}
type paimonPlan struct{}
type paimonRecordBatchReader struct{}
type paimonPredicate struct{}
+type paimonWriteBuilder struct{}
+type paimonTableWrite struct{}
+type paimonTableCommit struct{}
+type paimonCommitMessages struct{}
// Result types matching the C repr structs
type resultCatalogNew struct {
@@ -229,6 +270,26 @@ 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
+}
+
// 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 00000000..2e654960
--- /dev/null
+++ b/bindings/go/write.go
@@ -0,0 +1,293 @@
+/*
+ * 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.
+type WriteBuilder struct {
+ ctx context.Context
+ lib *libRef
+ inner *paimonWriteBuilder
+ overwrite bool
+ closeOnce sync.Once
+}
+
+// NewWriteBuilder creates a write builder with a 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 for merging messages across writers and for identifier retries.
+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 this builder's writers and
committers.
+func (wb *WriteBuilder) WithOverwrite() error {
+ if wb.inner == nil {
+ return ErrClosed
+ }
+ 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.
+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 identity and mode.
+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, overwrite:
wb.overwrite}, 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 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
+ }
+ 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 returns pending writes as commit messages; the writer can be
reused.
+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. It is a
+// process-local native handle and cannot be transferred between processes.
+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; 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
+ }
+ 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
+ overwrite bool
+ 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 persists messages using the builder's append or overwrite mode.
+func (tc *TableCommit) Commit(messages *CommitMessages) error {
+ operation := ffiTableCommitCommit.symbol(tc.ctx)
+ if tc.overwrite {
+ operation = ffiTableCommitOverwrite.symbol(tc.ctx)
+ }
+ return tc.withMessages(messages, operation)
+}
+
+// CommitWithIdentifier commits with a 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,
+ operation,
+ )
+}
+
+// FilterAndCommitWithIdentifier skips messages already committed under
+// 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 {
+ operation = ffiTableCommitOverwriteWithIdentifier.symbol(tc.ctx)
+ }
+ return tc.withMessagesAndIdentifier(
+ messages,
+ commitIdentifier,
+ operation,
+ )
+}
+
+// TruncateTable removes all table data.
+func (tc *TableCommit) TruncateTable() error {
+ if tc.inner == nil {
+ return ErrClosed
+ }
+ return ffiTableCommitTruncateTable.symbol(tc.ctx)(tc.inner)
+}
+
+// TruncateTableWithIdentifier truncates 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 00000000..3916b4c6
--- /dev/null
+++ b/bindings/go/write_ffi.go
@@ -0,0 +1,286 @@
+/*
+ * 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 =
newCommitFFI("paimon_table_commit_truncate_table")
+var ffiTableCommitTruncateTableWithIdentifier = newCommitWithIdentifierFFI(
+ "paimon_table_commit_truncate_table_with_identifier",
+)
+
+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,
+) *FFI[func(*paimonTableCommit, int64) error] {
+ return newFFI(ffiOpts{
+ sym: symbol,
+ rType: &ffi.TypePointer,
+ 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
+ ffiCall(
+ unsafe.Pointer(&ffiError),
+ unsafe.Pointer(&commit),
+ unsafe.Pointer(&identifier),
+ )
+ return parseError(ctx, ffiError)
+ }
+ })
+}
diff --git a/docs/src/go-binding.md b/docs/src/go-binding.md
index 41571112..3524e844 100644
--- a/docs/src/go-binding.md
+++ b/docs/src/go-binding.md
@@ -19,11 +19,15 @@ 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 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
- Go 1.22.4 or later
+- CGO enabled with a C toolchain
- Supported platforms: Linux (amd64, arm64), macOS (amd64, arm64)
## Installation
@@ -32,7 +36,8 @@ The Go integration is a binding built on top of Apache Paimon
Rust, allowing you
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`.
## Creating a Catalog
@@ -72,6 +77,54 @@ catalog, err := paimon.NewCatalog(map[string]string{
})
```
+## Writing a Table
+
+Use `NewWriteBuilder` for ordinary and fixed-bucket tables. The `arrow.Record`
+schema must match the table schema.
+
+```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)
+}
+```
+
+Call `WriteArrowBatch` multiple times before `PrepareCommit`. `WithOverwrite`
+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
+writing; merging messages does not establish ownership.
+
## 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 +337,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)