DEV Community

Andrey Karazhev
Andrey Karazhev

Posted on

Why a Numaflow watermark wouldn’t move

I started looking into issue #3645 in Numaflow, a Kubernetes-native stream processing platform, because of a workaround that seemed odd.

The pipeline was still processing messages. The vertices were healthy, and there were no errors in the logs. But pipeline_processing_lag kept growing, first for hours and then for days. The accumulator vertex’s watermark wasn’t moving.

The reporter, th0ger, was using Numaflow 1.8.3 with Python SDK 0.14.0. They had already found a way to get things moving again: emit an empty message from the accumulator, tag it DROP, and give it a watermark field.

That was the part I wanted to understand. Why would sending a message just to drop it make a difference?

A watermark tells the rest of the pipeline how far processing has progressed in event time. Windows and timeouts depend on it, and so does the lag metric. So a pipeline can keep handling messages while its watermark stays stuck. From the outside, it looks like processing is falling further and further behind.

Skipping messages was supposed to work

There was already some guidance on this in numaflow-go#228. The maintainers had explained that, since version 1.8.1, an accumulator didn’t need to return a message just to discard it:

“you can simply skip the message you don't want to process and the windows will be closed after the specified timeout”

In September, the same point came up again:

“any particular reason you need DROP semantics now that skipping is possible and timeout kicking in to do close-of-book?”

That sounded reasonable. Skip the message, wait for the timeout, and let the engine close the window.

But it didn’t explain the reporter’s experience. If the timeout handled cleanup, why did they need to send a response themselves?

Rather than assume the advice was wrong, I wanted to find out what had to be true for it to work.

Following the watermark through the code

I started with the watermark the vertex publishes.

In watermark/isb.rs, the published value is constrained by the oldest tracked timestamp, with a one-millisecond offset. That means one old timestamp can hold back the watermark for the entire vertex.

The accumulator’s window manager gets that oldest timestamp by looking across both active and closed windows. Each key’s window tracks a timestamp for every message it receives.

So the next question was: what removes those timestamps?

The deletion code removes timestamps older than the end of the supplied window. Anything at or after that boundary stays.

The reducer calls this cleanup for its tracked windows, but the important detail is where that list comes from: it is populated by responses from the UDF.

No response means no cleanup through that path.

There was still the timeout, though. A window could close, the engine could ask the UDF to finish it, and the resulting response could trigger cleanup.

I checked the close condition. The input watermark has to pass the key’s latest event time plus the timeout. And that latest event time keeps moving forward as newer messages arrive.

That was the detail I’d been missing. This wasn’t a timeout measured from when the window first opened. New events could keep pushing the closing boundary forward.

Finally, I checked how the Python SDK constructs responses. For a normal response, it uses the highest watermark among the messages returned by the user as the window end. For a close response, it echoes the window being closed.

At that point, there was one case where I couldn’t see cleanup happening: a key that kept receiving messages but never returned anything.

The key never goes quiet

Suppose an accumulator filters out all messages for a particular key. Messages keep arriving for that key, but the UDF intentionally produces no output.

As long as those messages arrive often enough, the latest event time keeps advancing before the timeout condition is reached. The window stays open.

Because it stays open, there is no close request and no close response. Because the UDF is filtering everything for that key, there are no normal responses either.

Neither path supplies a window boundary for cleanup.

Meanwhile, each incoming message adds another timestamp. The first one remains the oldest, and that timestamp continues to hold back the vertex watermark. The lag grows, and the retained state grows with it.

So “just skip the message” works for a key that eventually goes quiet, provided the input watermark advances far enough to close its window. It doesn’t cover a key that remains active and never emits a response.

The distinction wasn’t simply a version number or a setting. It was the traffic pattern.

This also explained the workaround. When the UDF emits an empty DROP message with a watermark, the SDK turns that into a response whose window ends at the supplied watermark. Cleanup can then remove the timestamps older than that boundary.

The useful part of the message wasn’t its payload. It was the progress information carried by the response. The UDF was telling the engine how far it could safely move forward for that key.

Checking the window logic

Reading the code gave me a possible explanation, but I wanted to check it.

I added a small test in a copy of the repository at commit 372bd73, where the relevant logic matched v1.8.3. I simulated a key receiving one message every ten seconds for an hour, with a sixty-second timeout. The input watermark stayed one second behind each event, and there were no UDF responses.

#[test]
fn probe_3645_active_key_without_responses_pins_watermark() {
    let timeout = Duration::from_secs(60);
    let windower = AccumulatorWindowManager::new(timeout);
    let base = DateTime::from_timestamp_millis(1_700_000_000_000).unwrap();
    let mk = |t: DateTime<Utc>| Message {
        keys: Arc::from(vec!["k".to_string()]),
        event_time: t,
        ..Default::default()
    };
    for i in 0..360 {
        let t = base + chrono::Duration::seconds(10 * i);
        windower.assign_windows(mk(t));
        let closed = windower.close_windows(t - chrono::Duration::seconds(1));
        assert!(closed.is_empty(), "window closed at step {i}");
    }
    assert_eq!(windower.oldest_window_end_time().unwrap(), base);
    let wm = base + chrono::Duration::seconds(3589);
    windower.delete_window(Window::new(DateTime::from_timestamp_millis(0).unwrap(), wm,
        Arc::from(vec!["k".to_string()])));
    assert_eq!(windower.oldest_window_end_time().unwrap(), base + chrono::Duration::seconds(3590));
}
Enter fullscreen mode Exit fullscreen mode

No window closed during the simulated hour. All 360 message timestamps remained, and the oldest was still the timestamp of the first message.

The last part of the test calls delete_window directly, using the window boundary that a response carrying the watermark would provide. That removed everything except the final timestamp.

This wasn’t an end-to-end pipeline test. I didn’t run the Python SDK or exercise the full response path. It checked the window manager behavior: continuous input prevented the window from closing, the oldest timestamp stayed put, and supplying a cleanup boundary released the old timestamps.

That was enough to take the hypothesis back to the issue.

What happened in the discussion

I posted the analysis and asked the reporter whether the affected keys ever stopped receiving messages.

They confirmed the pattern. Some keys produced output. Others never did because they were filtered out by design. Messages kept arriving for all of them.

Initially, a maintainer said the watermark should advance without an explicit DROP and that he would look into a fix. On 21 September, he came back with a different conclusion:

“My comment above is incorrect… without any input from the user, we cannot determine the watermark to publish downstream… This is expected behaviour”

The project lead made the same point: “we cannot move the WM unless there is some response from the UDF”.

That changed how the issue was understood. The engine couldn’t infer the UDF’s progress for an active key that never responded. The missing response wasn’t something the timeout would necessarily supply.

The issue was closed. Two older issues about a to_drop API, which had been closed in June as unnecessary, were reopened. On 22 September, my documentation fix, #3660, was merged to explain when an accumulator UDF needs to respond for the watermark to advance.

The reporter’s reaction summed up the change in understanding:

“that was an interesting design turn! That changes this bug into a feature”

What I found useful about the discussion was that everyone had a different part of the explanation. The reporter had the symptom and a working workaround. The earlier guidance described what happens when a key goes quiet. Tracing the code showed why that guidance didn’t cover a continuously active key with no output.

Once that distinction was clear, the discussion moved from fixing an apparent timeout failure to documenting the response the engine needs.

What I took away from it

The thing I’ll remember is how easy it is to read a cleanup condition and assume it will eventually run.

There was a timeout. There was code to close windows. There was code to delete old timestamps. None of those pieces was obviously broken on its own.

But for this traffic pattern, the events needed to reach cleanup never happened.

That also explains why the problem could sit there without producing an error. Messages were being processed. The UDF was filtering them as intended. The logs had nothing unusual to report. The growing lag was the visible sign that old state wasn’t being released.

I’d look for the same pattern in other systems: a queue waiting for an acknowledgement, a buffer waiting for a flush, or a cache whose cleanup depends on a later access. Having a cleanup path isn’t enough. The workload has to reach it.

The question I’d ask now is: what actually removes this state, and can the system keep accepting data without that ever happening?

Top comments (0)