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