Dual-Stream Correlation and Absence Alerting (Stream-Stream JOIN)
# Dual-Stream Correlation and Absence Alerting (Stream-Stream JOIN)
This case demonstrates the stream-stream JOIN (WITHIN syntax, since v1.3.0): correlating two live streams by matching keys plus time proximity, across four scenarios — INNER signal corroboration, LEFT absence alerting, event-time alignment, and a 3-stream cascade. Syntax and semantics: Stream-Stream JOIN.
Running these cases differs from the other cases
A stream-stream JOIN has multiple input streams: EmitSync is rejected; the FROM side may use plain Emit, while every other stream must be fed by name via EmitTo. The correct pattern:
ssql := streamsql.New()
_ = ssql.Execute(sql) // the case SQL
ssql.AddSink(func(rows []map[string]interface{}) { /* collect results */ })
_ = ssql.EmitTo("streamName", row) // feed rows by stream name
ssql.Stop() // on stop, unmatched LEFT pending rows are NULL-complemented (Flush)
2
3
4
5
# Scenario A: Dual-Source Signal Corroboration (INNER JOIN)
# Business Goal
One device reports over two channels: temperature and vibration. A lone high reading may be interference — only when both appear within 30 seconds is it a real anomaly. Join the two channels with INNER JOIN; WHERE filters each side to the anomalous readings.
# 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
2
3
4
5
# Input and Output
tempStream:
{"deviceId": "d1", "temperature": 80.5} // t+0.0s
{"deviceId": "d2", "temperature": 76.0} // t+0.3s
{"deviceId": "d3", "temperature": 24.0} // t+0.6s
2
3
vibrationStream:
{"deviceId": "d1", "vibration": 35.2} // t+0.9s
{"deviceId": "d2", "vibration": 41.0} // t+1.2s
{"deviceId": "d9", "vibration": 33.0} // t+1.5s
{"deviceId": "d1", "vibration": 28.0} // t+1.8s
2
3
4
Output (as executed):
{"deviceId":"d1","temp":80.5,"vib":35.2}
{"deviceId":"d2","temp":76,"vib":41}
2
# Behavior Notes
- d1 and d2 satisfy both side filters within the window → emitted;
d3's temperature of 24 is removed by WHERE;d9has only a vibration row and no temperature row for the same key within 30s → nothing; - d1's second vibration reading of 28.0 ≤ 30 is filtered by WHERE;
- Matched rows are not removed: if another qualifying vibration for d1 arrives within the 30s window, it joins with the buffered 80.5 again (rows stay matchable for the whole window);
- The right row lives under its alias: the output shape is left fields flattened +
v.vibrationreferencing the right row.
# Scenario B: Command Sent but No ACK (LEFT JOIN Absence Detection)
# Business Goal
The platform sends commands to devices; a healthy device ACKs within seconds. What we want is "no ACK within 10 seconds" — unfindable with WHERE alone (a missing row has no record). LEFT JOIN's "no match by window close → emit a NULL row" semantics does it, with IS NULL surfacing the absent ones.
# SQL
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
# Input and Output
cmdStream:
{"cmdId": "c1", "deviceId": "d1"} // t+0.0s
{"cmdId": "c2", "deviceId": "d2"} // t+0.3s
2
ackStream:
{"cmdId": "c1", "ackCode": 0} // t+0.6s (c1's ACK arrives on time)
Output (as executed; appears around t+10~15s — the sweeper period is WITHIN/2, see notes):
{"ackCode":null,"cmdId":"c2","deviceId":"d2"}
# Behavior Notes
- c1's ACK arrives within the window → matched, and
ackCode=0failsIS NULL, so WHERE drops it — on-time commands produce zero output; - c2's ACK never comes → at window close a row with
r = {}(empty object) is emitted;r.ackCodeis NULL and passes WHERE, producing the alert; - In processing-time mode the sweeper advances watermarks to wall clock, so the complement arrives after ≈ WITHIN plus one sweep period (WITHIN/2); pending rows not yet complemented are flushed at
Stop(); - A per-row marker guarantees each pending row is complemented exactly once — no duplicate alerts.
# Scenario C: Vehicle CAN × GPS Trajectory (Event-Time Alignment)
# Business Goal
The CAN stream and the GPS stream each carry device-stamped event timestamps. The two channels have different network latency, so arrival order is unreliable and alignment must use event time: declare the field with WITH (TIMESTAMP = 'ts'), correlate by VIN, and accept only pairs whose event times are within 5 seconds.
# SQL
SELECT can.vin, can.speed, gps.lon
FROM canStream AS can
JOIN gpsStream AS gps WITHIN 5 SECONDS
ON can.vin = gps.vin
WITH (TIMESTAMP = 'ts')
2
3
4
5
# Input and Output
canStream (ts is a millisecond epoch):
{"vin": "V1", "speed": 72, "ts": 1698700000000} // event time t0
{"vin": "V2", "speed": 90, "ts": 1698705000000} // event time t0+5000s (another device V2, independent timeline)
2
gpsStream:
{"vin": "V1", "lon": 116.1, "ts": 1698700000200} // 0.2s from V1's CAN row → matches
{"vin": "V1", "lon": 116.2, "ts": 1698700006000} // 6.0s from V1's CAN row > WITHIN → no match
{"vin": "V2", "lon": 117.0, "ts": 1698705001000} // 1.0s from V2's CAN row → matches
2
3
Output (as executed):
{"lon":116.1,"speed":72,"vin":"V1"}
{"lon":117,"speed":90,"vin":"V2"}
2
# Behavior Notes
- Matching is by event-time distance, not arrival order: the second GPS row (lon=116.2) arrives right after the CAN row, but its event time is 6 seconds away > WITHIN 5s, so it does not match — the core difference from processing-time mode;
- Timestamp values are unit-auto-detected (ns/μs/ms/s epoch): this case uses millisecond epochs (~1.7×10¹²), recognized as ms. Note that values below 10⁹ are compared as raw nanoseconds — hand-made sequence numbers (1, 2, 3, …) would always look "adjacent" to each other; prefer real epoch timestamps in production;
- In event-time mode, if one side stops sending, its watermark freezes and LEFT NULLs are delayed;
WITH (IDLETIMEOUT = '30s')advances the idle side's watermark to wall clock (at the cost of dropping rows that arrive late after recovery). - Rows later than WITHIN are dropped outright (no match, no alert; counted as
join_stage{i}_late_dropped): if one side's device clock drifts beyond WITHIN, that side's rows are systematically discarded — leave clock-drift slack in WITHIN, or watch the metric viaGetStats().
# Scenario D: 3-Stream Cascade (a ⋈ b) LEFT ⋈ c
# Business Goal
A two-stage join over three streams: device state stream a joins parameter stream b on key k (within 5s); the merged condition then joins the alert stream c on k2 (within 3s). The second stage is a LEFT: conditions without a paired alert are still emitted (vc = NULL), meaning "this condition has no associated alert".
# SQL
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
# Input and Output
aStream:
{"k": "k1", "va": 1} // t+0.0s
{"k": "k2", "va": 2} // t+0.9s
2
bStream:
{"k": "k1", "k2": "x", "vb": "B"} // t+0.3s
{"k": "k2", "k2": "y", "vb": "B2"} // t+1.2s
2
cStream:
{"k2": "x", "vc": "C1"} // t+0.6s
{"k2": "z", "vc": "C3"} // t+1.5s (no k2=y partner)
2
Output (as executed; the second row is complemented about 3~4.5s after the condition row):
{"va":1,"vb":"B","vc":"C1"}
{"va":2,"vb":"B2","vc":null}
2
# Behavior Notes
- N streams = N−1 binary JOINs chained left-deep: the (a⋈b) merged row feeds stage two's LEFT side, and ON refers to the
b.k2prefix directly; - Stage-two LEFT: the
{k2:"z"}alert row finds nok2=ycondition (buffered only, never emitted); the condition row (va=2, vb=B2) closes its window unmatched →c = {}is emitted withc.vcNULL; - The intermediate merged row's timestamp (event-time mode) = max(the two matched rows' ts) — a compound event happens at its latest constituent row, keeping the distance to c bounded;
- Each stage's WITHIN is independent (this stage's 3s does not touch stage one's 5s); metrics are separated by stage prefix (
join_stage1_*/join_stage2_*); Stop()flushes in cascade order: upper-stage pending rows flow into the next stage before it flushes.