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

  1. Reorder the TestStream events so watermark advances for each tag are non-decreasing over time
  2. Split into separate TestStream tests if you need independent watermark timelines per tag
  3. Audit test helper code that appends watermark events; sort by watermark before building the stream
  4. 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

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


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)