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 b116e665 feat(go): expose shared resource reservations (#949)
b116e665 is described below

commit b116e66547ca2551f37eab7a9b6f6110995a8fb2
Author: QuakeWang <[email protected]>
AuthorDate: Fri Sep 25 20:49:50 2026 +0800

    feat(go): expose shared resource reservations (#949)
---
 bindings/go/postpone_fixed_bucket_write.go     |  17 ++
 bindings/go/postpone_fixed_bucket_write_ffi.go |  12 ++
 bindings/go/read_builder.go                    |  29 ++++
 bindings/go/resource.go                        | 124 +++++++++++++++
 bindings/go/tests/resource_test.go             | 209 +++++++++++++++++++++++++
 bindings/go/types.go                           |  20 +++
 bindings/go/write.go                           |  17 ++
 bindings/go/write_ffi.go                       |  12 ++
 docs/src/go-binding.md                         |  25 +++
 9 files changed, 465 insertions(+)

diff --git a/bindings/go/postpone_fixed_bucket_write.go 
b/bindings/go/postpone_fixed_bucket_write.go
index 5b036d85..24ea290b 100644
--- a/bindings/go/postpone_fixed_bucket_write.go
+++ b/bindings/go/postpone_fixed_bucket_write.go
@@ -79,6 +79,23 @@ func (wb *PostponeFixedBucketWriteBuilder) Close() {
        })
 }
 
+// WithResources shares the context's reservation budget with writers from 
this builder.
+// The builder retains the context after the caller closes its handle.
+func (wb *PostponeFixedBucketWriteBuilder) WithResources(resources 
*ResourceContext) error {
+       if wb.inner == nil {
+               return ErrClosed
+       }
+       if resources == nil {
+               return errNilResourceContext
+       }
+       resources.mu.RLock()
+       defer resources.mu.RUnlock()
+       if resources.inner == nil {
+               return ErrClosed
+       }
+       return 
ffiPostponeFixedBucketWriteBuilderWithResources.symbol(wb.ctx)(wb.inner, 
resources.inner)
+}
+
 // WithOverwrite enables overwrite mode for both writers and committers created
 // by this builder.
 func (wb *PostponeFixedBucketWriteBuilder) WithOverwrite() error {
diff --git a/bindings/go/postpone_fixed_bucket_write_ffi.go 
b/bindings/go/postpone_fixed_bucket_write_ffi.go
index fff3bf0a..7a61764b 100644
--- a/bindings/go/postpone_fixed_bucket_write_ffi.go
+++ b/bindings/go/postpone_fixed_bucket_write_ffi.go
@@ -86,6 +86,18 @@ var ffiPostponeFixedBucketWriteBuilderFree = newFFI(ffiOpts{
        }
 })
 
+var ffiPostponeFixedBucketWriteBuilderWithResources = newFFI(ffiOpts{
+       sym:    "paimon_postpone_fixed_bucket_write_builder_with_resources",
+       rType:  &ffi.TypePointer,
+       aTypes: []*ffi.Type{&ffi.TypePointer, &ffi.TypePointer},
+}, func(ctx context.Context, ffiCall ffiCall) 
func(*paimonPostponeFixedBucketWriteBuilder, *paimonResourceContext) error {
+       return func(builder *paimonPostponeFixedBucketWriteBuilder, resources 
*paimonResourceContext) error {
+               var ffiError *paimonError
+               ffiCall(unsafe.Pointer(&ffiError), unsafe.Pointer(&builder), 
unsafe.Pointer(&resources))
+               return parseError(ctx, ffiError)
+       }
+})
+
 var ffiPostponeFixedBucketWriteBuilderWithOverwrite = newFFI(ffiOpts{
        sym:    "paimon_postpone_fixed_bucket_write_builder_with_overwrite",
        rType:  &ffi.TypePointer,
diff --git a/bindings/go/read_builder.go b/bindings/go/read_builder.go
index 95b2d0c6..3f52c3c7 100644
--- a/bindings/go/read_builder.go
+++ b/bindings/go/read_builder.go
@@ -45,6 +45,23 @@ func (rb *ReadBuilder) Close() {
        })
 }
 
+// WithResources shares the context's reservation budget with reads from this 
builder.
+// The builder retains the context after the caller closes its handle.
+func (rb *ReadBuilder) WithResources(resources *ResourceContext) error {
+       if rb.inner == nil {
+               return ErrClosed
+       }
+       if resources == nil {
+               return errNilResourceContext
+       }
+       resources.mu.RLock()
+       defer resources.mu.RUnlock()
+       if resources.inner == nil {
+               return ErrClosed
+       }
+       return ffiReadBuilderWithResources.symbol(rb.ctx)(rb.inner, 
resources.inner)
+}
+
 // WithProjection sets column projection by name. Output order follows the
 // caller-specified order. A name that matches no schema column under any case
 // sensitivity is rejected immediately; case-dependent errors and duplicate 
names
@@ -196,6 +213,18 @@ var ffiReadBuilderWithProjection = newFFI(ffiOpts{
        }
 })
 
+var ffiReadBuilderWithResources = newFFI(ffiOpts{
+       sym:    "paimon_read_builder_with_resources",
+       rType:  &ffi.TypePointer,
+       aTypes: []*ffi.Type{&ffi.TypePointer, &ffi.TypePointer},
+}, func(ctx context.Context, ffiCall ffiCall) func(*paimonReadBuilder, 
*paimonResourceContext) error {
+       return func(builder *paimonReadBuilder, resources 
*paimonResourceContext) error {
+               var ffiError *paimonError
+               ffiCall(unsafe.Pointer(&ffiError), unsafe.Pointer(&builder), 
unsafe.Pointer(&resources))
+               return parseError(ctx, ffiError)
+       }
+})
+
 // The trailing `bool` is passed as a 1-byte integer written through boolByte; 
see
 // that function for why.
 var ffiReadBuilderWithCaseSensitive = newFFI(ffiOpts{
diff --git a/bindings/go/resource.go b/bindings/go/resource.go
new file mode 100644
index 00000000..539b05bd
--- /dev/null
+++ b/bindings/go/resource.go
@@ -0,0 +1,124 @@
+/*
+ * 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"
+       "errors"
+       "sync"
+       "unsafe"
+
+       "github.com/jupiterrider/ffi"
+)
+
+// ResourceContext shares a memory reservation budget across readers and 
writers.
+// Reservations are accounting estimates, not a process memory limit.
+type ResourceContext struct {
+       ctx   context.Context
+       lib   *libRef
+       inner *paimonResourceContext
+       mu    sync.RWMutex
+}
+
+// ResourceMetrics reports current and peak reserved bytes.
+type ResourceMetrics struct {
+       ReservedMemoryBytes     uintptr
+       PeakReservedMemoryBytes uintptr
+}
+
+// NewResourceContext creates a shared reservation budget in bytes.
+// A zero limit rejects nonempty reservations.
+func NewResourceContext(memoryLimitBytes uintptr) (*ResourceContext, error) {
+       ctx, lib, err := ensureLoaded()
+       if err != nil {
+               return nil, err
+       }
+       inner, err := ffiResourceContextCreate.symbol(ctx)(memoryLimitBytes)
+       if err != nil {
+               return nil, err
+       }
+       lib.acquire()
+       return &ResourceContext{ctx: ctx, lib: lib, inner: inner}, nil
+}
+
+// Metrics reads the current and peak reservations. The counters are sampled 
independently.
+func (r *ResourceContext) Metrics() (ResourceMetrics, error) {
+       r.mu.RLock()
+       defer r.mu.RUnlock()
+       if r.inner == nil {
+               return ResourceMetrics{}, ErrClosed
+       }
+       return ffiResourceContextMetrics.symbol(r.ctx)(r.inner)
+}
+
+// Close releases this handle. Builders retain their own resource context 
clones.
+func (r *ResourceContext) Close() {
+       r.mu.Lock()
+       defer r.mu.Unlock()
+       if r.inner == nil {
+               return
+       }
+       ffiResourceContextFree.symbol(r.ctx)(r.inner)
+       r.inner = nil
+       r.lib.release()
+}
+
+var errNilResourceContext = errors.New("paimon: resource context must not be 
nil")
+
+var ffiResourceContextCreate = newFFI(ffiOpts{
+       sym:    "paimon_resource_context_create",
+       rType:  &typeResultResourceContext,
+       aTypes: []*ffi.Type{&ffi.TypePointer},
+}, func(ctx context.Context, ffiCall ffiCall) func(uintptr) 
(*paimonResourceContext, error) {
+       return func(memoryLimitBytes uintptr) (*paimonResourceContext, error) {
+               var result resultResourceContext
+               ffiCall(unsafe.Pointer(&result), 
unsafe.Pointer(&memoryLimitBytes))
+               if result.error != nil {
+                       return nil, parseError(ctx, result.error)
+               }
+               return result.context, nil
+       }
+})
+
+var ffiResourceContextMetrics = newFFI(ffiOpts{
+       sym:    "paimon_resource_context_metrics",
+       rType:  &ffi.TypePointer,
+       aTypes: []*ffi.Type{&ffi.TypePointer, &ffi.TypePointer},
+}, func(ctx context.Context, ffiCall ffiCall) func(*paimonResourceContext) 
(ResourceMetrics, error) {
+       return func(resource *paimonResourceContext) (ResourceMetrics, error) {
+               var metrics paimonResourceMetrics
+               metricsPtr := &metrics
+               var ffiError *paimonError
+               ffiCall(unsafe.Pointer(&ffiError), unsafe.Pointer(&resource), 
unsafe.Pointer(&metricsPtr))
+               if err := parseError(ctx, ffiError); err != nil {
+                       return ResourceMetrics{}, err
+               }
+               return ResourceMetrics{metrics.reservedMemoryBytes, 
metrics.peakReservedMemoryBytes}, nil
+       }
+})
+
+var ffiResourceContextFree = newFFI(ffiOpts{
+       sym:    "paimon_resource_context_free",
+       rType:  &ffi.TypeVoid,
+       aTypes: []*ffi.Type{&ffi.TypePointer},
+}, func(_ context.Context, ffiCall ffiCall) func(*paimonResourceContext) {
+       return func(resource *paimonResourceContext) {
+               ffiCall(nil, unsafe.Pointer(&resource))
+       }
+})
diff --git a/bindings/go/tests/resource_test.go 
b/bindings/go/tests/resource_test.go
new file mode 100644
index 00000000..83c6e0cf
--- /dev/null
+++ b/bindings/go/tests/resource_test.go
@@ -0,0 +1,209 @@
+/*
+ * 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_test
+
+import (
+       "errors"
+       "os"
+       "path/filepath"
+       "strings"
+       "testing"
+
+       paimon "github.com/apache/paimon-rust/bindings/go"
+)
+
+func requireResourceExhausted(t *testing.T, err error) {
+       t.Helper()
+       var nativeErr *paimon.Error
+       if !errors.As(err, &nativeErr) || nativeErr.Code() != 
paimon.CodeResourceExhausted {
+               t.Fatalf("expected ResourceExhausted, got %v", err)
+       }
+}
+
+func TestResourceContextMetricsAndClose(t *testing.T) {
+       resources, err := paimon.NewResourceContext(0)
+       if err != nil {
+               t.Fatal(err)
+       }
+       metrics, err := resources.Metrics()
+       if err != nil {
+               t.Fatal(err)
+       }
+       if metrics.ReservedMemoryBytes != 0 || metrics.PeakReservedMemoryBytes 
!= 0 {
+               t.Fatalf("unexpected initial metrics: %+v", metrics)
+       }
+       resources.Close()
+       resources.Close()
+       if _, err := resources.Metrics(); !errors.Is(err, paimon.ErrClosed) {
+               t.Fatalf("expected ErrClosed after Close, got %v", err)
+       }
+}
+
+func TestReadBuilderWithResourceBudget(t *testing.T) {
+       source := filepath.Join("testdata", "map_blob_table")
+       warehouse := t.TempDir()
+       if err := copyDirectory(source, filepath.Join(warehouse, "default.db", 
"map_blob_table")); err != nil {
+               t.Fatal(err)
+       }
+       table := openTableAt(t, warehouse, "map_blob_table")
+       resources, err := paimon.NewResourceContext(0)
+       if err != nil {
+               t.Fatal(err)
+       }
+       defer resources.Close()
+       builder, err := 
table.NewReadBuilderWithOptions(map[string]string{"blob-as-descriptor": "true"})
+       if err != nil {
+               t.Fatal(err)
+       }
+       defer builder.Close()
+       if err := builder.WithResources(resources); err != nil {
+               t.Fatal(err)
+       }
+       resources.Close()
+       if err := builder.WithResources(resources); !errors.Is(err, 
paimon.ErrClosed) {
+               t.Fatalf("expected ErrClosed for a closed resource context, got 
%v", err)
+       }
+       scan, err := builder.NewScan()
+       if err != nil {
+               t.Fatal(err)
+       }
+       defer scan.Close()
+       plan, err := scan.Plan()
+       if err != nil {
+               t.Fatal(err)
+       }
+       defer plan.Close()
+       if len(plan.Splits()) == 0 {
+               t.Fatal("expected a nonempty read plan")
+       }
+       read, err := builder.NewRead()
+       if err != nil {
+               t.Fatal(err)
+       }
+       defer read.Close()
+       reader, err := read.NewRecordBatchReader(plan.Splits())
+       if err != nil {
+               requireResourceExhausted(t, err)
+               return
+       }
+       defer reader.Close()
+       record, err := reader.NextRecord()
+       if record != nil {
+               record.Release()
+       }
+       requireResourceExhausted(t, err)
+}
+
+func TestWriteBuildersShareResourceBudget(t *testing.T) {
+       table := openCopiedTestTable(t)
+       resources, err := paimon.NewResourceContext(1_000_000)
+       if err != nil {
+               t.Fatal(err)
+       }
+       defer resources.Close()
+       builders := make([]*paimon.WriteBuilder, 2)
+       writers := make([]*paimon.TableWrite, 2)
+       for i := range builders {
+               builders[i], err = table.NewWriteBuilder()
+               if err != nil {
+                       t.Fatal(err)
+               }
+               defer builders[i].Close()
+               if err := builders[i].WithResources(resources); err != nil {
+                       t.Fatal(err)
+               }
+               writers[i], err = builders[i].NewWrite()
+               if err != nil {
+                       t.Fatal(err)
+               }
+               defer writers[i].Close()
+       }
+
+       value := strings.Repeat("x", 600_000)
+       first := makeRecord(t, []row{{1, value}})
+       err = writers[0].WriteArrowBatch(first)
+       first.Release()
+       if err != nil {
+               t.Fatal(err)
+       }
+       metrics, err := resources.Metrics()
+       if err != nil {
+               t.Fatal(err)
+       }
+       if metrics.ReservedMemoryBytes == 0 || metrics.ReservedMemoryBytes > 
1_000_000 {
+               t.Fatalf("unexpected reserved bytes: %+v", metrics)
+       }
+       second := makeRecord(t, []row{{2, value}})
+       err = writers[1].WriteArrowBatch(second)
+       second.Release()
+       requireResourceExhausted(t, err)
+       writers[0].Close()
+       metrics, err = resources.Metrics()
+       if err != nil {
+               t.Fatal(err)
+       }
+       if metrics.ReservedMemoryBytes != 0 || metrics.PeakReservedMemoryBytes 
== 0 {
+               t.Fatalf("unexpected metrics after releasing writers: %+v", 
metrics)
+       }
+}
+
+func TestPostponeWriteBuilderWithResourceBudget(t *testing.T) {
+       warehouse := testWarehouse()
+       if _, err := os.Stat(filepath.Join(warehouse, "default.db", 
"postpone_fixed_bucket_pk_table")); os.IsNotExist(err) {
+               t.Skip("postpone fixed-bucket test table is unavailable")
+       }
+       table := openCopiedTable(t, "postpone_fixed_bucket_pk_table")
+       resources, err := paimon.NewResourceContext(0)
+       if err != nil {
+               t.Fatal(err)
+       }
+       defer resources.Close()
+       builder, err := table.NewPostponeFixedBucketWriteBuilder()
+       if err != nil {
+               t.Fatal(err)
+       }
+       defer builder.Close()
+       if err := builder.WithResources(resources); err != nil {
+               t.Fatal(err)
+       }
+       resources.Close()
+       plan := makePartitionedBucketPlan(t, []string{"2026-08-14"}, 1)
+       err = builder.WithBucketPlan(plan)
+       plan.Release()
+       if err != nil {
+               t.Fatal(err)
+       }
+       write, err := builder.NewWrite()
+       if err != nil {
+               t.Fatal(err)
+       }
+       defer write.Close()
+       record := makePartitionedRecord(t, partitionedRow{4, "dave", 
"2026-08-14"})
+       err = write.WriteArrowBatch(record)
+       record.Release()
+       if err == nil {
+               messages, prepareErr := write.PrepareCommit()
+               if messages != nil {
+                       messages.Close()
+               }
+               err = prepareErr
+       }
+       requireResourceExhausted(t, err)
+}
diff --git a/bindings/go/types.go b/bindings/go/types.go
index 5412bcd6..21be0b79 100644
--- a/bindings/go/types.go
+++ b/bindings/go/types.go
@@ -29,6 +29,15 @@ import (
 
 // FFI type definitions mirroring C repr structs from paimon-c.
 var (
+       typeResultResourceContext = ffi.Type{
+               Type: ffi.Struct,
+               Elements: &[]*ffi.Type{
+                       &ffi.TypePointer,
+                       &ffi.TypePointer,
+                       nil,
+               }[0],
+       }
+
        typeResultBlobReader = ffi.Type{
                Type: ffi.Struct,
                Elements: &[]*ffi.Type{
@@ -350,6 +359,7 @@ type paimonBlobReader struct{}
 type paimonBlobStream struct{}
 type paimonIdentifier struct{}
 type paimonTable struct{}
+type paimonResourceContext struct{}
 type paimonReadBuilder struct{}
 type paimonTableScan struct{}
 type paimonTableRead struct{}
@@ -366,6 +376,16 @@ type paimonPostponeFixedBucketTableCommit struct{}
 type paimonPostponeFixedBucketCommitMessages struct{}
 
 // Result types matching the C repr structs
+type resultResourceContext struct {
+       context *paimonResourceContext
+       error   *paimonError
+}
+
+type paimonResourceMetrics struct {
+       reservedMemoryBytes     uintptr
+       peakReservedMemoryBytes uintptr
+}
+
 type resultCatalogNew struct {
        catalog *paimonCatalog
        error   *paimonError
diff --git a/bindings/go/write.go b/bindings/go/write.go
index 2e654960..c5b793ab 100644
--- a/bindings/go/write.go
+++ b/bindings/go/write.go
@@ -72,6 +72,23 @@ func (wb *WriteBuilder) Close() {
        })
 }
 
+// WithResources shares the context's reservation budget with writers from 
this builder.
+// The builder retains the context after the caller closes its handle.
+func (wb *WriteBuilder) WithResources(resources *ResourceContext) error {
+       if wb.inner == nil {
+               return ErrClosed
+       }
+       if resources == nil {
+               return errNilResourceContext
+       }
+       resources.mu.RLock()
+       defer resources.mu.RUnlock()
+       if resources.inner == nil {
+               return ErrClosed
+       }
+       return ffiWriteBuilderWithResources.symbol(wb.ctx)(wb.inner, 
resources.inner)
+}
+
 // WithOverwrite enables overwrite mode for this builder's writers and 
committers.
 func (wb *WriteBuilder) WithOverwrite() error {
        if wb.inner == nil {
diff --git a/bindings/go/write_ffi.go b/bindings/go/write_ffi.go
index 3916b4c6..bcd046fc 100644
--- a/bindings/go/write_ffi.go
+++ b/bindings/go/write_ffi.go
@@ -74,6 +74,18 @@ var ffiWriteBuilderFree = newFFI(ffiOpts{
        }
 })
 
+var ffiWriteBuilderWithResources = newFFI(ffiOpts{
+       sym:    "paimon_write_builder_with_resources",
+       rType:  &ffi.TypePointer,
+       aTypes: []*ffi.Type{&ffi.TypePointer, &ffi.TypePointer},
+}, func(ctx context.Context, ffiCall ffiCall) func(*paimonWriteBuilder, 
*paimonResourceContext) error {
+       return func(builder *paimonWriteBuilder, resources 
*paimonResourceContext) error {
+               var ffiError *paimonError
+               ffiCall(unsafe.Pointer(&ffiError), unsafe.Pointer(&builder), 
unsafe.Pointer(&resources))
+               return parseError(ctx, ffiError)
+       }
+})
+
 var ffiWriteBuilderWithOverwrite = newFFI(ffiOpts{
        sym:    "paimon_write_builder_with_overwrite",
        rType:  &ffi.TypePointer,
diff --git a/docs/src/go-binding.md b/docs/src/go-binding.md
index a1fd1695..203db6bd 100644
--- a/docs/src/go-binding.md
+++ b/docs/src/go-binding.md
@@ -696,6 +696,31 @@ pred, _ := pb.Eq("amount", paimon.NewDecimal(12345, 10, 2))
 pred, _ := pb.Eq("ts", paimon.Timestamp{Millis: 1700000000000, Nanos: 0})
 ```
 
+## Shared Memory Reservations
+
+Create a `ResourceContext` to share one reservation budget across read and 
write
+builders. Attach it before creating a reader or writer. The builder keeps its 
own
+reference, so the original handle may be closed after attachment.
+
+```go
+resources, err := paimon.NewResourceContext(256 * 1024 * 1024)
+if err != nil { log.Fatal(err) }
+defer resources.Close()
+
+if err := readBuilder.WithResources(resources); err != nil { log.Fatal(err) }
+if err := writeBuilder.WithResources(resources); err != nil { log.Fatal(err) }
+
+metrics, err := resources.Metrics()
+if err != nil { log.Fatal(err) }
+fmt.Println(metrics.ReservedMemoryBytes, metrics.PeakReservedMemoryBytes)
+```
+
+`PostponeFixedBucketWriteBuilder.WithResources` accepts the same context. When 
a
+reservation exceeds the limit, the operation returns `CodeResourceExhausted`.
+These counters track estimated reservations held by Paimon, not total process
+memory. Current reservations fall as readers and writers release them; the peak
+remains available while the context handle is open.
+
 ## Resource Management
 
 Paimon objects with a `Close` method hold native resources and must be closed.

Reply via email to