Skip to content

Windowed Joins Guide

Windowed Joins Guide

Overview

Windowed joins enable correlation of events from multiple streams within time-based windows. HeliosDB Streaming supports all Apache Flink window types with optimized state management and encryption.

Join Types

1. Tumbling Window Join

Non-overlapping, fixed-size time windows.

Use Case: Hourly order-shipment correlation

SQL Example:

SELECT
orders.order_id,
orders.customer_id,
shipments.tracking_number,
shipments.carrier
FROM orders
INNER JOIN shipments
ON orders.order_id = shipments.order_id
WHERE
orders.event_time BETWEEN
TUMBLE_START(orders.event_time, INTERVAL '1' HOUR) AND
TUMBLE_END(orders.event_time, INTERVAL '1' HOUR)

Characteristics:

  • Simple, predictable windows
  • Low memory usage (one window at a time)
  • Aligned boundaries (0:00, 1:00, 2:00, etc.)
  • ❌ May miss events at boundary

2. Sliding Window Join

Overlapping windows with configurable slide interval.

Use Case: Real-time metrics with 10-minute windows, 1-minute updates

SQL Example:

SELECT
page_views.page_url,
COUNT(*) AS view_count,
AVG(purchases.amount) AS avg_purchase
FROM page_views
LEFT JOIN purchases
ON page_views.user_id = purchases.user_id
WHERE
page_views.event_time BETWEEN
HOP_START(page_views.event_time, INTERVAL '1' MINUTE, INTERVAL '10' MINUTE) AND
HOP_END(page_views.event_time, INTERVAL '1' MINUTE, INTERVAL '10' MINUTE)
GROUP BY page_views.page_url

Characteristics:

  • Smooth, continuous updates
  • Events can appear in multiple windows
  • ❌ Higher memory usage (overlapping state)
  • ❌ More CPU for window management

3. Session Window Join

Dynamic windows based on inactivity gaps.

Use Case: User session attribution (clicks → purchases)

SQL Example:

SELECT
clicks.user_id,
STRING_AGG(clicks.page_url, ', ') AS click_path,
purchases.product_id,
purchases.amount
FROM clicks
INNER JOIN purchases
ON clicks.user_id = purchases.user_id
GROUP BY SESSION(clicks.event_time, INTERVAL '30' MINUTE)

Characteristics:

  • Natural session boundaries
  • Handles variable-length sessions
  • Gap-based window closing
  • ❌ Unpredictable window sizes
  • ❌ Late events can reopen sessions

4. Interval Join

Time-based correlation with configurable before/after intervals.

Use Case: Click-to-conversion attribution within ±5 minutes

SQL Example:

SELECT
clicks.click_id,
clicks.ad_campaign,
conversions.conversion_id,
conversions.revenue,
(conversions.event_time - clicks.event_time) AS time_to_convert
FROM clicks
INNER JOIN conversions
ON clicks.user_id = conversions.user_id
WHERE
conversions.event_time BETWEEN
clicks.event_time - INTERVAL '5' MINUTE AND
clicks.event_time + INTERVAL '5' MINUTE

Characteristics:

  • Precise time-based correlation
  • Symmetric or asymmetric intervals
  • Low latency (no window wait)
  • ❌ Requires buffering both streams

Join Semantics

Inner Join

Only matching records from both streams.

Left Outer Join

All records from left stream, with nulls for non-matching right records.

Full Outer Join

All records from both streams, with nulls for non-matches.

Optimization Techniques

1. Join Reordering

Join smaller streams first to minimize state size.

2. State TTL (Time-To-Live)

Automatically remove expired join state to prevent memory bloat.

3. Key Distribution

Ensure even key distribution to avoid hotspots.

4. Watermark Configuration

Proper watermark handling prevents late data issues.

Performance Metrics

Key Metrics to Monitor

Critical Metrics:

  • State Size: Memory footprint of join state
  • Watermark Lag: Time between event time and watermark
  • Late Events: Events arriving after watermark
  • Join Throughput: Events processed per second

Troubleshooting

High Memory Usage

Symptom: OOM errors during joins

Solutions:

  1. Enable state TTL
  2. Reduce window size
  3. Increase parallelism to distribute state
  4. Use smaller join keys

Late Events Dropped

Symptom: Missing join results

Solutions:

  1. Increase allowed lateness
  2. Adjust watermark strategy
  3. Monitor late event metrics

Low Throughput

Symptom: Slow join processing

Solutions:

  1. Increase parallelism
  2. Optimize join predicates
  3. Use faster serialization (bincode vs JSON)
  4. Reduce window size

Best Practices

  1. Always Monitor Watermarks: Watermark lag directly impacts latency
  2. Set Appropriate TTL: Prevent unbounded state growth
  3. Test with Realistic Data: Validate with production-like volumes
  4. Use Session Windows Carefully: Can accumulate large state
  5. Optimize Join Keys: Simple, evenly distributed keys perform best
  6. Partition Smartly: High-cardinality partition keys enable parallelism

References


Version: 1.0