RuleGo RuleGo
🏠Home
  • Quick Start
  • Rule Chain
  • Standard Components
  • Extension Components
  • Custom Components
  • Visualization
  • RuleGo-Server
  • AOP
  • Trigger
  • Advanced Topics
  • Performance
  • Standard Components
  • Extension Components
  • Custom Components
  • Components Marketplace
  • Overview
  • Quick Start
  • Routing
  • DSL
  • API
  • Options
  • Components
🔥Editor (opens new window)
  • RuleGo Editor (opens new window)
  • RuleGo Server (opens new window)
  • StreamSQL
  • AI Agent Framework
  • GFlow Approval Workflow (opens new window)
  • TPCLAW Agent Platform (opens new window)
  • Github (opens new window)
  • Gitee (opens new window)
  • Changelog (opens new window)
  • English
  • 简体中文
🏠Home
  • Quick Start
  • Rule Chain
  • Standard Components
  • Extension Components
  • Custom Components
  • Visualization
  • RuleGo-Server
  • AOP
  • Trigger
  • Advanced Topics
  • Performance
  • Standard Components
  • Extension Components
  • Custom Components
  • Components Marketplace
  • Overview
  • Quick Start
  • Routing
  • DSL
  • API
  • Options
  • Components
🔥Editor (opens new window)
  • RuleGo Editor (opens new window)
  • RuleGo Server (opens new window)
  • StreamSQL
  • AI Agent Framework
  • GFlow Approval Workflow (opens new window)
  • TPCLAW Agent Platform (opens new window)
  • Github (opens new window)
  • Gitee (opens new window)
  • Changelog (opens new window)
  • English
  • 简体中文

广告采用随机轮播方式显示 ❤️成为赞助商
  • Quick Start

  • Rule Chain

  • Standard Components

  • Extension Components

  • Custom Components

  • Components marketplace

  • Visualization

  • AOP

  • Trigger

  • Advanced Topic

  • Agent Framework

  • RuleGo-Server

  • FAQ

  • Endpoint Module

  • Support

  • StreamSQL

    • Overview
    • Quick Start
    • Core Concepts
    • SQL Reference
    • API Reference
    • RuleGo Integration
    • Schema Validation
    • Advanced Examples
    • Pattern Matching (CEP)
    • Stream-Stream JOIN
    • functions

    • case-studies

      • Case Studies Overview
      • Stream-Table JOIN Metadata Enrichment
      • Session Window and Device Online Analysis
      • Change Data Capture Case Study
      • Sliding Window and Continuous Detection
      • Data Filtering and Transformation
      • IoT Temperature Alerting and Metrics Aggregation
      • Device Fault Pattern Recognition (MATCH_RECOGNIZE)
      • Dual-Stream Correlation and Absence Alerting (Stream-Stream JOIN)
        • Scenario A: Dual-Source Signal Corroboration (INNER JOIN)
          • Business Goal
          • SQL
          • Input and Output
          • Behavior Notes
        • Scenario B: Command Sent but No ACK (LEFT JOIN Absence Detection)
          • Business Goal
          • SQL
          • Input and Output
          • Behavior Notes
        • Scenario C: Vehicle CAN × GPS Trajectory (Event-Time Alignment)
          • Business Goal
          • SQL
          • Input and Output
          • Behavior Notes
        • Scenario D: 3-Stream Cascade (a ⋈ b) LEFT ⋈ c
          • Business Goal
          • SQL
          • Input and Output
          • Behavior Notes
目录

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)
1
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
1
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
1
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
1
2
3
4

Output (as executed):

{"deviceId":"d1","temp":80.5,"vib":35.2}
{"deviceId":"d2","temp":76,"vib":41}
1
2

# Behavior Notes

  • d1 and d2 satisfy both side filters within the window → emitted; d3's temperature of 24 is removed by WHERE; d9 has 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.vibration referencing 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
1
2
3
4
5

# Input and Output

cmdStream:

{"cmdId": "c1", "deviceId": "d1"}          // t+0.0s
{"cmdId": "c2", "deviceId": "d2"}          // t+0.3s
1
2

ackStream:

{"cmdId": "c1", "ackCode": 0}              // t+0.6s (c1's ACK arrives on time)
1

Output (as executed; appears around t+10~15s — the sweeper period is WITHIN/2, see notes):

{"ackCode":null,"cmdId":"c2","deviceId":"d2"}
1

# Behavior Notes

  • c1's ACK arrives within the window → matched, and ackCode=0 fails IS 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.ackCode is 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')
1
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)
1
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
1
2
3

Output (as executed):

{"lon":116.1,"speed":72,"vin":"V1"}
{"lon":117,"speed":90,"vin":"V2"}
1
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 via GetStats().

# 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
1
2
3
4

# Input and Output

aStream:

{"k": "k1", "va": 1}                       // t+0.0s
{"k": "k2", "va": 2}                       // t+0.9s
1
2

bStream:

{"k": "k1", "k2": "x", "vb": "B"}          // t+0.3s
{"k": "k2", "k2": "y", "vb": "B2"}         // t+1.2s
1
2

cStream:

{"k2": "x", "vc": "C1"}                    // t+0.6s
{"k2": "z", "vc": "C3"}                    // t+1.5s (no k2=y partner)
1
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}
1
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.k2 prefix directly;
  • Stage-two LEFT: the {k2:"z"} alert row finds no k2=y condition (buffered only, never emitted); the condition row (va=2, vb=B2) closes its window unmatched → c = {} is emitted with c.vc NULL;
  • 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.
Edit this page on GitHub (opens new window)
Last Updated: 2026/09/26, 14:57:53
Device Fault Pattern Recognition (MATCH_RECOGNIZE)

← Device Fault Pattern Recognition (MATCH_RECOGNIZE)

Theme by Vdoing | Copyright © 2023-2026 RuleGo Team | Apache 2.0 License

  • 跟随系统
  • 浅色模式
  • 深色模式
  • 阅读模式