This is an automated email from the ASF dual-hosted git repository.
lostluck 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 020707df9d6 Fix flaky Go tests (#40296)
020707df9d6 is described below
commit 020707df9d6d0276a936d98e4f79f4792e761005
Author: Reuven Lax <[email protected]>
AuthorDate: Sun Sep 27 04:51:34 2026 -0700
Fix flaky Go tests (#40296)
---
.../beam/runners/prism/internal/worker/worker.go | 30 ++++++++++++++++++----
1 file changed, 25 insertions(+), 5 deletions(-)
diff --git a/sdks/go/pkg/beam/runners/prism/internal/worker/worker.go
b/sdks/go/pkg/beam/runners/prism/internal/worker/worker.go
index 33c8c3a7de5..da6601814d0 100644
--- a/sdks/go/pkg/beam/runners/prism/internal/worker/worker.go
+++ b/sdks/go/pkg/beam/runners/prism/internal/worker/worker.go
@@ -446,7 +446,15 @@ func (wk *W) Data(data fnpb.BeamFnData_DataServer) error {
//
// State requests come from SDKs, and the runner responds.
func (wk *W) State(state fnpb.BeamFnState_StateServer) error {
- responses := make(chan *fnpb.StateResponse)
+ responses := make(chan *fnpb.StateResponse, 64)
+ sendResp := func(resp *fnpb.StateResponse) bool {
+ select {
+ case responses <- resp:
+ return true
+ case <-state.Context().Done():
+ return false
+ }
+ }
go func() {
// This go routine creates all responses to state requests from
the worker
// so we want to close the State handler when it's all done.
@@ -472,6 +480,12 @@ func (wk *W) State(state fnpb.BeamFnState_StateServer)
error {
wk.mu.Unlock()
if !ok {
slog.Warn("state request after bundle
inactive", "instruction", req.GetInstructionId(), "worker", wk)
+ if !sendResp(&fnpb.StateResponse{
+ Id: req.GetId(),
+ Error: fmt.Sprintf("state request for
inactive bundle %v", req.GetInstructionId()),
+ }) {
+ return
+ }
continue
}
switch req.GetRequest().(type) {
@@ -558,13 +572,15 @@ func (wk *W) State(state fnpb.BeamFnState_StateServer)
error {
// Encode the runner iterable (no length, just
consecutive elements), and send it out.
// This is also where we can handle things like
State Backed Iterables.
- responses <- &fnpb.StateResponse{
+ if !sendResp(&fnpb.StateResponse{
Id: req.GetId(),
Response: &fnpb.StateResponse_Get{
Get: &fnpb.StateGetResponse{
Data: bytes.Join(data,
[]byte{}),
},
},
+ }) {
+ return
}
case *fnpb.StateRequest_Append:
@@ -585,11 +601,13 @@ func (wk *W) State(state fnpb.BeamFnState_StateServer)
error {
panic(fmt.Sprintf("unsupported StateKey
Append type: %T: %v", key.GetType(), prototext.Format(key)))
}
- responses <- &fnpb.StateResponse{
+ if !sendResp(&fnpb.StateResponse{
Id: req.GetId(),
Response: &fnpb.StateResponse_Append{
Append:
&fnpb.StateAppendResponse{},
},
+ }) {
+ return
}
case *fnpb.StateRequest_Clear:
@@ -602,7 +620,7 @@ func (wk *W) State(state fnpb.BeamFnState_StateServer)
error {
mmkey := key.GetMultimapUserState()
b.OutputData.ClearMultimapState(engine.LinkID{Transform:
mmkey.GetTransformId(), Local: mmkey.GetUserStateId()}, mmkey.GetWindow(),
mmkey.GetKey(), mmkey.GetMapKey())
case *fnpb.StateKey_MultimapKeysUserState_:
- mmkey := key.GetMultimapUserState()
+ mmkey := key.GetMultimapKeysUserState()
b.OutputData.ClearMultimapKeysState(engine.LinkID{Transform:
mmkey.GetTransformId(), Local: mmkey.GetUserStateId()}, mmkey.GetWindow(),
mmkey.GetKey())
case *fnpb.StateKey_OrderedListUserState_:
olkey := key.GetOrderedListUserState()
@@ -611,11 +629,13 @@ func (wk *W) State(state fnpb.BeamFnState_StateServer)
error {
default:
panic(fmt.Sprintf("unsupported StateKey
Clear type: %T: %v", key.GetType(), prototext.Format(key)))
}
- responses <- &fnpb.StateResponse{
+ if !sendResp(&fnpb.StateResponse{
Id: req.GetId(),
Response: &fnpb.StateResponse_Clear{
Clear:
&fnpb.StateClearResponse{},
},
+ }) {
+ return
}
default: