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 db063b8de70 Fix data race on Prism runner's artifact cache map (#40197)
db063b8de70 is described below
commit db063b8de70e81a0080d3421beed9022270c877f
Author: waneal <[email protected]>
AuthorDate: Tue Sep 22 00:17:59 2026 +0900
Fix data race on Prism runner's artifact cache map (#40197)
s.artifacts was read/written without a lock from
ReverseArtifactRetrievalService and GetArtifact, causing "concurrent
map writes" under the race detector. Guard it with a dedicated
sync.RWMutex.
Fixes #32656
---
CHANGES.md | 1 +
.../pkg/beam/runners/prism/internal/jobservices/artifact.go | 7 ++++---
.../go/pkg/beam/runners/prism/internal/jobservices/server.go | 12 +++++++-----
3 files changed, 12 insertions(+), 8 deletions(-)
diff --git a/CHANGES.md b/CHANGES.md
index 5480ddc1c21..4866998bb05 100644
--- a/CHANGES.md
+++ b/CHANGES.md
@@ -80,6 +80,7 @@
## Bugfixes
+* (Go) Fixed a data race on the Prism runner's artifact cache map in
JobServices ([#32656](https://github.com/apache/beam/issues/32656)).
* Fixed X (Java/Python) ([#X](https://github.com/apache/beam/issues/X)).
## Security Fixes
diff --git a/sdks/go/pkg/beam/runners/prism/internal/jobservices/artifact.go
b/sdks/go/pkg/beam/runners/prism/internal/jobservices/artifact.go
index f3e97406027..babbcbf562b 100644
--- a/sdks/go/pkg/beam/runners/prism/internal/jobservices/artifact.go
+++ b/sdks/go/pkg/beam/runners/prism/internal/jobservices/artifact.go
@@ -83,10 +83,9 @@ func (s *Server) ReverseArtifactRetrievalService(stream
jobpb.ArtifactStagingSer
return err
}
}
- if len(s.artifacts) == 0 {
- s.artifacts = map[string][]byte{}
- }
+ s.artifactsMu.Lock()
s.artifacts[string(dep.GetTypePayload())] = buf.Bytes()
+ s.artifactsMu.Unlock()
}
}
return nil
@@ -100,7 +99,9 @@ func (s *Server) ResolveArtifacts(_ context.Context, req
*jobpb.ResolveArtifacts
func (s *Server) GetArtifact(req *jobpb.GetArtifactRequest, stream
jobpb.ArtifactRetrievalService_GetArtifactServer) error {
info := req.GetArtifact()
+ s.artifactsMu.RLock()
buf, ok := s.artifacts[string(info.GetTypePayload())]
+ s.artifactsMu.RUnlock()
if !ok {
pt := prototext.Format(info)
slog.Warn("unable to provide artifact to worker",
"artifact_info", pt)
diff --git a/sdks/go/pkg/beam/runners/prism/internal/jobservices/server.go
b/sdks/go/pkg/beam/runners/prism/internal/jobservices/server.go
index 2399fd726da..f6c88205146 100644
--- a/sdks/go/pkg/beam/runners/prism/internal/jobservices/server.go
+++ b/sdks/go/pkg/beam/runners/prism/internal/jobservices/server.go
@@ -61,7 +61,8 @@ type Server struct {
execute func(*Job)
// Artifact hack
- artifacts map[string][]byte
+ artifacts map[string][]byte
+ artifactsMu sync.RWMutex
mw *worker.MultiplexW
}
@@ -73,10 +74,11 @@ func NewServer(port int, execute func(*Job)) *Server {
panic(fmt.Sprintf("failed to listen: %v", err))
}
s := &Server{
- lis: lis,
- jobs: make(map[string]*Job),
- execute: execute,
- logger: slog.Default(), // TODO substitute with a configured
logger.
+ lis: lis,
+ jobs: make(map[string]*Job),
+ execute: execute,
+ artifacts: map[string][]byte{},
+ logger: slog.Default(), // TODO substitute with a configured
logger.
}
s.logger.Info("Serving JobManagement", slog.String("endpoint",
s.Endpoint()))
opts := []grpc.ServerOption{