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.