From feddf873e88678f7f4eaf21f683e9a808e9c6ab9 Mon Sep 17 00:00:00 2001 From: Eliaazzz Date: Sat, 1 Aug 2026 01:46:44 +1000 Subject: [PATCH] [Prism] Schedule consumers of a self checkpointing source Fixes #39446. An unbounded source SDF that returns a process continuation residual has that residual re-queued as a pending element carrying the input element's event time, MinTimestamp for an Impulse rooted source, so the stage's watermark never advances. updateWatermarks returns no refreshes when the output watermark does not move, and PersistBundle marked only the producing stage as changed, so a consumer that was just handed pending elements was never surfaced to the scheduler. Its elements accumulated forever. checkForQuiescence cannot catch this, because the source stays schedulable, so the job live locks instead of failing fast. PersistBundle now records the consumers that accepted data, and updateWatermarks hands them to the scheduler on the path where the output watermark does not advance. Recording is gated on any bundle having returned a residual, so a pipeline that never self checkpoints keeps its previous bundle scheduling. bundleReady additionally lets a stateful stage with no side inputs run on pending data alone, under the same gate, since statefulStageKind.buildEventTimeBundle takes data at any watermark and gates only timers. Its stillSchedulable now requires buildable work, and a key that supplies nothing no longer consumes the OneKeyPerBundle slot, is not marked in progress, and does not hold the bundle's minimum timestamp. A consumer that reads a side input still waits on the watermark, since side input readiness is derived from it. --- .../prism/internal/engine/elementmanager.go | 91 +++++- .../elementmanager_continuation_test.go | 299 ++++++++++++++++++ 2 files changed, 377 insertions(+), 13 deletions(-) create mode 100644 sdks/go/pkg/beam/runners/prism/internal/engine/elementmanager_continuation_test.go 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 2b502c679db9..7ebf174d9ec0 100644 --- a/sdks/go/pkg/beam/runners/prism/internal/engine/elementmanager.go +++ b/sdks/go/pkg/beam/runners/prism/internal/engine/elementmanager.go @@ -235,6 +235,10 @@ type ElementManager struct { livePending atomic.Int64 // An accessible live pending count. DEBUG USE ONLY pendingElements sync.WaitGroup // pendingElements counts all unprocessed elements in a job. Jobs with no pending elements terminate successfully. + // Latched once any bundle returns a residual, after which a stage watermark + // can be pinned indefinitely and can't be relied on to schedule consumers. + sawResidual atomic.Bool + processTimeEvents *stageRefreshQueue // Manages sequence of stage updates when interfacing with processing time. Callers must hold refreshCond.L lock. testStreamHandler *testStreamHandler // Optional test stream handler when a test stream is in the pipeline. } @@ -841,6 +845,17 @@ func reElementResiduals(residuals []Residual, inputInfo PColInfo, rb RunBundle) // input elements, and the committed output elements. func (em *ElementManager) PersistBundle(rb RunBundle, col2Coders map[string]PColInfo, d TentativeData, inputInfo PColInfo, residuals Residuals) { stage := em.stages[rb.StageID] + // Consumers that received data from this bundle, recorded so they can still + // be scheduled when this stage's output watermark is held back. Only needed + // once something self checkpoints, so pipelines that never do keep their + // previous bundle scheduling exactly. + if len(residuals.Data) > 0 { + em.sawResidual.Store(true) + } + var changedConsumers set[string] + if em.sawResidual.Load() { + changedConsumers = set[string]{} + } var seq int for output, data := range d.Raw { info := col2Coders[output] @@ -910,6 +925,9 @@ func (em *ElementManager) PersistBundle(rb RunBundle, col2Coders map[string]PCol count = consumer.AddPending(em, newPending) } em.addPending(count) + if changedConsumers != nil && count > 0 { + changedConsumers.insert(sID) + } } for _, link := range sideConsumers { consumer := em.stages[link.Global] @@ -942,6 +960,12 @@ func (em *ElementManager) PersistBundle(rb RunBundle, col2Coders map[string]PCol // even if a panic occurs during `em.addPending`. This prevents potential deadlocks // if the waitgroup unexpectedly drops below zero due to a runner bug. defer stage.mu.Unlock() + if len(changedConsumers) > 0 { + if stage.consumersWithNewData == nil { + stage.consumersWithNewData = set[string]{} + } + stage.consumersWithNewData.merge(changedConsumers) + } completed := stage.inprogress[rb.BundleID] em.addPending(-len(completed.es)) delete(stage.inprogress, rb.BundleID) @@ -1101,6 +1125,9 @@ func (em *ElementManager) ReturnResiduals(rb RunBundle, firstRsIndex int, inputI stage := em.stages[rb.StageID] stage.splitBundle(rb, firstRsIndex, em) + if len(residuals.Data) > 0 { + em.sawResidual.Store(true) + } unprocessedElements := reElementResiduals(residuals.Data, inputInfo, rb) if len(unprocessedElements) > 0 { slog.Debug("ReturnResiduals: unprocessed elements", "bundle", rb, "count", len(unprocessedElements)) @@ -1211,6 +1238,9 @@ type stageState struct { estimatedOutput mtime.Time // Estimated watermark output from DoFns previousInput mtime.Time // input watermark before the latest watermark refresh + // Consumers handed data by a bundle, pending delivery to the scheduler. + consumersWithNewData set[string] + pending elementHeap // pending input elements for this stage that are to be processesd inprogress map[string]elements // inprogress elements by active bundles, keyed by bundle sideInputs map[LinkID]map[typex.Window][][]byte // side input data for this stage, from {tid, inputID} -> window @@ -1822,12 +1852,6 @@ keysPerBundle: if ss.inprogressKeys.present(k) { continue } - newKeys.insert(k) - // Track the min-timestamp for later watermark handling. - if dnt.elements[0].timestamp < minTs { - minTs = dnt.elements[0].timestamp - } - dataInBundle := false var toProcessForKey []element @@ -1874,22 +1898,51 @@ keysPerBundle: break } } + if len(toProcessForKey) > 0 { + newKeys.insert(k) + // Track the min-timestamp for later watermark handling. Elements pop + // in timestamp order, so the first selected one is the earliest. + if ts := toProcessForKey[0].timestamp; ts < minTs { + minTs = ts + } + } toProcess = append(toProcess, toProcessForKey...) if dnt.elements.Len() == 0 { delete(ss.pendingByKeys, k) } - if OneKeyPerBundle { + // A key that yielded nothing, such as one headed by a timer above the + // watermark, must not consume the single key slot, or the bundle is empty + // and the stage keeps rescheduling on it. + if OneKeyPerBundle && len(toProcessForKey) > 0 { break keysPerBundle } } - // If we're out of data, and timers were not cleared then the watermark is accurate. - stillSchedulable := !(len(ss.pendingByKeys) == 0 && !timerCleared) + // Reschedule only when a later bundle could build something, or a cleared + // timer may have held back the minimum pending timestamp. + stillSchedulable := timerCleared || ss.hasBuildableDataLocked(watermark) return toProcess, minTs, newKeys, holdsInBundle, nil, stillSchedulable, 0 } +// hasBuildableDataLocked reports whether a key that isn't in progress heads its +// heap with data, or with a timer the watermark has reached. Callers hold ss.mu. +func (ss *stageState) hasBuildableDataLocked(watermark mtime.Time) bool { + for k, dnt := range ss.pendingByKeys { + if ss.inprogressKeys.present(k) { + continue + } + if dnt.elements.Len() == 0 { + continue + } + if e := dnt.elements[0]; e.IsData() || e.timestamp <= watermark { + return true + } + } + return false +} + // buildEventTimeBundle for aggregation stages, processes all elements that are within the watermark for completed windows. func (*aggregateStageKind) buildEventTimeBundle(ss *stageState, watermark mtime.Time) (toProcess elementHeap, _ mtime.Time, _ set[string], _ map[mtime.Time]int, panesInBundle []bundlePane, schedulable bool, pendingAdjustment int) { minTs := mtime.MaxTimestamp @@ -2286,9 +2339,13 @@ func (ss *stageState) updateWatermarks(em *ElementManager) set[string] { if minWatermarkHold < newOut { newOut = minWatermarkHold } - // If the newOut is smaller, then don't change downstream watermarks. + // If the newOut is smaller, then don't change downstream watermarks. Any + // consumer that received data still needs scheduling, since an unadvancing + // watermark is otherwise the only thing that would surface it. if newOut <= ss.output { - return nil + refreshes := ss.consumersWithNewData + ss.consumersWithNewData = nil + return refreshes } // If bigger, advance the output watermark @@ -2328,6 +2385,7 @@ func (ss *stageState) updateWatermarks(em *ElementManager) set[string] { // Update this stage's output watermark, and then propagate that to downstream stages refreshes := set[string]{} + ss.consumersWithNewData = nil ss.output = newOut for _, outputCol := range ss.outputIDs { consumers := em.consumers[outputCol] @@ -2435,13 +2493,20 @@ func (ss *stageState) bundleReady(em *ElementManager, emNow mtime.Time) (mtime.T previousInputW := ss.previousInput _, isOrdinaryStage := ss.kind.(*ordinaryStageKind) - if isOrdinaryStage && len(ss.sides) == 0 { + _, isStatefulStage := ss.kind.(*statefulStageKind) + switch { + case isOrdinaryStage && len(ss.sides) == 0: // For ordinary stage with no side inputs, we use whether there are pending elements to determine // whether a bundle is ready or not. if len(ss.pending) == 0 { return mtime.MinTimestamp, false, ptimeEventsReady, injectedReady } - } else if inputW == upstreamW && previousInputW == inputW { + case isStatefulStage && len(ss.sides) == 0 && em.sawResidual.Load() && ss.hasBuildableDataLocked(upstreamW): + // A stateful stage processes pending data at whatever the current + // watermark is, so data alone makes it ready once something is self + // checkpointing. Side input readiness comes from the watermark, so + // stages that read one keep waiting. + case inputW == upstreamW && previousInputW == inputW: // Otherwise, use the progression of watermark to determine the bundle readiness. slog.Debug("bundleReady: unchanged upstream watermark", slog.String("stage", ss.ID), diff --git a/sdks/go/pkg/beam/runners/prism/internal/engine/elementmanager_continuation_test.go b/sdks/go/pkg/beam/runners/prism/internal/engine/elementmanager_continuation_test.go new file mode 100644 index 000000000000..2c0eac52e5e1 --- /dev/null +++ b/sdks/go/pkg/beam/runners/prism/internal/engine/elementmanager_continuation_test.go @@ -0,0 +1,299 @@ +// 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 engine + +import ( + "bytes" + "context" + "fmt" + "io" + "testing" + + "github.com/apache/beam/sdks/v2/go/pkg/beam/core/graph/coder" + "github.com/apache/beam/sdks/v2/go/pkg/beam/core/graph/mtime" + "github.com/apache/beam/sdks/v2/go/pkg/beam/core/graph/window" + "github.com/apache/beam/sdks/v2/go/pkg/beam/core/runtime/exec" + "github.com/apache/beam/sdks/v2/go/pkg/beam/core/typex" +) + +// bundleBudget bounds how many bundles a continuation test drives. The source +// never terminates, so these tests stop as soon as the consumer is scheduled and +// use the budget only to declare starvation. +const bundleBudget = 50 + +// continuationInfo is a global window PColInfo, keyed for stateful consumers. +func continuationInfo(t *testing.T, keyed bool) PColInfo { + t.Helper() + readAll := func(r io.Reader) []byte { + b, err := io.ReadAll(r) + if err != nil { + t.Fatalf("error decoding element: %v", err) + } + return b + } + info := PColInfo{ + GlobalID: "continuation_info", + WDec: exec.MakeWindowDecoder(coder.NewGlobalWindow()), + WEnc: exec.MakeWindowEncoder(coder.NewGlobalWindow()), + EDec: readAll, + } + if keyed { + info.KeyDec = readAll + } + return info +} + +// encodeElement produces a global window element at the given event time. +func encodeElement(t *testing.T, info PColInfo, et mtime.Time) []byte { + t.Helper() + var buf bytes.Buffer + if err := exec.EncodeWindowedValueHeader(info.WEnc, []typex.Window{window.GlobalWindow{}}, et, typex.NoFiringPane(), &buf); err != nil { + t.Fatalf("EncodeWindowedValueHeader: %v", err) + } + buf.Write([]byte{3, 65, 66, 67}) // "ABC" + return buf.Bytes() +} + +// TestPersistBundle_ContinuationResidualConsumers covers issue #39446: a source +// whose residual pins its watermark must still get its consumers scheduled. +func TestPersistBundle_ContinuationResidualConsumers(t *testing.T) { + for _, test := range []struct { + name string + keyed bool + stateful bool + }{ + {name: "ordinary consumer"}, + {name: "stateful consumer", keyed: true, stateful: true}, + } { + t.Run(test.name, func(t *testing.T) { + srcInfo := continuationInfo(t, false) + outInfo := continuationInfo(t, test.keyed) + + ctx, cancelFn := context.WithCancelCause(context.Background()) + defer cancelFn(nil) + + em := NewElementManager(Config{}) + em.AddStage("impulse", nil, []string{"src_in"}, nil) + em.AddStage("src", []string{"src_in"}, []string{"sink_in"}, nil) + em.AddStage("sink", []string{"sink_in"}, nil, nil) + if test.stateful { + em.StageStateful("sink", nil) + } + em.Impulse("impulse") + + var i int + ch := em.Bundles(ctx, cancelFn, func() string { + defer func() { i++ }() + return fmt.Sprintf("%v", i) + }) + + src := em.stages["src"] + + // The source emits a record and self checkpoints every round, so it + // always has a residual outstanding and never terminates. + var srcBundles, sinkBundles int + for b := 0; b < bundleBudget && sinkBundles == 0; b++ { + rb, ok := <-ch + if !ok { + t.Fatalf("bundle %d: bundles channel closed early", b) + } + switch rb.StageID { + case "sink": + sinkBundles++ + em.PersistBundle(rb, nil, TentativeData{}, outInfo, Residuals{}) + case "src": + srcBundles++ + td := TentativeData{} + td.WriteData("sink_in", encodeElement(t, outInfo, mtime.Time(100*srcBundles))) + // No reported estimate means MIN_TIMESTAMP, so only the + // arriving data can drive the consumer. + em.PersistBundle(rb, map[string]PColInfo{"sink_in": outInfo}, td, srcInfo, Residuals{ + TransformID: "src", + InputID: "i0", + Data: []Residual{{Element: encodeElement(t, srcInfo, mtime.MinTimestamp)}}, + }) + default: + t.Fatalf("bundle %d: unexpected stage %v", b, rb.StageID) + } + } + + if sinkBundles == 0 { + t.Errorf("consumer stage was never scheduled across %d source bundles, src output watermark = %v; its pending elements are starved", + srcBundles, src.OutputWatermark()) + } + }) + } +} + +// TestPersistBundle_ContinuationResidualTransitive covers a consumer two stages +// below the source. The middle stage returns no residual of its own, yet its +// watermark is still pinned by the source's. +func TestPersistBundle_ContinuationResidualTransitive(t *testing.T) { + info := continuationInfo(t, false) + ctx, cancelFn := context.WithCancelCause(context.Background()) + defer cancelFn(nil) + + em := NewElementManager(Config{}) + em.AddStage("impulse", nil, []string{"src_in"}, nil) + em.AddStage("src", []string{"src_in"}, []string{"mid_in"}, nil) + em.AddStage("mid", []string{"mid_in"}, []string{"sink_in"}, nil) + em.AddStage("sink", []string{"sink_in"}, nil, nil) + em.Impulse("impulse") + + var i int + ch := em.Bundles(ctx, cancelFn, func() string { + defer func() { i++ }() + return fmt.Sprintf("%v", i) + }) + + var srcBundles, midBundles, sinkBundles int + for b := 0; b < bundleBudget && sinkBundles == 0; b++ { + rb, ok := <-ch + if !ok { + t.Fatalf("bundle %d: bundles channel closed early", b) + } + switch rb.StageID { + case "src": + srcBundles++ + td := TentativeData{} + td.WriteData("mid_in", encodeElement(t, info, mtime.Time(100*srcBundles))) + em.PersistBundle(rb, map[string]PColInfo{"mid_in": info}, td, info, Residuals{ + TransformID: "src", + InputID: "i0", + Data: []Residual{{Element: encodeElement(t, info, mtime.MinTimestamp)}}, + }) + case "mid": + midBundles++ + td := TentativeData{} + td.WriteData("sink_in", encodeElement(t, info, mtime.Time(100*midBundles))) + em.PersistBundle(rb, map[string]PColInfo{"sink_in": info}, td, info, Residuals{}) + case "sink": + sinkBundles++ + em.PersistBundle(rb, nil, TentativeData{}, info, Residuals{}) + default: + t.Fatalf("bundle %d: unexpected stage %v", b, rb.StageID) + } + } + + if midBundles == 0 { + t.Error("middle stage was never scheduled") + } + if sinkBundles == 0 { + t.Errorf("stage two below the source was never scheduled across %d source and %d middle bundles; its pending elements are starved", + srcBundles, midBundles) + } +} + +// TestPersistBundle_ContinuationResidualWatermark pins the BundleApplication +// output_watermarks contract: an unreported estimate means MIN_TIMESTAMP. +func TestPersistBundle_ContinuationResidualWatermark(t *testing.T) { + for _, test := range []struct { + name string + report bool + wantHeld bool + }{ + {name: "reported estimate advances the watermark", report: true}, + {name: "no estimate holds the watermark", wantHeld: true}, + } { + t.Run(test.name, func(t *testing.T) { + info := continuationInfo(t, false) + ctx, cancelFn := context.WithCancelCause(context.Background()) + defer cancelFn(nil) + + em := NewElementManager(Config{}) + em.AddStage("impulse", nil, []string{"src_in"}, nil) + em.AddStage("src", []string{"src_in"}, []string{"sink_in"}, nil) + em.AddStage("sink", []string{"sink_in"}, nil, nil) + em.Impulse("impulse") + + var i int + ch := em.Bundles(ctx, cancelFn, func() string { + defer func() { i++ }() + return fmt.Sprintf("%v", i) + }) + + src := em.stages["src"] + + for round := 0; round < 3; round++ { + rb, ok := <-ch + if !ok { + t.Fatalf("round %d: bundles channel closed early", round) + } + residuals := Residuals{ + TransformID: "src", + InputID: "i0", + Data: []Residual{{Element: encodeElement(t, info, mtime.MinTimestamp)}}, + } + if test.report { + residuals.MinOutputWatermarks = map[string]mtime.Time{"sink_in": mtime.Time(1000 * (round + 1))} + } + em.PersistBundle(rb, nil, TentativeData{}, info, residuals) + } + + got := src.OutputWatermark() + if test.wantHeld && got != mtime.MinTimestamp { + t.Errorf("src.OutputWatermark() = %v, want %v: an unreported estimate defaults to MIN_TIMESTAMP", got, mtime.MinTimestamp) + } + if !test.wantHeld && got == mtime.MinTimestamp { + t.Errorf("src.OutputWatermark() = %v, want it to follow the reported estimate", got) + } + }) + } +} + +// TestStatefulBuildEventTimeBundle_OneKeyPerBundle checks that a key holding +// only a timer above the watermark does not consume the single key slot, which +// would build an empty bundle and reschedule the stage on the same key forever. +func TestStatefulBuildEventTimeBundle_OneKeyPerBundle(t *testing.T) { + OneKeyPerBundle = true + t.Cleanup(func() { OneKeyPerBundle = false }) + + // Key iteration order is randomized, so repeat until the timer key leads. + for i := 0; i < 20; i++ { + em := NewElementManager(Config{}) + ss := makeStageState("stateful", []string{"input"}, nil, nil) + ss.kind = &statefulStageKind{} + ss.AddPending(em, []element{{ + window: window.GlobalWindow{}, + timestamp: mtime.MaxTimestamp - 1, + holdTimestamp: mtime.MaxTimestamp - 1, + pane: typex.NoFiringPane(), + transform: "stateful", + family: "timer", + keyBytes: []byte("timerkey"), + sequence: 0, + }, { + window: window.GlobalWindow{}, + timestamp: 10, + pane: typex.NoFiringPane(), + elmBytes: []byte{3, 65, 66, 67}, + keyBytes: []byte("datakey"), + }}) + + toProcess, minTs, newKeys, _, _, _, _ := ss.kind.buildEventTimeBundle(ss, mtime.Time(100)) + if len(toProcess) == 0 { + t.Fatalf("iteration %d: built an empty bundle while a data key was pending", i) + } + // The skipped timer key must not be marked in progress, nor hold the + // bundle's minimum timestamp. + if len(newKeys) != 1 || !newKeys.present("datakey") { + t.Fatalf("iteration %d: newKeys = %v, want only datakey", i, newKeys) + } + if want := mtime.Time(10); minTs != want { + t.Fatalf("iteration %d: minTs = %v, want %v", i, minTs, want) + } + } +}