{"record":{"id":"922489e921e90108","repo":"apache/beam","slug":"test-stream-event-decreases-watermark-watermarks-cannot-go","errorCode":null,"errorMessage":"test stream event decreases watermark. Watermarks cannot go backwards.","messagePattern":"test stream event decreases watermark\\. Watermarks cannot go backwards\\.","errorType":"panic","errorClass":null,"httpStatus":null,"severity":"error","filePath":"sdks/go/pkg/beam/runners/prism/internal/engine/teststream.go","lineNumber":222,"sourceCode":"\tfor _, link := range em.sideConsumers[t.pcollection] {\n\t\tss := em.stages[link.Global]\n\t\tss.AddPendingSide(pending, link.Transform, link.Local)\n\t\tem.changedStages.insert(link.Global)\n\t}\n}\n\n// tsWatermarkEvent sets the watermark for the new stage.\ntype tsWatermarkEvent struct {\n\tTag          string\n\tNewWatermark mtime.Time\n}\n\n// Execute this WatermarkEvent by updating the watermark for the tag, and notify affected downstream stages.\nfunc (ev tsWatermarkEvent) Execute(em *ElementManager) {\n\tt := em.testStreamHandler.tagState[ev.Tag]\n\n\tif ev.NewWatermark < t.watermark {\n\t\tpanic(\"test stream event decreases watermark. Watermarks cannot go backwards.\")\n\t}\n\tt.watermark = ev.NewWatermark\n\tem.testStreamHandler.tagState[ev.Tag] = t\n\n\t// Update the upstream watermarks in the consumers.\n\tfor _, sID := range em.consumers[t.pcollection] {\n\t\tss := em.stages[sID]\n\t\tss.updateUpstreamWatermark(ss.inputID, t.watermark)\n\t\tem.changedStages.insert(sID)\n\t}\n\t// Clear the default hold after the inserts have occured.\n\tem.testStreamHandler.UpdateHold(em, t.watermark)\n}\n\n// tsProcessingTimeEvent implements advancing the synthetic processing time.\ntype tsProcessingTimeEvent struct {\n\tAdvanceBy time.Duration\n}","sourceCodeStart":204,"sourceCodeEnd":240,"githubUrl":"https://github.com/apache/beam/blob/12126d8942aaf848030c478b4c6a28c6af861c66/sdks/go/pkg/beam/runners/prism/internal/engine/teststream.go#L204-L240","documentation":"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.","triggerScenarios":"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.","commonSituations":"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.","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"],"exampleFix":"// before: regressing watermark\nts.AdvanceWatermarkTo(inf) // ... later:\nts.AdvanceWatermarkTo(100) // earlier\n// after: keep per-tag watermarks monotonic\nts.AdvanceWatermarkTo(100)\nts.AdvanceWatermarkTo(inf)","handlingStrategy":"validation","validationCode":"// Validate test stream watermarks are monotonic per tag before running\nlast := map[string]int64{}\nfor _, ev := range events {\n    if wm, ok := ev.(WatermarkEvent); ok && wm.NewWatermark < last[wm.Tag] {\n        return fmt.Errorf(\"tag %v regresses watermark\", wm.Tag)\n    }\n    last[wm.Tag] = wm.NewWatermark\n}","typeGuard":null,"tryCatchPattern":null,"preventionTips":["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"],"tags":["go","beam","prism","teststream","watermark"],"backgroundTag":"invalid-state-transition","analyzedSha":"12126d8942aaf848030c478b4c6a28c6af861c66","analyzedAt":"2026-09-13T01:50:10.254Z","contentChangedAt":"2026-09-13T01:50:10.254Z","schemaVersion":2},"datasetVersion":"2026-09-20T03:17:13.778Z"}