apache/beam · error
test stream event decreases watermark. Watermarks cannot go…
Error message
test stream event decreases watermark. Watermarks cannot go backwards.
What it means
When executing a test-stream WatermarkEvent, Prism requires watermarks to be monotonically non-decreasing per tag. If the event's NewWatermark is earlier than the tag's current watermark, the runner panics, since watermark regression would invalidate all downstream scheduling and triggering logic. This only occurs when using the TestStream primitive.
Solutions
- Reorder the TestStream events so watermark advances for each tag are non-decreasing over time
- Split into separate TestStream tests if you need independent watermark timelines per tag
- Audit test helper code that appends watermark events; sort by watermark before building the stream
- Remember TestStream is test-only; for production-like behavior use regular source watermarking
Example fix
// before: regressing watermark ts.AdvanceWatermarkTo(inf) // ... later: ts.AdvanceWatermarkTo(100) // earlier // after: keep per-tag watermarks monotonic ts.AdvanceWatermarkTo(100) ts.AdvanceWatermarkTo(inf)
Defensive patterns
Strategy: validation
Validate before calling
// Validate test stream watermarks are monotonic per tag before running
last := map[string]int64{}
for _, ev := range events {
if wm, ok := ev.(WatermarkEvent); ok && wm.NewWatermark < last[wm.Tag] {
return fmt.Errorf("tag %v regresses watermark", wm.Tag)
}
last[wm.Tag] = wm.NewWatermark
} Prevention
- Build TestStream events in strictly increasing watermark order per tag
- Centralize test-stream construction in a helper that enforces monotonicity
- Split tests needing different watermark timelines into separate streams
When it happens
Trigger: A TestStream pipeline (typical in unit tests of streaming behavior) that advances a tag's watermark and later emits an earlier watermark for the same tag — e.g. watermark elements authored out of order in the test stream.
Common situations: Hand-written TestStream scenarios where watermark advance events are appended in a non-monotonic order, or where two code paths share a test stream builder and each advances the same tag.
Understand the failure class
Background: "Invalid state transition" errors: "status must be X, actually Y", "already rejected/charging/uninstalled", "cannot ... while running" — what they mean when a library rejects your call — this error's family across 31 libraries.
Related errors
- generating bundle for stage
- prism error: negative watermark hold count
- couldn't decode characteristic for variant
- CreateWatermarkEstimator fn
- error decoding append bag user state window key
AI-assisted analysis of apache/beam@12126d8942 (2026-09-13).
Data as JSON: /api/errors/922489e921e90108.
Report an issue: GitHub.
Appendix: source
Thrown at sdks/go/pkg/beam/runners/prism/internal/engine/teststream.go:222
for _, link := range em.sideConsumers[t.pcollection] {
ss := em.stages[link.Global]
ss.AddPendingSide(pending, link.Transform, link.Local)
em.changedStages.insert(link.Global)
}
}
// tsWatermarkEvent sets the watermark for the new stage.
type tsWatermarkEvent struct {
Tag string
NewWatermark mtime.Time
}
// Execute this WatermarkEvent by updating the watermark for the tag, and notify affected downstream stages.
func (ev tsWatermarkEvent) Execute(em *ElementManager) {
t := em.testStreamHandler.tagState[ev.Tag]
if ev.NewWatermark < t.watermark {
panic("test stream event decreases watermark. Watermarks cannot go backwards.")
}
t.watermark = ev.NewWatermark
em.testStreamHandler.tagState[ev.Tag] = t
// Update the upstream watermarks in the consumers.
for _, sID := range em.consumers[t.pcollection] {
ss := em.stages[sID]
ss.updateUpstreamWatermark(ss.inputID, t.watermark)
em.changedStages.insert(sID)
}
// Clear the default hold after the inserts have occured.
em.testStreamHandler.UpdateHold(em, t.watermark)
}
// tsProcessingTimeEvent implements advancing the synthetic processing time.
type tsProcessingTimeEvent struct {
AdvanceBy time.Duration
}View on GitHub (pinned to 12126d8942)