Copilot commented on code in PR #722: URL: https://github.com/apache/paimon-rust/pull/722#discussion_r3804262063
########## 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 + } +}) Review Comment: Several FFI bindings specify `rType` as the *standard* result types (`typeResultWriteBuilder`, `typeResultPrepareCommit`, etc.) while the call sites use the new postpone fixed-bucket result structs (`resultPostponeFixedBucketWriteBuilder`, `resultPostponeFixedBucketPrepareCommit`, etc.). This is an ABI mismatch risk (wrong size/layout) that can cause memory corruption or incorrect reads. Define and use dedicated ffi.Type descriptors for the new postpone fixed-bucket result structs (e.g., `typeResultPostponeFixedBucketWriteBuilder`, `typeResultPostponeFixedBucketTableWrite`, `typeResultPostponeFixedBucketTableCommit`, `typeResultPostponeFixedBucketPrepareCommit`) and reference them in each `newFFI` opts. ########## bindings/go/postpone_fixed_bucket_write.go: ########## @@ -0,0 +1,329 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package paimon + +import ( + "context" + "sync" + "unsafe" + + "github.com/apache/arrow-go/v18/arrow" +) + +// PostponeFixedBucketWriteBuilder creates fixed-bucket writers for bucket=-2 +// tables. A resolved bucket plan must be supplied before NewWrite. +type PostponeFixedBucketWriteBuilder struct { + ctx context.Context + lib *libRef + inner *paimonPostponeFixedBucketWriteBuilder + closeOnce sync.Once +} + +// NewPostponeFixedBucketWriteBuilder creates an explicitly selected +// fixed-bucket builder for a postpone table. +func (t *Table) NewPostponeFixedBucketWriteBuilder() (*PostponeFixedBucketWriteBuilder, error) { + if t.inner == nil { + return nil, ErrClosed + } + inner, err := ffiTableNewPostponeFixedBucketWriteBuilder.symbol(t.ctx)(t.inner) + if err != nil { + return nil, err + } + t.lib.acquire() + return &PostponeFixedBucketWriteBuilder{ctx: t.ctx, lib: t.lib, inner: inner}, nil +} + +// NewPostponeFixedBucketWriteBuilderWithCommitUser creates a fixed-bucket +// builder with a stable commit identity. +func (t *Table) NewPostponeFixedBucketWriteBuilderWithCommitUser( + commitUser string, +) (*PostponeFixedBucketWriteBuilder, error) { + if t.inner == nil { + return nil, ErrClosed + } + inner, err := ffiTableNewPostponeFixedBucketWriteBuilderWithCommitUser.symbol(t.ctx)( + t.inner, + commitUser, + ) + if err != nil { + return nil, err + } + t.lib.acquire() + return &PostponeFixedBucketWriteBuilder{ctx: t.ctx, lib: t.lib, inner: inner}, nil +} + +// Close releases the builder resources. Safe to call multiple times. +func (wb *PostponeFixedBucketWriteBuilder) Close() { + wb.closeOnce.Do(func() { + ffiPostponeFixedBucketWriteBuilderFree.symbol(wb.ctx)(wb.inner) + wb.inner = nil + wb.lib.release() + }) +} + +// WithOverwrite enables overwrite mode for both writers and committers created +// by this builder. +func (wb *PostponeFixedBucketWriteBuilder) WithOverwrite() error { + if wb.inner == nil { + return ErrClosed + } + return ffiPostponeFixedBucketWriteBuilderWithOverwrite.symbol(wb.ctx)(wb.inner) +} + +// WithBucketPlan sets a resolved partition-to-bucket-count plan. The plan must +// contain the table partition columns followed by a non-null Int32 +// total_buckets column. The caller retains ownership of plan. +func (wb *PostponeFixedBucketWriteBuilder) WithBucketPlan(plan arrow.Record) error { + if wb.inner == nil { + return ErrClosed + } + return withOwnedArrowRecord( + plan, + "paimon: bucket plan must not be nil", + func(array, schema unsafe.Pointer) error { + return ffiPostponeFixedBucketWriteBuilderWithBucketPlan.symbol(wb.ctx)( + wb.inner, + array, + schema, + ) + }, + ) +} + +// NewWrite creates a fixed-bucket writer. WithBucketPlan must be called first. +func (wb *PostponeFixedBucketWriteBuilder) NewWrite() (*PostponeFixedBucketTableWrite, error) { + if wb.inner == nil { + return nil, ErrClosed + } + inner, err := ffiPostponeFixedBucketWriteBuilderNewWrite.symbol(wb.ctx)(wb.inner) + if err != nil { + return nil, err + } + wb.lib.acquire() + return &PostponeFixedBucketTableWrite{ctx: wb.ctx, lib: wb.lib, inner: inner}, nil +} + +// NewCommit creates a committer using this builder's commit identity and mode. +func (wb *PostponeFixedBucketWriteBuilder) NewCommit() (*PostponeFixedBucketTableCommit, error) { + if wb.inner == nil { + return nil, ErrClosed + } + inner, err := ffiPostponeFixedBucketWriteBuilderNewCommit.symbol(wb.ctx)(wb.inner) + if err != nil { + return nil, err + } + wb.lib.acquire() + return &PostponeFixedBucketTableCommit{ctx: wb.ctx, lib: wb.lib, inner: inner}, nil +} + +// PostponeFixedBucketTableWrite writes rows according to a resolved bucket plan. +type PostponeFixedBucketTableWrite struct { + ctx context.Context + lib *libRef + inner *paimonPostponeFixedBucketTableWrite + closeOnce sync.Once +} + +// Close releases the writer resources. Safe to call multiple times. +func (tw *PostponeFixedBucketTableWrite) Close() { + tw.closeOnce.Do(func() { + ffiPostponeFixedBucketTableWriteFree.symbol(tw.ctx)(tw.inner) + tw.inner = nil + tw.lib.release() + }) +} + +// WriteArrowBatch writes one Arrow record batch. The caller retains ownership. +func (tw *PostponeFixedBucketTableWrite) WriteArrowBatch(record arrow.Record) error { + if tw.inner == nil { + return ErrClosed + } + return withOwnedArrowRecord( + record, + "paimon: record batch must not be nil", + func(array, schema unsafe.Pointer) error { + return ffiPostponeFixedBucketTableWriteWriteArrowBatch.symbol(tw.ctx)( + tw.inner, + array, + schema, + ) + }, + ) +} + +// PrepareCommit closes current writers and returns fixed-bucket messages. +// The writer is single-use; create a new writer for the next batch. +func (tw *PostponeFixedBucketTableWrite) PrepareCommit() (*PostponeFixedBucketCommitMessages, error) { + if tw.inner == nil { + return nil, ErrClosed + } + inner, err := ffiPostponeFixedBucketTableWritePrepareCommit.symbol(tw.ctx)(tw.inner) + if err != nil { + return nil, err + } + tw.lib.acquire() Review Comment: The docstring states `PrepareCommit` \"closes\" the writer and that the writer is single-use, but the method does not invalidate `tw.inner` (or otherwise mark the writer closed). This makes it easy for callers (including the doc/test pattern that defers `Close()`) to accidentally double-close/free, or to call `WriteArrowBatch` after `PrepareCommit` despite the 'single-use' contract. Consider making `PrepareCommit` explicitly consume the writer: on success set `tw.inner = nil` and release the writer’s lib ref (or internally call `tw.Close()` after the FFI call), so subsequent operations return `ErrClosed` and deferred closes are safe no-ops. ########## 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 + } +}) Review Comment: Several FFI bindings specify `rType` as the *standard* result types (`typeResultWriteBuilder`, `typeResultPrepareCommit`, etc.) while the call sites use the new postpone fixed-bucket result structs (`resultPostponeFixedBucketWriteBuilder`, `resultPostponeFixedBucketPrepareCommit`, etc.). This is an ABI mismatch risk (wrong size/layout) that can cause memory corruption or incorrect reads. Define and use dedicated ffi.Type descriptors for the new postpone fixed-bucket result structs (e.g., `typeResultPostponeFixedBucketWriteBuilder`, `typeResultPostponeFixedBucketTableWrite`, `typeResultPostponeFixedBucketTableCommit`, `typeResultPostponeFixedBucketPrepareCommit`) and reference them in each `newFFI` opts. ########## bindings/go/postpone_fixed_bucket_write.go: ########## @@ -0,0 +1,329 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package paimon + +import ( + "context" + "sync" + "unsafe" + + "github.com/apache/arrow-go/v18/arrow" +) + +// PostponeFixedBucketWriteBuilder creates fixed-bucket writers for bucket=-2 +// tables. A resolved bucket plan must be supplied before NewWrite. +type PostponeFixedBucketWriteBuilder struct { + ctx context.Context + lib *libRef + inner *paimonPostponeFixedBucketWriteBuilder + closeOnce sync.Once +} + +// NewPostponeFixedBucketWriteBuilder creates an explicitly selected +// fixed-bucket builder for a postpone table. +func (t *Table) NewPostponeFixedBucketWriteBuilder() (*PostponeFixedBucketWriteBuilder, error) { + if t.inner == nil { + return nil, ErrClosed + } + inner, err := ffiTableNewPostponeFixedBucketWriteBuilder.symbol(t.ctx)(t.inner) + if err != nil { + return nil, err + } + t.lib.acquire() + return &PostponeFixedBucketWriteBuilder{ctx: t.ctx, lib: t.lib, inner: inner}, nil +} + +// NewPostponeFixedBucketWriteBuilderWithCommitUser creates a fixed-bucket +// builder with a stable commit identity. +func (t *Table) NewPostponeFixedBucketWriteBuilderWithCommitUser( + commitUser string, +) (*PostponeFixedBucketWriteBuilder, error) { + if t.inner == nil { + return nil, ErrClosed + } + inner, err := ffiTableNewPostponeFixedBucketWriteBuilderWithCommitUser.symbol(t.ctx)( + t.inner, + commitUser, + ) + if err != nil { + return nil, err + } + t.lib.acquire() + return &PostponeFixedBucketWriteBuilder{ctx: t.ctx, lib: t.lib, inner: inner}, nil +} + +// Close releases the builder resources. Safe to call multiple times. +func (wb *PostponeFixedBucketWriteBuilder) Close() { + wb.closeOnce.Do(func() { + ffiPostponeFixedBucketWriteBuilderFree.symbol(wb.ctx)(wb.inner) + wb.inner = nil + wb.lib.release() + }) +} + +// WithOverwrite enables overwrite mode for both writers and committers created +// by this builder. +func (wb *PostponeFixedBucketWriteBuilder) WithOverwrite() error { + if wb.inner == nil { + return ErrClosed + } + return ffiPostponeFixedBucketWriteBuilderWithOverwrite.symbol(wb.ctx)(wb.inner) +} + +// WithBucketPlan sets a resolved partition-to-bucket-count plan. The plan must +// contain the table partition columns followed by a non-null Int32 +// total_buckets column. The caller retains ownership of plan. +func (wb *PostponeFixedBucketWriteBuilder) WithBucketPlan(plan arrow.Record) error { + if wb.inner == nil { + return ErrClosed + } + return withOwnedArrowRecord( + plan, + "paimon: bucket plan must not be nil", + func(array, schema unsafe.Pointer) error { + return ffiPostponeFixedBucketWriteBuilderWithBucketPlan.symbol(wb.ctx)( + wb.inner, + array, + schema, + ) + }, + ) +} + +// NewWrite creates a fixed-bucket writer. WithBucketPlan must be called first. +func (wb *PostponeFixedBucketWriteBuilder) NewWrite() (*PostponeFixedBucketTableWrite, error) { + if wb.inner == nil { + return nil, ErrClosed + } + inner, err := ffiPostponeFixedBucketWriteBuilderNewWrite.symbol(wb.ctx)(wb.inner) + if err != nil { + return nil, err + } + wb.lib.acquire() + return &PostponeFixedBucketTableWrite{ctx: wb.ctx, lib: wb.lib, inner: inner}, nil +} + +// NewCommit creates a committer using this builder's commit identity and mode. +func (wb *PostponeFixedBucketWriteBuilder) NewCommit() (*PostponeFixedBucketTableCommit, error) { + if wb.inner == nil { + return nil, ErrClosed + } + inner, err := ffiPostponeFixedBucketWriteBuilderNewCommit.symbol(wb.ctx)(wb.inner) + if err != nil { + return nil, err + } + wb.lib.acquire() + return &PostponeFixedBucketTableCommit{ctx: wb.ctx, lib: wb.lib, inner: inner}, nil +} + +// PostponeFixedBucketTableWrite writes rows according to a resolved bucket plan. +type PostponeFixedBucketTableWrite struct { + ctx context.Context + lib *libRef + inner *paimonPostponeFixedBucketTableWrite + closeOnce sync.Once +} + +// Close releases the writer resources. Safe to call multiple times. +func (tw *PostponeFixedBucketTableWrite) Close() { + tw.closeOnce.Do(func() { + ffiPostponeFixedBucketTableWriteFree.symbol(tw.ctx)(tw.inner) + tw.inner = nil + tw.lib.release() + }) +} + +// WriteArrowBatch writes one Arrow record batch. The caller retains ownership. +func (tw *PostponeFixedBucketTableWrite) WriteArrowBatch(record arrow.Record) error { + if tw.inner == nil { + return ErrClosed + } + return withOwnedArrowRecord( + record, + "paimon: record batch must not be nil", + func(array, schema unsafe.Pointer) error { + return ffiPostponeFixedBucketTableWriteWriteArrowBatch.symbol(tw.ctx)( + tw.inner, + array, + schema, + ) + }, + ) +} + +// PrepareCommit closes current writers and returns fixed-bucket messages. +// The writer is single-use; create a new writer for the next batch. +func (tw *PostponeFixedBucketTableWrite) PrepareCommit() (*PostponeFixedBucketCommitMessages, error) { + if tw.inner == nil { + return nil, ErrClosed + } + inner, err := ffiPostponeFixedBucketTableWritePrepareCommit.symbol(tw.ctx)(tw.inner) + if err != nil { + return nil, err + } + tw.lib.acquire() + return &PostponeFixedBucketCommitMessages{ctx: tw.ctx, lib: tw.lib, inner: inner}, nil +} + +// PostponeFixedBucketCommitMessages contains files produced by fixed-bucket +// writers. It is a process-local native handle and cannot be transferred +// between processes or passed to a standard TableCommit. +type PostponeFixedBucketCommitMessages struct { + ctx context.Context + lib *libRef + inner *paimonPostponeFixedBucketCommitMessages + closeOnce sync.Once +} + +// Close releases the messages. Safe to call multiple times. +func (m *PostponeFixedBucketCommitMessages) Close() { + m.closeOnce.Do(func() { + ffiPostponeFixedBucketCommitMessagesFree.symbol(m.ctx)(m.inner) + m.inner = nil + m.lib.release() + }) +} + +// Merge appends a copy of source's messages. Both handles must belong to the +// same process, and both builders must use the same table, commit user, and +// overwrite mode. +func (m *PostponeFixedBucketCommitMessages) Merge( + source *PostponeFixedBucketCommitMessages, +) error { + if m.inner == nil { + return ErrClosed + } + if source == nil || source.inner == nil { + return ErrClosed + } + return ffiPostponeFixedBucketCommitMessagesMerge.symbol(m.ctx)(m.inner, source.inner) +} Review Comment: Passing `nil` for `source` currently returns `ErrClosed`, which is misleading (the handle isn’t 'closed', it’s missing). Returning a more specific error (e.g., `errors.New(\"paimon: source messages must not be nil\")`) would make debugging easier and align better with the explicit error strings used elsewhere (e.g., the `withOwnedArrowRecord` messages). ########## docs/src/go-binding.md: ########## @@ -125,6 +125,63 @@ process-local and cannot be sent to another process. For primary-key fixed-bucket tables, assign each `(partition, bucket)` to one writer before writing; merging messages does not establish ownership. +### Postpone Fixed-Bucket Writes + +For a `bucket = -2` table, the ordinary builder writes postpone files. To write +real buckets directly, use the dedicated builder with a plan mapping each +partition to its bucket count: + +```go +// Plan schema: partition keys in order, then a non-null Int32 total_buckets. +schema := arrow.NewSchema([]arrow.Field{ + {Name: "dt", Type: arrow.BinaryTypes.String, Nullable: true}, + {Name: "total_buckets", Type: arrow.PrimitiveTypes.Int32, Nullable: false}, +}, nil) +rb := array.NewRecordBuilder(memory.DefaultAllocator, schema) +defer rb.Release() +rb.Field(0).(*array.StringBuilder).AppendValues([]string{"2026-08-14", "2026-08-15"}, nil) +rb.Field(1).(*array.Int32Builder).AppendValues([]int32{1, 1}, nil) +plan := rb.NewRecord() +defer plan.Release() + +builder, err := table.NewPostponeFixedBucketWriteBuilderWithCommitUser("job-1") +if err != nil { + log.Fatal(err) +} +defer builder.Close() +if err := builder.WithBucketPlan(plan); err != nil { + log.Fatal(err) +} +writer, err := builder.NewWrite() // fails without a plan Review Comment: This inline comment is slightly ambiguous because the example has already set the plan right above. Consider rewording to clarify it’s a precondition (e.g., \"requires a plan; errors if WithBucketPlan wasn't called\") so readers don’t misread it as saying this call should fail in the example. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
