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
      • Stream-Table vs Stream-Stream
      • Quick Start
        • Running It in a RuleGo Component (no Go code needed)
      • Syntax
      • Matching Semantics
      • 3+ Streams: Left-Deep Cascade
      • Scenario Examples
      • Two-Stage Join-then-Aggregate
      • Go API
      • Resource Bounds and Memory Estimation
      • Observability
      • Compile-Time Validation
      • Mainstream Engine Comparison
      • Roadmap (v1.3.x / v2)
    • functions

    • case-studies

目录

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
1
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()
1
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"}}
1
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 ...
1
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
1
2
3
4
  • The intermediate merged row (a⋈b) feeds the next stage's LEFT side; ON refers to existing a.x / b.y prefixes 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| ≤ W2 and |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
1
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
1
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')
1
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
1
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] ─┘
1
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"}
    ]
  }
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16

Two notes before copying:

  • streamKey names the metadata key that carries the target stream name — each input message in the example carries streamName=tempStream or streamName=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}");
  • joinOut in B's FROM joinOut is B's own input stream name and can be anything — B's data arrives via A's stream_event relation, independent of that name.

Two rules to follow (verified by integration tests):

  1. Wire A→B via the stream_event relation only. A's Success passthrough carries the raw telemetry; connecting it to B would push un-correlated raw rows into B's windows (polluting counts/averages);
  2. For an event-time window on B, A's SQL must project the timestamp (e.g. SELECT ..., s.ts), and B declares WITH (TIMESTAMP = 'ts') — JOIN matched rows do not carry a timestamp automatically; without the projection, B aggregates by arrival time (CountingWindow or 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)
1
2
3
4
5
6
7
  • Emit(row) = EmitTo(the FROM stream, row);
  • EmitSync returns 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:

  1. WITHIN is mandatory — the retention period is the time dimension of the state bound;
  2. 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 in keys_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;
  3. WithJoinMaxRows(n) — row-buffer cap per side per stage (default 0 = unbounded); over-limit drops the newest row (counted in rows_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
1
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+30s interval-predicate sugar (parsed down to the same WITHIN config);
  • Out-of-orderness margin config (MaxOutOfOrderness combined 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).
Edit this page on GitHub (opens new window)
Last Updated: 2026/09/26, 14:57:53
Pattern Matching (CEP)
Aggregate Functions

← Pattern Matching (CEP) Aggregate Functions→

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

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