From 10653c038236397d95a132bf717516fb489785c0 Mon Sep 17 00:00:00 2001 From: Shunping Huang Date: Thu, 2 Oct 2025 21:49:59 -0400 Subject: [PATCH 1/3] Fix race condition when injecting processing-time bundle --- .../go/pkg/beam/runners/prism/internal/engine/elementmanager.go | 2 ++ 1 file changed, 2 insertions(+) diff --git a/sdks/go/pkg/beam/runners/prism/internal/engine/elementmanager.go b/sdks/go/pkg/beam/runners/prism/internal/engine/elementmanager.go index d03d906e47de..291ce606d187 100644 --- a/sdks/go/pkg/beam/runners/prism/internal/engine/elementmanager.go +++ b/sdks/go/pkg/beam/runners/prism/internal/engine/elementmanager.go @@ -1366,7 +1366,9 @@ func (ss *stageState) injectTriggeredBundlesIfReady(em *ElementManager, window t // TODO: how to deal with watermark holds for this implicit processing time timer // ss.watermarkHolds.Add(timer.holdTimestamp, 1) ss.processingTimeTimers.Persist(firingTime, timer, notYetHolds) + em.refreshCond.L.Lock() em.processTimeEvents.Schedule(firingTime, ss.ID) + em.refreshCond.L.Unlock() em.wakeUpAt(firingTime) } } From a2b78f6ed09878a52ccf215b22d9345633dd34b1 Mon Sep 17 00:00:00 2001 From: Shunping Huang Date: Fri, 3 Oct 2025 00:02:43 -0400 Subject: [PATCH 2/3] Fix an issue of dereferencing nil. --- .../pkg/beam/runners/prism/internal/engine/elementmanager.go | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/sdks/go/pkg/beam/runners/prism/internal/engine/elementmanager.go b/sdks/go/pkg/beam/runners/prism/internal/engine/elementmanager.go index 291ce606d187..a4f13df47630 100644 --- a/sdks/go/pkg/beam/runners/prism/internal/engine/elementmanager.go +++ b/sdks/go/pkg/beam/runners/prism/internal/engine/elementmanager.go @@ -384,6 +384,7 @@ func (em *ElementManager) Bundles(ctx context.Context, upstreamCancelFn context. defer func() { // In case of panics in bundle generation, fail and cancel the job. if e := recover(); e != nil { + slog.Error("panic in ElementManager.Bundles watermark evaluation goroutine", "error", e, "traceback", string(debug.Stack())) upstreamCancelFn(fmt.Errorf("panic in ElementManager.Bundles watermark evaluation goroutine: %v\n%v", e, string(debug.Stack()))) } }() @@ -1568,6 +1569,9 @@ func (ss *stageState) savePanes(bundID string, panesInBundle []bundlePane) { func (ss *stageState) buildTriggeredBundle(em *ElementManager, key string, win typex.Window) ([]element, int) { var toProcess []element dnt := ss.pendingByKeys[key] + if dnt == nil { + return toProcess, 0 + } var notYet []element // Look at all elements for this key, and only for this window. From 3f88c2776154c96426d4869fce1a3bea02753858 Mon Sep 17 00:00:00 2001 From: Shunping Huang Date: Fri, 3 Oct 2025 00:18:55 -0400 Subject: [PATCH 3/3] Add comments to explain the nil pointer case. --- .../pkg/beam/runners/prism/internal/engine/elementmanager.go | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/sdks/go/pkg/beam/runners/prism/internal/engine/elementmanager.go b/sdks/go/pkg/beam/runners/prism/internal/engine/elementmanager.go index a4f13df47630..f77844b6f6ca 100644 --- a/sdks/go/pkg/beam/runners/prism/internal/engine/elementmanager.go +++ b/sdks/go/pkg/beam/runners/prism/internal/engine/elementmanager.go @@ -1570,6 +1570,10 @@ func (ss *stageState) buildTriggeredBundle(em *ElementManager, key string, win t var toProcess []element dnt := ss.pendingByKeys[key] if dnt == nil { + // If we set an after-processing-time trigger, but some other triggers fire or + // the end of window is reached before the first trigger could fire, then + // the pending elements are processed in other bundles, leaving a nil when + // we try to build this triggered bundle. return toProcess, 0 } var notYet []element