This is an automated email from the ASF dual-hosted git repository.

jrmccluskey pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git


The following commit(s) were added to refs/heads/master by this push:
     new 7a5f89199d5 [Claude] Plumb through pipeline-level resource hints to 
cross-lanugage transforms (#40379)
7a5f89199d5 is described below

commit 7a5f89199d5ec6512cd293846f2a3ea5ef914aca
Author: Jack McCluskey <[email protected]>
AuthorDate: Thu Oct 8 11:06:48 2026 -0400

    [Claude] Plumb through pipeline-level resource hints to cross-lanugage 
transforms (#40379)
---
 CHANGES.md                                         |   1 +
 sdks/go/pkg/beam/core/runtime/xlangx/expand.go     |  47 ++++++-
 .../go/pkg/beam/core/runtime/xlangx/expand_test.go | 148 +++++++++++++++++++++
 sdks/go/pkg/beam/options/jobopts/options.go        |   2 +
 sdks/go/pkg/beam/options/jobopts/options_test.go   |  24 +++-
 sdks/go/pkg/beam/options/resource/hint.go          |  39 ++++++
 sdks/go/pkg/beam/options/resource/hint_test.go     |  45 +++++++
 7 files changed, 302 insertions(+), 4 deletions(-)

diff --git a/CHANGES.md b/CHANGES.md
index 5c5370006d6..997ee9adb49 100644
--- a/CHANGES.md
+++ b/CHANGES.md
@@ -71,6 +71,7 @@
 * (Go) Added `wait.On`, which delays each input window until the corresponding 
windows in its signal PCollections have closed 
([#39909](https://github.com/apache/beam/issues/39909)).
 * (Python) Expanded the SDK worker heap dump 
(`--experiments=enable_heap_dump`) with process RSS, CPython allocator/GC 
stats, and glibc `mallinfo2` native-heap/fragmentation stats to help 
distinguish native-heap from Python-object memory growth 
([#39244](https://github.com/apache/beam/issues/39244)).
 * The `disableCounterMetrics`, `disableStringSetMetrics` and 
`disableBoundedTrieMetrics` experiments are now honored by the Python SDK, as 
they already were in Java (Python) 
([#38746](https://github.com/apache/beam/issues/38746)).
+* (Go) Pipeline-level `--resource_hints` are now forwarded to expansion 
services, so they apply to cross-language transforms as they do in the Java and 
Python SDKs. The `max_active_bundles_per_worker` hint name is now also accepted 
by `--resource_hints` ([#23893](https://github.com/apache/beam/issues/23893)).
 * ReadFromBigQuery now supports Lakehouse runtime catalog (BigLake metastore) 
tables with `method=DIRECT_READ`, using 4-part 
`project.catalog.namespace.table` identifiers. Previously 
`project:catalog.namespace.table` was silently mis-parsed and tables that 
report no `numBytes` failed to split (Python) 
([#39597](https://github.com/apache/beam/issues/39597)).
 
 ## Breaking Changes
diff --git a/sdks/go/pkg/beam/core/runtime/xlangx/expand.go 
b/sdks/go/pkg/beam/core/runtime/xlangx/expand.go
index 94dda75e805..de7aae3341d 100644
--- a/sdks/go/pkg/beam/core/runtime/xlangx/expand.go
+++ b/sdks/go/pkg/beam/core/runtime/xlangx/expand.go
@@ -29,11 +29,15 @@ import (
        "github.com/apache/beam/sdks/v2/go/pkg/beam/core/runtime/pipelinex"
        
"github.com/apache/beam/sdks/v2/go/pkg/beam/core/runtime/xlangx/expansionx"
        "github.com/apache/beam/sdks/v2/go/pkg/beam/internal/errors"
+       "github.com/apache/beam/sdks/v2/go/pkg/beam/log"
        jobpb 
"github.com/apache/beam/sdks/v2/go/pkg/beam/model/jobmanagement_v1"
        pipepb "github.com/apache/beam/sdks/v2/go/pkg/beam/model/pipeline_v1"
+       "github.com/apache/beam/sdks/v2/go/pkg/beam/options/jobopts"
+       "github.com/apache/beam/sdks/v2/go/pkg/beam/options/resource"
        "github.com/apache/beam/sdks/v2/go/pkg/beam/transforms/xlang"
        "github.com/avast/retry-go/v4"
        "google.golang.org/grpc"
+       "google.golang.org/protobuf/types/known/structpb"
 )
 
 // maxRetries is the maximum number of retries to attempt connecting to
@@ -85,7 +89,12 @@ func Expand(edge *graph.MultiEdge, ext 
*graph.ExternalTransform) error {
        delete(transforms, extTransformID)
 
        // Querying the expansion service
-       res, err := expand(context.Background(), p.GetComponents(), 
extTransform, edge, ext)
+       ctx := context.Background()
+       pipelineOpts, err := expansionPipelineOptions(ctx, 
jobopts.GetPipelineResourceHints())
+       if err != nil {
+               return errors.Wrapf(err, "unable to generate pipeline options 
for expansion of %v", ext)
+       }
+       res, err := expand(ctx, p.GetComponents(), extTransform, edge, ext, 
pipelineOpts)
        if err != nil {
                return err
        }
@@ -104,12 +113,45 @@ func Expand(edge *graph.MultiEdge, ext 
*graph.ExternalTransform) error {
        return nil
 }
 
+// resourceHintsOptionURN is the portable pipeline option key for resource 
hints,
+// as understood by the PipelineOptions translation of the Java and Python 
SDKs.
+const resourceHintsOptionURN = "beam:option:resource_hints:v1"
+
+// expansionPipelineOptions produces the pipeline options to be sent with an
+// ExpansionRequest, or nil if there are none to send.
+//
+// Currently only the pipeline level resource hints are forwarded. Expansion
+// services use the request options as the defaults for the expansion pipeline,
+// which applies the hints to the environments of the expanded transforms. This
+// matches the behavior of the Java and Python SDKs, which forward their 
pipeline
+// options when expanding cross-language transforms.
+//
+// Other options are intentionally not forwarded, as Go SDK flags don't 
necessarily
+// share names or semantics with options in other SDKs.
+func expansionPipelineOptions(ctx context.Context, hints resource.Hints) 
(*structpb.Struct, error) {
+       opts, omitted := hints.OptionStrings()
+       if len(omitted) > 0 {
+               log.Warnf(ctx, "resource hints %v have no portable 
representation and won't be forwarded to expansion services", omitted)
+       }
+       if len(opts) == 0 {
+               return nil, nil
+       }
+       vals := make([]any, 0, len(opts))
+       for _, o := range opts {
+               vals = append(vals, o)
+       }
+       return structpb.NewStruct(map[string]any{
+               resourceHintsOptionURN: vals,
+       })
+}
+
 func expand(
        ctx context.Context,
        comps *pipepb.Components,
        transform *pipepb.PTransform,
        edge *graph.MultiEdge,
-       ext *graph.ExternalTransform) (*jobpb.ExpansionResponse, error) {
+       ext *graph.ExternalTransform,
+       pipelineOpts *structpb.Struct) (*jobpb.ExpansionResponse, error) {
 
        h, config := defaultReg.getHandlerFunc(transform.GetSpec().GetUrn(), 
ext.ExpansionAddr)
        // Overwrite expansion address if changed due to override for service 
or URN.
@@ -141,6 +183,7 @@ func expand(
                        Transform:           transform,
                        Namespace:           ext.Namespace,
                        OutputCoderRequests: outputCoderID,
+                       PipelineOptions:     pipelineOpts,
                },
                edge: edge,
                ext:  ext,
diff --git a/sdks/go/pkg/beam/core/runtime/xlangx/expand_test.go 
b/sdks/go/pkg/beam/core/runtime/xlangx/expand_test.go
new file mode 100644
index 00000000000..7f45f3cc866
--- /dev/null
+++ b/sdks/go/pkg/beam/core/runtime/xlangx/expand_test.go
@@ -0,0 +1,148 @@
+// 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 xlangx
+
+import (
+       "context"
+       "testing"
+
+       "github.com/apache/beam/sdks/v2/go/pkg/beam/core/graph"
+       jobpb 
"github.com/apache/beam/sdks/v2/go/pkg/beam/model/jobmanagement_v1"
+       pipepb "github.com/apache/beam/sdks/v2/go/pkg/beam/model/pipeline_v1"
+       "github.com/apache/beam/sdks/v2/go/pkg/beam/options/resource"
+       "github.com/google/go-cmp/cmp"
+       "google.golang.org/protobuf/testing/protocmp"
+       "google.golang.org/protobuf/types/known/structpb"
+)
+
+// nonPortableHint is a resource hint without a portable option representation.
+type nonPortableHint struct{}
+
+func (nonPortableHint) URN() string                                    { 
return "beam:resources:not_portable:v1" }
+func (nonPortableHint) Payload() []byte                                { 
return []byte("value") }
+func (h nonPortableHint) MergeWithOuter(_ resource.Hint) resource.Hint { 
return h }
+
+func mustStruct(t *testing.T, m map[string]any) *structpb.Struct {
+       t.Helper()
+       s, err := structpb.NewStruct(m)
+       if err != nil {
+               t.Fatalf("structpb.NewStruct(%v) = %v", m, err)
+       }
+       return s
+}
+
+func TestExpansionPipelineOptions(t *testing.T) {
+       tests := []struct {
+               name  string
+               hints resource.Hints
+               want  map[string]any // nil indicates no options are expected.
+       }{
+               {
+                       name: "noHints",
+               }, {
+                       name:  "onlyNonPortableHints",
+                       hints: resource.NewHints(nonPortableHint{}),
+               }, {
+                       name:  "standardHints",
+                       hints: resource.NewHints(resource.ParseMinRAM("16GB"), 
resource.Accelerator("type:nvidia-l4;count:1"), resource.CPUCount(4), 
resource.MaxActiveBundlesPerWorker(2)),
+                       want: map[string]any{
+                               resourceHintsOptionURN: []any{
+                                       "accelerator=type:nvidia-l4;count:1",
+                                       "cpu_count=4",
+                                       "max_active_bundles_per_worker=2",
+                                       "min_ram=16000000000B",
+                               },
+                       },
+               }, {
+                       name:  "nonPortableHintsDropped",
+                       hints: resource.NewHints(resource.CPUCount(4), 
nonPortableHint{}),
+                       want: map[string]any{
+                               resourceHintsOptionURN: []any{"cpu_count=4"},
+                       },
+               },
+       }
+       for _, test := range tests {
+               t.Run(test.name, func(t *testing.T) {
+                       got, err := 
expansionPipelineOptions(context.Background(), test.hints)
+                       if err != nil {
+                               t.Fatalf("expansionPipelineOptions(%v) error = 
%v", test.hints, err)
+                       }
+                       if test.want == nil {
+                               if got != nil {
+                                       t.Errorf("expansionPipelineOptions(%v) 
= %v, want nil", test.hints, got)
+                               }
+                               return
+                       }
+                       if d := cmp.Diff(mustStruct(t, test.want), got, 
protocmp.Transform()); d != "" {
+                               t.Errorf("expansionPipelineOptions(%v) diff 
(-want, +got):\n%v", test.hints, d)
+                       }
+               })
+       }
+}
+
+func TestExpand_SetsPipelineOptions(t *testing.T) {
+       // Swap in a fresh registry so the test handler doesn't leak to other 
tests.
+       oldReg := defaultReg
+       defaultReg = newRegistry()
+       defer func() { defaultReg = oldReg }()
+
+       const ns = "capturepipelineoptions"
+       var gotReq *jobpb.ExpansionRequest
+       if err := defaultReg.RegisterHandler(ns, func(_ context.Context, p 
*HandlerParams) (*jobpb.ExpansionResponse, error) {
+               gotReq = p.Req
+               return &jobpb.ExpansionResponse{}, nil
+       }); err != nil {
+               t.Fatalf("RegisterHandler(%q) = %v", ns, err)
+       }
+
+       tests := []struct {
+               name  string
+               hints resource.Hints
+               want  *structpb.Struct
+       }{
+               {
+                       name: "noHints",
+               }, {
+                       name:  "withHints",
+                       hints: resource.NewHints(resource.ParseMinRAM("2GB")),
+                       want: mustStruct(t, map[string]any{
+                               resourceHintsOptionURN: 
[]any{"min_ram=2000000000B"},
+                       }),
+               },
+       }
+       for _, test := range tests {
+               t.Run(test.name, func(t *testing.T) {
+                       gotReq = nil
+                       opts, err := 
expansionPipelineOptions(context.Background(), test.hints)
+                       if err != nil {
+                               t.Fatalf("expansionPipelineOptions(%v) error = 
%v", test.hints, err)
+                       }
+                       ext := &graph.ExternalTransform{ExpansionAddr: ns, 
Namespace: "test"}
+                       edge := &graph.MultiEdge{External: ext}
+                       transform := &pipepb.PTransform{Spec: 
&pipepb.FunctionSpec{Urn: "beam:transform:test:v1"}}
+
+                       if _, err := expand(context.Background(), 
&pipepb.Components{}, transform, edge, ext, opts); err != nil {
+                               t.Fatalf("expand() error = %v", err)
+                       }
+                       if gotReq == nil {
+                               t.Fatal("expand() didn't call the registered 
handler")
+                       }
+                       if d := cmp.Diff(test.want, 
gotReq.GetPipelineOptions(), protocmp.Transform()); d != "" {
+                               t.Errorf("ExpansionRequest.PipelineOptions diff 
(-want, +got):\n%v", d)
+                       }
+               })
+       }
+}
diff --git a/sdks/go/pkg/beam/options/jobopts/options.go 
b/sdks/go/pkg/beam/options/jobopts/options.go
index 8456c8fd061..f382d72b532 100644
--- a/sdks/go/pkg/beam/options/jobopts/options.go
+++ b/sdks/go/pkg/beam/options/jobopts/options.go
@@ -209,6 +209,8 @@ func GetPipelineResourceHints() resource.Hints {
                        h = resource.Accelerator(val)
                case "cpu_count", "beam:resources:cpu_count:v1":
                        h = resource.ParseCPUCount(val)
+               case "max_active_bundles_per_worker", 
"beam:resources:max_active_bundles_per_worker:v1":
+                       h = resource.ParseMaxActiveBundlesPerWorker(val)
                default:
                        if strings.HasPrefix(name, "beam:resources:") {
                                h = stringHint{urn: name, value: val}
diff --git a/sdks/go/pkg/beam/options/jobopts/options_test.go 
b/sdks/go/pkg/beam/options/jobopts/options_test.go
index 5bcdf39ea59..f72455bf1e6 100644
--- a/sdks/go/pkg/beam/options/jobopts/options_test.go
+++ b/sdks/go/pkg/beam/options/jobopts/options_test.go
@@ -151,15 +151,35 @@ func TestGetPipelineResourceHints(t *testing.T) {
        hints.Set("min_ram=1GB")
        hints.Set("cpu_count=1")
        hints.Set("beam:resources:cpu_count:v1=4")
+       hints.Set("max_active_bundles_per_worker=1")
+       hints.Set("beam:resources:max_active_bundles_per_worker:v1=2")
        ResourceHints = hints
+       defer func() { ResourceHints = nil }()
 
-       want := resource.NewHints(resource.ParseMinRAM("1GB"), 
resource.Accelerator("pedal_to_the_metal"), resource.ParseCPUCount("4"), 
stringHint{
+       want := resource.NewHints(resource.ParseMinRAM("1GB"), 
resource.Accelerator("pedal_to_the_metal"), resource.ParseCPUCount("4"), 
resource.MaxActiveBundlesPerWorker(2), stringHint{
                urn:   "beam:resources:novel_execution:v1",
                value: "jaguar",
        })
-       if got := GetPipelineResourceHints(); !got.Equal(want) {
+       got := GetPipelineResourceHints()
+       if !got.Equal(want) {
                t.Errorf("GetPipelineResourceHints() = %v, want %v", got, want)
        }
+
+       // Equal only compares payloads, so validate standard hints were parsed 
into
+       // their typed representations, which have portable option strings.
+       opts, omitted := got.OptionStrings()
+       wantOpts := []string{
+               "accelerator=pedal_to_the_metal",
+               "cpu_count=4",
+               "max_active_bundles_per_worker=2",
+               "min_ram=1000000000B",
+       }
+       if !reflect.DeepEqual(opts, wantOpts) {
+               t.Errorf("GetPipelineResourceHints().OptionStrings() opts = %v, 
want %v", opts, wantOpts)
+       }
+       if wantOmitted := []string{"beam:resources:novel_execution:v1"}; 
!reflect.DeepEqual(omitted, wantOmitted) {
+               t.Errorf("GetPipelineResourceHints().OptionStrings() omitted = 
%v, want %v", omitted, wantOmitted)
+       }
 }
 
 func TestGetExperiements(t *testing.T) {
diff --git a/sdks/go/pkg/beam/options/resource/hint.go 
b/sdks/go/pkg/beam/options/resource/hint.go
index 1d3f1cd37ec..84cfa9781e9 100644
--- a/sdks/go/pkg/beam/options/resource/hint.go
+++ b/sdks/go/pkg/beam/options/resource/hint.go
@@ -22,6 +22,7 @@ package resource
 import (
        "bytes"
        "fmt"
+       "sort"
        "strconv"
 
        "github.com/dustin/go-humanize"
@@ -97,6 +98,44 @@ func NewHints(hs ...Hint) Hints {
        return hints
 }
 
+// OptionStrings returns the hints formatted as "name=value" strings, in the 
format
+// accepted by the resource_hints pipeline option of the Beam SDKs. This 
allows the
+// hints to be forwarded to other SDKs, such as to expansion services for 
cross-language
+// transforms.
+//
+// Only standard hints with a well known short name are converted, since not 
every SDK
+// accepts arbitrary hint URNs. The URNs of any hints that couldn't be 
converted are
+// returned in omitted. Both lists are sorted for determinism.
+func (hs Hints) OptionStrings() (opts, omitted []string) {
+       for urn, h := range hs.h {
+               if s, ok := optionString(h); ok {
+                       opts = append(opts, s)
+               } else {
+                       omitted = append(omitted, urn)
+               }
+       }
+       sort.Strings(opts)
+       sort.Strings(omitted)
+       return opts, omitted
+}
+
+// optionString converts a known standard hint into its "name=value" option 
form.
+func optionString(h Hint) (string, bool) {
+       switch h := h.(type) {
+       case minRAMHint:
+               // Use an explicit byte unit suffix, as SDKs require a unit 
when parsing this hint.
+               return fmt.Sprintf("min_ram=%dB", h.value), true
+       case acceleratorHint:
+               return "accelerator=" + h.value, true
+       case CPUCountHint:
+               return fmt.Sprintf("cpu_count=%d", h.value), true
+       case maxActiveBundlesPerWorkerHint:
+               return fmt.Sprintf("max_active_bundles_per_worker=%d", 
h.value), true
+       default:
+               return "", false
+       }
+}
+
 // Hint contains all the information about a given resource hint.
 type Hint interface {
        // URN returns the name for this hint.
diff --git a/sdks/go/pkg/beam/options/resource/hint_test.go 
b/sdks/go/pkg/beam/options/resource/hint_test.go
index 35dc5fa0dd7..190b0e8fd2f 100644
--- a/sdks/go/pkg/beam/options/resource/hint_test.go
+++ b/sdks/go/pkg/beam/options/resource/hint_test.go
@@ -405,3 +405,48 @@ func TestHints_NilHints(t *testing.T) {
                t.Errorf("nil equal test: (nil).Equal(hs) = %v, want %v", got, 
want)
        }
 }
+
+func TestHints_OptionStrings(t *testing.T) {
+       tests := []struct {
+               name                string
+               hints               Hints
+               wantOpts, wantOmits []string
+       }{
+               {
+                       name: "empty",
+               }, {
+                       name:     "minRAM",
+                       hints:    NewHints(ParseMinRAM("2GB")),
+                       wantOpts: []string{"min_ram=2000000000B"},
+               }, {
+                       name:  "allStandard",
+                       hints: NewHints(MinRAMBytes(2e9), 
Accelerator("type:jeans;count1;"), CPUCount(4), MaxActiveBundlesPerWorker(2)),
+                       wantOpts: []string{
+                               "accelerator=type:jeans;count1;",
+                               "cpu_count=4",
+                               "max_active_bundles_per_worker=2",
+                               "min_ram=2000000000B",
+                       },
+               }, {
+                       name:      "customOmitted",
+                       hints:     NewHints(CPUCount(8), customHint{}),
+                       wantOpts:  []string{"cpu_count=8"},
+                       wantOmits: []string{"top:secret:custom:urn"},
+               }, {
+                       name:      "onlyCustom",
+                       hints:     NewHints(customHint{}),
+                       wantOmits: []string{"top:secret:custom:urn"},
+               },
+       }
+       for _, test := range tests {
+               t.Run(test.name, func(t *testing.T) {
+                       gotOpts, gotOmits := test.hints.OptionStrings()
+                       if !reflect.DeepEqual(gotOpts, test.wantOpts) {
+                               t.Errorf("OptionStrings() opts = %v, want %v", 
gotOpts, test.wantOpts)
+                       }
+                       if !reflect.DeepEqual(gotOmits, test.wantOmits) {
+                               t.Errorf("OptionStrings() omitted = %v, want 
%v", gotOmits, test.wantOmits)
+                       }
+               })
+       }
+}

Reply via email to