Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 4 additions & 3 deletions .claude/skills/slow-soak/slow.yml
Original file line number Diff line number Diff line change
Expand Up @@ -33,8 +33,6 @@ tables:
size_seconds: 60
grace_seconds: 60
idle_close_seconds: 60
late_rows: drop
poll_interval_seconds: 10
emit_sql: |
SELECT bucket, sum(count)::BIGINT AS count, max(max_latency_s) AS max_latency_s
FROM closed
Expand Down Expand Up @@ -64,13 +62,16 @@ pipeline:
auto_offset_reset: earliest
topics:
- "{{ SQLFLOW_TOPIC|default('slow-soak') }}"
event_time:
path: sent_at
format: rfc3339
handler:
type: 'handlers.InferredMemBatch'
sql: |
INSERT INTO agg_slow
BY NAME
SELECT
date_trunc('minute', CAST(sent_at AS TIMESTAMPTZ)) AS bucket,
time_bucket(INTERVAL '1 minute', event_time) AS bucket,
count(*) AS count,
max(epoch(now() - CAST(sent_at AS TIMESTAMPTZ))) AS max_latency_s
FROM batch
Expand Down
28 changes: 28 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,34 @@

### Changed

- **The window's manager no longer polls, and lateness is decided when a
record arrives.** The engine tells the manager when a window's watermark
moved; the manager publishes what closed, republishes what a late row
landed in, and deletes what is past its lateness. It has no clock.
`poll_interval_seconds` is gone: nothing is left for it to pace, and a
config that sets it fails validation naming the change.

`late_rows` is gone too, replaced by `allowed_lateness_seconds` (default 0):
Flink's `allowedLateness`. With 0, a record for a bucket that closed is
refused before the handler and counted in
`window_late_rows_total{outcome="refused"}` -- it never reaches the window
table, which is what `drop` did after the fact. With a positive value a
closed bucket's rows are kept that long, a late record within it is written,
and the bucket is republished as a **whole value** -- `emit_sql` over every
row it has -- so the sink must replace by key; validate refuses an
append-only sink. This replaces `reemit`, which published the late rows
alone as a delta that only an additive sink could use and that a replacing
sink turned into loss. The sink now always receives the exact value, and
`window_recomputes_total` counts each republication.

`time_column` must be `time_bucket(INTERVAL '<size>', event_time)`, which
validate now refuses rather than warns: the engine decides lateness from
that bucket, and the six shipped windowed examples are updated to it. The
Render template follows when its image pin moves to a release that carries
this. The wire bundle's `late_rows_reemitted` is gone and
`late_rows_recomputed` is new; `late_rows_dropped` now counts refusals.
See `docs/superpowers/specs/2026-09-26-watermark-driven-close-design.md`.

- **Windowing is now decided by a watermark the engine asserts.** A window
manager needs one thing, how far the stream has got, and nothing used to
tell it: it inferred the answer from the newest bucket in its own table and
Expand Down
84 changes: 52 additions & 32 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -1203,8 +1203,7 @@ tables:
size_seconds: 3600
grace_seconds: 0
idle_close_seconds: 60
late_rows: drop # or reemit; required
poll_interval_seconds: 10 # optional
allowed_lateness_seconds: 0 # optional; see Late rows
emit_sql: | # optional; default SELECT * FROM closed
SELECT bucket, city, sum(count)::INT AS count
FROM closed
Expand All @@ -1216,25 +1215,38 @@ tables:
topic: output-tumbling-window-1

pipeline:
source:
type: kafka
kafka:
brokers: [localhost:9092]
group_id: tumbling
topics: [tumbling-window]
event_time: {path: timestamp, format: rfc3339}
handler:
type: handlers.InferredMemBatch
sql: |
INSERT INTO agg_cities_count BY NAME
SELECT
date_trunc('hour', CAST(timestamp AS TIMESTAMPTZ)) AS bucket,
time_bucket(INTERVAL '1 hour', event_time) AS bucket,
properties.city AS city,
count(*) AS count
FROM batch
GROUP BY bucket, city
```

The handler must cut `time_column` as `time_bucket(INTERVAL '<size>',
event_time)`: the engine decides each record's lateness from that bucket, and
`sqlflow validate` refuses any other expression. `event_time` is the column
the source fills from the record's own time; see **Event time** under each
source.

**The watermark.** The engine asserts it, the window follows it. One instant
per window, in event time, written to `sqlflow_watermarks` in the same commit
as the rows it describes, never moving backwards. Its meaning is a promise:
as of this commit, every row this pipeline will ever write for this window has
`time_column` at or after the watermark. A bucket is closed when its end,
`time_column + size_seconds`, is at or before it; `sqlflow_windows` holds what
each window has closed up to, and a poll moves that to the assertion and
each window has closed up to, and a pass moves that to the assertion and
publishes the buckets between them. That is the whole decision -- no clock, no
liveness reading, one comparison in one clock.

Expand Down Expand Up @@ -1265,33 +1277,38 @@ websocket reconnecting, are the first case: the window holds through the
outage, and the clock restarts on the return. A source that can never deliver
holds every bucket open until it does, which the log says.

`poll_interval_seconds` only bounds how soon a closed bucket is noticed; on
shutdown the drain asserts first, so the final poll sees the watermark up to
the signal. A pipeline's idle ticks write nothing unless a watermark moved --
no statement, no WAL append, no fsync -- so a quiet stream costs one write,
when its last partition goes idle. `sqlflow_progress` is liveness for `/stats`
and `/healthz` and for SQL to read; no window decision rests on it.
Nothing polls. The engine signals the window's manager after the commit that
moved its watermark, and the manager passes: publishes what closed, deletes
what is past its lateness, records where it closed up to. It passes once when
it starts and once on the drain, which asserts first so the final pass sees
the watermark up to the signal, and otherwise only when told. A pipeline's
idle ticks write nothing unless a watermark moved -- no statement, no WAL
append, no fsync -- so a quiet stream costs one write, when its last partition
goes idle. `sqlflow_progress` is liveness for `/stats` and `/healthz` and for
SQL to read; no window decision rests on it.

Every decision a window makes is a row in one of two tables, rendered from
the code to [docs/windows/decisions.md](docs/windows/decisions.md).

**Late rows.** A row for a bucket below the watermark arrived after that
bucket was published. `late_rows` is required, because the two policies are
different promises to the sink. `drop` deletes the row and counts it in
`window_late_rows_total`, so the sink sees each bucket once, as it closed.
`reemit` runs `emit_sql` over the late rows alone and publishes the result,
because the bucket's other rows were deleted when it closed. The sink has to
add that result to the bucket it holds. An upsert that replaces on the
bucket's key overwrites the bucket's count with the late rows' count, so pair
a replacing upsert with `drop`. An upsert that adds counts a republished close
twice; see **Guarantees**. `sqlflow validate` warns when `reemit` is paired
with the Iceberg or Kafka sink. A rising drop count means the grace is too
short for the stream.

The watermark is one value for the whole table, which makes it the fastest
partition's clock. A topic whose partitions run at uneven rates has a slow
partition whose rows arrive late, and the grace is the allowance for them.
Size it from `window_late_rows_total`.
**Late rows.** A record whose bucket ended at or before the watermark is
late. The engine decides that when the record arrives, from its `event_time`
and the window's size, before the handler runs -- the way Flink's window
operator does. With `allowed_lateness_seconds` absent or `0` a late record is
refused and counted in `window_late_rows_total{outcome="refused"}`; the
window table never holds it. With a positive value a closed bucket's rows are
kept that long, a late record within it is written, and the bucket is
republished as a **whole value** -- `emit_sql` over every row it has -- so the
sink must replace by key; `sqlflow validate` refuses an append-only sink and
warns for one it cannot classify. Past `end + allowed_lateness_seconds` the
rows are deleted and later records are refused. Each republication counts in
`window_recomputes_total`, and the rows that caused it in
`window_late_rows_total{outcome="recomputed"}`. A rising refused count means
the grace is too short for the stream, or the lateness is.

The watermark is the minimum over the partitions, so a slow partition holds
the buckets it shares open rather than arriving late; the grace is the
allowance for reordering within a partition. Size it from
`window_late_rows_total`.

**emit_sql** shapes the rows before the sink. It reads one relation, `closed`,
holding every row of every bucket that just closed. Cast a `sum` back to the
Expand Down Expand Up @@ -1418,14 +1435,17 @@ joins, with or without a state path:
|---|---|---|---|
| `window_watermark_seconds` | `window_watermark_seconds` | gauge | `window` |
| `window_closed` | `window_closed_total` | counter | `window` |
| `window_late_rows` | `window_late_rows_total` | counter | `window`, `policy` |
| `window_late_rows` | `window_late_rows_total` | counter | `window`, `outcome` |
| `window_recomputes` | `window_recomputes_total` | counter | `window` |
| `window_close_lag_seconds` | `window_close_lag_seconds` | gauge | `window` |
| `window_newest_bucket_start_seconds` | `window_newest_bucket_start_seconds` | gauge | `window` |

`window_closed_total` is present from startup at zero, so a windowed pipeline
is never mistaken for one without windows before its first close.
`window_late_rows_total` counts only after the close that dropped or reemitted
the rows commits.
`window_late_rows_total` is the engine's: `outcome="refused"` counts records
refused before the handler, `outcome="recomputed"` counts records admitted
within the lateness, each counted once the commit that admitted it lands.
`window_recomputes_total` counts the buckets republished whole for them.

`window_close_lag_seconds` is how far the window's closes trail the data it
holds, in event time, and it is the number to alert on: zero while closes keep
Expand All @@ -1434,8 +1454,8 @@ uses no clock, so a sparse stream does not read as stalled and a host whose
clock is wrong reports it correctly; a stream going quiet shows in
`pipeline_last_message_timestamp` instead. `window_newest_bucket_start_seconds`
ahead of a trusted clock means rows are stamped in the future, which moves the
watermark past every correctly-timed row and, under `late_rows: drop`, deletes
them. It is updated on every poll that finds rows, even one whose close fails.
watermark past every correctly-timed row and refuses them as late. It is
updated on every pass that finds rows, even one whose close fails.
`window_watermark_seconds` trails the newest bucket by `grace_seconds` by
design, which makes its age the number to size `grace_seconds` from rather
than one to alert on.
Expand Down
16 changes: 7 additions & 9 deletions dev/bench/bluesky/demo-10x-noop-sink.yml
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,6 @@
# shrinks by the same factor, and the ratios between messages, batches,
# polls and flushes stay what they are live:
#
# poll_interval_seconds 10 -> 1
# flush_interval_seconds 30 (the default) -> 3
#
# grace_seconds is stream time and stays. idle_close_seconds only fires when
Expand All @@ -17,6 +16,7 @@ commands:
sql: |
CREATE TABLE IF NOT EXISTS posts (
time_us BIGINT,
event_time TIMESTAMPTZ,
commit STRUCT(
operation TEXT,
record STRUCT(langs TEXT[])
Expand Down Expand Up @@ -50,18 +50,13 @@ tables:
# land. After idle_close_seconds with nothing arriving, every open
# bucket closes.
#
# late_rows: drop, because the sink replaces. A row for a bucket that
# already closed is discarded and counted in window_late_rows_total.
# reemit would run emit_sql over the late rows alone, because the
# bucket's other rows were deleted when it closed, and the upsert would
# replace the bucket's count with theirs. validate refuses that pairing.
# No lateness: a post for a bucket that already closed is refused
# before the handler and counted in window_late_rows_total.
window:
time_column: bucket
size_seconds: 60
grace_seconds: 60
idle_close_seconds: 60
late_rows: drop
poll_interval_seconds: 1
emit_sql: |
SELECT bucket, lang, sum(posts)::INTEGER AS posts
FROM closed
Expand All @@ -80,14 +75,17 @@ pipeline:
type: websocket
websocket:
uri: "{{ SQLFLOW_JETSTREAM_URI|default('wss://jetstream2.us-east.bsky.network/subscribe?wantedCollections=app.bsky.feed.post') }}"
event_time:
path: time_us
format: unix_us

handler:
type: handlers.StructuredBatch
table: posts
sql: |
INSERT INTO posts_per_minute_by_lang
SELECT
time_bucket(INTERVAL '1 minute', to_timestamp(time_us / 1000000)) AS bucket,
time_bucket(INTERVAL '1 minute', event_time) AS bucket,
coalesce(commit.record.langs[1], 'unknown') AS lang,
count(*) AS posts
FROM posts
Expand Down
16 changes: 7 additions & 9 deletions dev/bench/bluesky/demo-10x-sqlcommand.yml
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,6 @@
# shrinks by the same factor, and the ratios between messages, batches,
# polls and flushes stay what they are live:
#
# poll_interval_seconds 10 -> 1
# flush_interval_seconds 30 (the default) -> 3
#
# grace_seconds is stream time and stays. idle_close_seconds only fires when
Expand All @@ -27,6 +26,7 @@ commands:
sql: |
CREATE TABLE IF NOT EXISTS posts (
time_us BIGINT,
event_time TIMESTAMPTZ,
commit STRUCT(
operation TEXT,
record STRUCT(langs TEXT[])
Expand Down Expand Up @@ -60,18 +60,13 @@ tables:
# land. After idle_close_seconds with nothing arriving, every open
# bucket closes.
#
# late_rows: drop, because the sink replaces. A row for a bucket that
# already closed is discarded and counted in window_late_rows_total.
# reemit would run emit_sql over the late rows alone, because the
# bucket's other rows were deleted when it closed, and the upsert would
# replace the bucket's count with theirs. validate refuses that pairing.
# No lateness: a post for a bucket that already closed is refused
# before the handler and counted in window_late_rows_total.
window:
time_column: bucket
size_seconds: 60
grace_seconds: 60
idle_close_seconds: 60
late_rows: drop
poll_interval_seconds: 1
emit_sql: |
SELECT bucket, lang, sum(posts)::INTEGER AS posts
FROM closed
Expand Down Expand Up @@ -104,14 +99,17 @@ pipeline:
type: websocket
websocket:
uri: "{{ SQLFLOW_JETSTREAM_URI|default('wss://jetstream2.us-east.bsky.network/subscribe?wantedCollections=app.bsky.feed.post') }}"
event_time:
path: time_us
format: unix_us

handler:
type: handlers.StructuredBatch
table: posts
sql: |
INSERT INTO posts_per_minute_by_lang
SELECT
time_bucket(INTERVAL '1 minute', to_timestamp(time_us / 1000000)) AS bucket,
time_bucket(INTERVAL '1 minute', event_time) AS bucket,
coalesce(commit.record.langs[1], 'unknown') AS lang,
count(*) AS posts
FROM posts
Expand Down
18 changes: 8 additions & 10 deletions dev/bench/bluesky/demo-10x.yml
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,6 @@
# shrinks by the same factor, and the ratios between messages, batches,
# polls and flushes stay what they are live:
#
# poll_interval_seconds 10 -> 1
# flush_interval_seconds 30 (the default) -> 3
#
# grace_seconds is stream time and stays. idle_close_seconds only fires when
Expand All @@ -18,6 +17,7 @@ commands:
sql: |
CREATE TABLE IF NOT EXISTS posts (
time_us BIGINT,
event_time TIMESTAMPTZ,
commit STRUCT(
operation TEXT,
record STRUCT(langs TEXT[])
Expand Down Expand Up @@ -51,26 +51,21 @@ tables:
# land. After idle_close_seconds with nothing arriving, every open
# bucket closes.
#
# late_rows: drop, because the sink replaces. A row for a bucket that
# already closed is discarded and counted in window_late_rows_total.
# reemit would run emit_sql over the late rows alone, because the
# bucket's other rows were deleted when it closed, and the upsert would
# replace the bucket's count with theirs. validate refuses that pairing.
# No lateness: a post for a bucket that already closed is refused
# before the handler and counted in window_late_rows_total.
window:
time_column: bucket
size_seconds: 60
grace_seconds: 60
idle_close_seconds: 60
late_rows: drop
poll_interval_seconds: 1
emit_sql: |
SELECT bucket, lang, sum(posts)::INTEGER AS posts
FROM closed
GROUP BY ALL
sink:
# A keyed sink replaces the row a (bucket, lang) already identifies.
# The engine deletes a window only after the sink accepts it, so a
# crash between the two republishes that window on the next poll,
# crash between the two republishes that window on the next pass,
# and the republish lands harmlessly. Each flush costs the batch: a
# COPY into a staging table and a server-side merge, in one
# transaction. The target needs PRIMARY KEY (bucket, lang) or a
Expand All @@ -91,14 +86,17 @@ pipeline:
type: websocket
websocket:
uri: "{{ SQLFLOW_JETSTREAM_URI|default('wss://jetstream2.us-east.bsky.network/subscribe?wantedCollections=app.bsky.feed.post') }}"
event_time:
path: time_us
format: unix_us

handler:
type: handlers.StructuredBatch
table: posts
sql: |
INSERT INTO posts_per_minute_by_lang
SELECT
time_bucket(INTERVAL '1 minute', to_timestamp(time_us / 1000000)) AS bucket,
time_bucket(INTERVAL '1 minute', event_time) AS bucket,
coalesce(commit.record.langs[1], 'unknown') AS lang,
count(*) AS posts
FROM posts
Expand Down
Loading
Loading