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.carrierFROM ordersINNER JOIN shipments ON orders.order_id = shipments.order_idWHERE 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_purchaseFROM page_viewsLEFT JOIN purchases ON page_views.user_id = purchases.user_idWHERE 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_urlCharacteristics:
- 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.amountFROM clicksINNER JOIN purchases ON clicks.user_id = purchases.user_idGROUP 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_convertFROM clicksINNER JOIN conversions ON clicks.user_id = conversions.user_idWHERE conversions.event_time BETWEEN clicks.event_time - INTERVAL '5' MINUTE AND clicks.event_time + INTERVAL '5' MINUTECharacteristics:
- 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:
- Enable state TTL
- Reduce window size
- Increase parallelism to distribute state
- Use smaller join keys
Late Events Dropped
Symptom: Missing join results
Solutions:
- Increase allowed lateness
- Adjust watermark strategy
- Monitor late event metrics
Low Throughput
Symptom: Slow join processing
Solutions:
- Increase parallelism
- Optimize join predicates
- Use faster serialization (bincode vs JSON)
- Reduce window size
Best Practices
- Always Monitor Watermarks: Watermark lag directly impacts latency
- Set Appropriate TTL: Prevent unbounded state growth
- Test with Realistic Data: Validate with production-like volumes
- Use Session Windows Carefully: Can accumulate large state
- Optimize Join Keys: Simple, evenly distributed keys perform best
- Partition Smartly: High-cardinality partition keys enable parallelism
References
Version: 1.0