Stream-Stream JOIN
# Stream-Stream JOIN
The stream-stream JOIN correlates two live streams by matching keys and time proximity (since v1.3.0): ksqlDB-style WITHIN syntax, the time-proximity semantics of Flink Interval Join / Kafka Streams JoinWindows, LEFT JOIN absence detection, and left-deep cascades for 3+ streams. State is bounded (WITHIN is the retention period), matching the edge memory profile.
# Stream-Table vs Stream-Stream
StreamSQL has two kinds of JOIN, distinguished by the presence of WITHIN:
| Stream-Table JOIN | Stream-Stream JOIN | |
|---|---|---|
| Correlates | stream × static/cached metadata table | stream × stream (both live) |
| Syntax marker | JOIN table ON ... (no WITHIN) | JOIN stream WITHIN 30 SECONDS ON ... |
| Time constraint | none (lookup against the current table snapshot) | must satisfy \|L.ts − R.ts\| ≤ WITHIN |
| Typical use | enrich with location/model attributes | dual-source signal corroboration, command-ACK absence alerting |
| How to feed | RegisterTable | EmitTo by stream name |
| Docs | Case: Stream-Table JOIN Enrichment | this page + Case: Dual-Stream Correlation and Absence Alerting |
If WITHIN is missing, the JOIN target is treated as an unregistered table and the engine reports join table %q is not registered with a hint to add WITHIN.
# Quick Start
SELECT s.deviceId, s.temperature AS temp, v.vibration AS vib
FROM tempStream AS s
JOIN vibrationStream AS v WITHIN 30 SECONDS
ON s.deviceId = v.deviceId
WHERE s.temperature > 75 AND v.vibration > 30
2
3
4
5
ssql := streamsql.New()
_ = ssql.Execute(`...the SQL above...`)
ssql.AddSink(func(rows []map[string]interface{}) {
for _, r := range rows {
fmt.Println(r)
}
})
// Feed both sides by stream name (= the name after FROM / JOIN, case-sensitive)
_ = ssql.EmitTo("tempStream", map[string]interface{}{"deviceId": "d1", "temperature": 80.5})
_ = ssql.EmitTo("vibrationStream", map[string]interface{}{"deviceId": "d1", "vibration": 35.2})
// On Stop, unmatched pending rows of a LEFT JOIN are NULL-complemented before exit (Flush)
ssql.Stop()
2
3
4
5
6
7
8
9
10
11
12
13
14
15
Single-stream code does not change: Emit(row) is equivalent to EmitTo(the FROM stream, row).
# Running It in a RuleGo Component (no Go code needed)
RuleGo users skip the Go API — host the JOIN SQL in an x/streamAggregator node, configure streamKey as the stream-name source, and route multiple upstream feeds into the same node (each message's metadata carries its own stream name):
{"id": "join1", "type": "x/streamAggregator",
"configuration": {
"sql": "SELECT s.deviceId, s.temperature AS temp, v.vibration AS vib FROM tempStream AS s JOIN vibrationStream AS v WITHIN 30 SECONDS ON s.deviceId = v.deviceId",
"streamKey": "streamName"}}
2
3
4
Matched rows flow out via the stream_event relation (Success is the raw-data passthrough); to aggregate the join output further, use the two-stage composition — see Two-Stage Join-then-Aggregate below.
# Syntax
SELECT s.deviceId, s.temperature, v.vibration
FROM tempStream AS s
JOIN vibrationStream AS v WITHIN 30 SECONDS
ON s.deviceId = v.deviceId
WHERE ...
2
3
4
5
| Syntax | Rule |
|---|---|
| WITHIN | Required — the marker of a stream-stream JOIN; WITHIN 30 SECONDS / WITHIN (30 SECONDS) / WITHIN '30s' / WITHIN 100 MS are all accepted |
| WITHIN position | after the JOIN table name (ksqlDB position) or after the ON clause — exactly one of the two; writing both is an error |
| JOIN types | INNER JOIN / LEFT JOIN; RIGHT / FULL / CROSS are compile-time errors |
| ON | equality conditions; composite keys chained with AND (ON a.x = b.x AND a.y = b.y); non-equality is rejected |
| Event time | WITH (TIMESTAMP = 'ts') names the event-time field (numeric values auto-detected as ns/μs/ms/s; values below 10⁹ are compared as raw nanoseconds — hand-made sequence numbers would always look "adjacent"; use real epoch timestamps); without it, arrival time is used |
| Idle advance | WITH (IDLETIMEOUT = '30s'): in event-time mode, a side idle beyond the timeout has its watermark advanced to wall clock |
| Output shape | same shape as stream-table JOIN: left fields flattened + right row under its alias (v.vibration); an alias.field projection (either side) without AS may emit both the original and the de-prefixed key (f.score → f.score and score; the same applies to cascade outputs and LEFT NULL rows, same as stream-table JOIN) — always prefer an explicit AS; SELECT * yields {left fields…, s: {left row}, v: {right row}} |
| Overflow | follows the drop/block strategy configured via the Go API WithOverflowStrategy (not a SQL WITH option); expand does not apply to JOIN inputs (downgraded to drop with a warning at build time) |
WITHIN in a JOIN is not WITHIN in CEP
MATCH_RECOGNIZE ... WITHIN '1h' (Pattern Matching) bounds the active lifetime of a pattern match; JOIN ... WITHIN 30 SECONDS (this page) bounds the timestamp distance between two rows. The two never coexist in one query.
# Matching Semantics
Unified time model: each row's ts is the event time (taken mandatorily when TIMESTAMP is configured; rows missing the field are dropped and counted) or the arrival time (processing-time). Both modes share the same matching and eviction path.
First time here? Focus on four rows
Match predicate, Emission, Late rows, LEFT absence — these four decide what output you will see; the remaining rows are operational/edge details to consult when troubleshooting.
| Semantic | Rule |
|---|---|
| Match predicate | \|L.ts − R.ts\| ≤ WITHIN (symmetric interval, same as Kafka Streams JoinWindows / Flink Interval Join) |
| Emission | emit as soon as a row arrives and matches (append-only); rows are not removed after matching and stay matchable for the whole window (one-to-many Cartesian happens naturally) |
| NULL join keys | missing/NULL keys normalize to a <nil> key, and NULL == NULL is treated as equal (inherited from stream-table JOIN, unlike the SQL standard); keep join keys non-null to avoid surprises |
| Watermark & retention | each side keeps a maxSeenTs minimal watermark; rows with maxSeenTs − row.ts > WITHIN expire and no longer participate in matching (they leave via eviction / LEFT NULL complement) |
| Late rows | in event-time mode, rows later than WITHIN are dropped, counted as join_stage{i}_late_dropped with a throttled warning (never silent); leave headroom in WITHIN for clock drift. MaxOutOfOrderness/AllowedLateness combined with JOIN is a compile-time error (no out-of-orderness margin in v1) |
| Idle streams | a side that stops sending freezes its watermark and delays LEFT NULLs; in processing-time mode the sweeper advances watermarks to wall clock; in event-time mode configure IDLETIMEOUT to do the same (trade-off: rows arriving after recovery are dropped as late) |
| LEFT absence | a left row with no match is parked as pending; when the window closes (watermark / idle advance / Stop-Flush) exactly one NULL row is emitted (an empty object under the right alias), guaranteed once per row |
| Lifecycle | Stop() flushes in cascade order: upper-stage pending LEFT rows flow into the next stage before each stage flushes |
# 3+ Streams: Left-Deep Cascade
N streams = N−1 binary JOINs chained left-deep (SQL standard associativity; mainstream engines do the same). Each stage has its own WITHIN; INNER/LEFT freely combined:
SELECT a.va, b.vb, c.vc
FROM aStream AS a
JOIN bStream AS b WITHIN 5 SECONDS ON a.k = b.k
LEFT JOIN cStream AS c WITHIN 3 SECONDS ON b.k2 = c.k2
2
3
4
- The intermediate merged row (a⋈b) feeds the next stage's LEFT side; ON refers to existing
a.x/b.yprefixes directly; - Intermediate-row timestamp = max(the two matched rows' ts) (event-time mode): a compound event "happens" at its latest constituent row, keeping the next stage's distance bounded (
|c−b| ≤ W2and|c−a| ≤ W1+W2); in processing-time mode it is the arrival time at that stage; - Memory gates apply per stage; metrics carry the stage prefix
join_stage{i}_*(i starts at 1).
# Scenario Examples
① Dual-source signal corroboration (INNER) — same device "high temperature then high vibration" within 30s is a real anomaly:
SELECT s.deviceId, s.temperature AS temp, v.vibration AS vib
FROM tempStream AS s
JOIN vibrationStream AS v WITHIN 30 SECONDS
ON s.deviceId = v.deviceId
WHERE s.temperature > 75 AND v.vibration > 30
2
3
4
5
② Command sent but no ACK (LEFT absence detection) — commands with no ACK within 10 seconds (NULL-complemented, then surfaced by IS NULL):
SELECT c.cmdId, c.deviceId, r.ackCode
FROM cmdStream AS c
LEFT JOIN ackStream AS r WITHIN 10 SECONDS
ON c.cmdId = r.cmdId
WHERE r.ackCode IS NULL
2
3
4
5
③ Vehicle CAN × GPS trajectory (event time) — CAN and GPS streams carry their own event timestamps; correlate by VIN within 5 seconds:
SELECT can.vin, can.speed, gps.lon, gps.lat
FROM canStream AS can
JOIN gpsStream AS gps WITHIN 5 SECONDS
ON can.vin = gps.vin
WITH (TIMESTAMP = 'ts')
2
3
4
5
④ Access control: card + face (composite key) — badge swipes and face captures joined on gate ID + user ID; f.score is the face-match score:
SELECT d.gateId, d.userId, f.score
FROM cardStream AS d
JOIN faceStream AS f WITHIN 5 SECONDS
ON d.gateId = f.gateId AND d.userId = f.userId
2
3
4
Full input/output walkthroughs (actually executed) for ①②③ and a 3-stream cascade: Case: Dual-Stream Correlation and Absence Alerting.
# Two-Stage Join-then-Aggregate
A WITHIN JOIN cannot share one SQL with GROUP BY / aggregate functions (compile-time error). When you need "correlate then aggregate" (per-device windowed counts, averages over matched rows, etc.), use the two-stage composition: the JOIN node's stream_event output feeds a downstream aggregation node — the recommended pattern in the RuleGo component context:
[source A route] ─┐
├─→ A: x/streamAggregator (JOIN SQL, streamKey=streamName) ─stream_event→ B: x/streamAggregator (aggregation SQL) ─stream_event→ alerts/storage
[source B route] ─┘
2
3
Rule chain JSON (the connection type is the key):
{
"metadata": {
"nodes": [
{"id": "joinA", "type": "x/streamAggregator", "name": "dual-stream join",
"configuration": {
"sql": "SELECT s.deviceId, s.temperature AS temp, v.vibration AS vib FROM tempStream AS s JOIN vibrationStream AS v WITHIN 30 SECONDS ON s.deviceId = v.deviceId WHERE s.temperature > 75 AND v.vibration > 30",
"streamKey": "streamName"}},
{"id": "aggB", "type": "x/streamAggregator", "name": "windowed aggregation",
"configuration": {
"sql": "SELECT deviceId, COUNT(*) AS cnt FROM joinOut GROUP BY deviceId, CountingWindow(10)"}}
],
"connections": [
{"fromId": "joinA", "toId": "aggB", "type": "stream_event"}
]
}
}
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
Two notes before copying:
streamKeynames the metadata key that carries the target stream name — each input message in the example carriesstreamName=tempStreamorstreamName=vibrationStream(both the key and the value are case-sensitive; the value must match the stream name after FROM/JOIN in the SQL; a${}expression that resolves to the stream name directly also works, e.g."${msg.stream}");joinOutin B'sFROM joinOutis B's own input stream name and can be anything — B's data arrives via A'sstream_eventrelation, independent of that name.
Two rules to follow (verified by integration tests):
- Wire A→B via the
stream_eventrelation only. A'sSuccesspassthrough carries the raw telemetry; connecting it to B would push un-correlated raw rows into B's windows (polluting counts/averages); - For an event-time window on B, A's SQL must project the timestamp (e.g.
SELECT ..., s.ts), and B declaresWITH (TIMESTAMP = 'ts')— JOIN matched rows do not carry a timestamp automatically; without the projection, B aggregates by arrival time (CountingWindowor the default window).
B receives matched-row arrays as stream_event payloads (one message per match) and feeds them row by row into its own stream; B's window/aggregation semantics are identical to a single-stream query. Single-SQL join-then-aggregate remains on the roadmap (see the validation checklist below); two-stage is the recommended pattern today.
# Go API
All new APIs are additive; the 7 red-line APIs (New/Execute/Emit/EmitSync/AddSink/Stop/IsAggregationQuery) keep their signatures and behavior unchanged.
ssql := streamsql.New(
streamsql.WithJoinMaxKeys(20000), // optional: buffered key cap per side per stage (LRU, default 10000)
streamsql.WithJoinMaxRows(500000), // optional: buffered row cap per side per stage (default 0 = unbounded)
)
err := ssql.Execute(sql) // WITHIN is mandatory
err = ssql.EmitTo("cmdStream", row) // feed by name; unknown name returns an error (lists known names)
ok := ssql.IsStreamJoinQuery() // true for stream-stream JOIN queries (for RuleGo component routing)
2
3
4
5
6
7
Emit(row)=EmitTo(the FROM stream, row);EmitSyncreturns an error on JOIN queries (same as CEP queries);- Results come through
AddSink/ToChannel, exactly like single-stream queries.
# Resource Bounds and Memory Estimation
Three gates, each backing up the previous:
- WITHIN is mandatory — the retention period is the time dimension of the state bound;
WithJoinMaxKeys(n)— LRU cap on normalized keys per side per stage (default 10000). Eviction of pending LEFT rows does not emit their NULL rows (counted inkeys_evicted+ 10s throttled warning); under overload, which row gets evicted is not deterministic (two independent inputs, scheduler-dependent) — with maxKeys ≥ peak active keys there is no eviction and no impact;WithJoinMaxRows(n)— row-buffer cap per side per stage (default 0 = unbounded); over-limit drops the newest row (counted inrows_dropped+ throttled warning).
Memory estimation (per side per stage):
peak rows ≈ input rate(msg/s) × WITHIN(s) × active-key ratio
peak memory ≈ peak rows × ~440 B/row (same magnitude as window buffers)
cascade total ≈ Σ per-stage (left + right) peaks; stage i's left input rate = stage i-1's output rate
2
3
Measured references: two streams at 1k msg/s each, WITHIN=30s → ~26MB both sides; 200k distinct keys pushed with maxKeys=10000 → buffer converges at the gate, total allocation bounded at ~109MB.
# Observability
GetStats() / Metrics() expose stage-prefixed metrics (i = stage number, starting at 1):
| Metric | Meaning |
|---|---|
join_stage{i}_matches_emitted | matched rows emitted |
join_stage{i}_left_timeout_emitted | LEFT NULL rows emitted at window close |
join_stage{i}_late_dropped | rows dropped for lateness / missing TIMESTAMP field |
join_stage{i}_input_dropped | dropped on input chan overflow (drop strategy or block timeout) |
join_stage{i}_rows_dropped | rows dropped by the maxRows gate |
join_stage{i}_keys_evicted | keys evicted by the maxKeys LRU |
join_stage{i}_buffered_rows | currently buffered rows (gauge) |
join_stage{i}_left_watermark / _right_watermark | both sides' watermarks (gauge; debug idle streams) |
Every drop category also emits a 10s-throttled Warn log.
# Compile-Time Validation
The following combinations make Execute fail immediately — no silent no-ops:
| Combination | Disposition |
|---|---|
| WITHIN JOIN + GROUP BY / window / aggregate functions | error: single-SQL aggregate-after-join not supported (use the two-stage composition, see "Two-Stage Join-then-Aggregate" above; actual error text: WITHIN JOIN cannot be combined with GROUP BY/window/aggregation (...aggregate the JOIN output in a downstream query)) |
| WITHIN JOIN + HAVING / DISTINCT / LIMIT | error: the direct tail does not consume these clauses |
| WITHIN JOIN + analytic functions (OVER) | error: evaluation order vs two-stream matching is undefined |
WITHIN JOIN + MaxOutOfOrderness / AllowedLateness > 0 | error: no out-of-orderness margin in v1; IdleTimeout is allowed and effective |
| WITHIN JOIN mixed with stream-table JOIN | error: v1.3.x roadmap ("join then enrich") |
| WITHIN JOIN + unnest | error: row expansion is direct-path only |
| RIGHT / FULL / CROSS JOIN | parse error |
| WITHIN written twice (after table and after ON) | parse error |
| WITHIN ≤ 0 | parse error |
| FROM/JOIN stream name collision | Execute error |
EmitSync on a JOIN query | error |
EmitTo with an unknown stream name | error (lists known names) |
# Mainstream Engine Comparison
| StreamSQL | Flink | ksqlDB | Kafka Streams | |
|---|---|---|---|---|
| Syntax | JOIN s2 WITHIN 30 SECONDS ON ... | ON a.ts BETWEEN b.ts±30s | JOIN s2 WITHIN 30 SECONDS ON ... | JoinWindows.of(...) (Java API) |
| Time semantics | event time (TIMESTAMP) / arrival time by default | event time + watermark | stream time | stream time |
| Unmatched LEFT | NULL at window close | NULL when watermark passes the interval | NULL at window close | NULL at window close |
| Update streams / retract | ❌ (append-only end to end) | ✅ | ✅ | ✅ |
| 3+ streams | ✅ left-deep cascade, per-stage WITHIN | ✅ arbitrary plans | ❌ | ❌ |
| State retention | retention (WITHIN) + LRU/row-count gates | state TTL | window close | retention |
# Roadmap (v1.3.x / v2)
- RIGHT / FULL JOIN;
ON a.ts BETWEEN b.ts-30s AND b.ts+30sinterval-predicate sugar (parsed down to the same WITHIN config);- Out-of-orderness margin config (
MaxOutOfOrdernesscombined with JOIN); - Stream-table enrichment after a stream-stream JOIN (mixing WITHIN JOIN with stream-table JOIN);
- RuleGo component example for feeding two streams (component repo release cadence; "multiple sources into one node" is an existing RuleGo capability).