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)
+ }
+ })
+ }
+}