diff --git a/bindings/go/arrow_ffi.go b/bindings/go/arrow_ffi.go new file mode 100644 index 000000000..c755f347a --- /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 e9ee196f5..ed5054ce0 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 7d943cb84..a8cec54e4 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 000000000..2e6549602 --- /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 000000000..3916b4c66 --- /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 41571112d..3524e8446 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)