Streaming joins concept page with three flavors: stream-stream windowed join (click+impression in 5-min window with symmetric hash join), stream-table enrichment (async lookup with Redis cache + LRU), and temporal join (FX rate as-of trade time using versioned table). Includes state explosion warning scenario and interval join. Two ADRs: stream-stream join state cost vs precomputed enrichment, and async lookup vs co-partitioned KTable.
Streaming join — это объединение двух (или более) бесконечных потоков событий по общему ключу. В отличие от batch SQL, где обе таблицы статичны, в стриминге обе стороны живут во времени: правая запись может прийти через миллисекунды или часы после левой, может вообще не прийти, может прийти раньше watermark'а и опоздать. Поэтому streaming join — это не «как SQL JOIN, только быстро», а отдельный класс алгоритмов с window, state retention и watermark-driven cleanup.
Три канонических флавора, которые покрывает диаграмма:
«Stream-Stream = храни обе стороны в hashmap, пробуй match на каждом event, чисти state по watermark. Stream-Table = lookup по ключу в дешёвом storage (Redis / RocksDB / GlobalKTable). Temporal = lookup в версионированной таблице, выбирая version с max(ts) ≤ event_time.»
Главный риск общий для всех: state explosion. Join по природе stateful. Capacity-планирование = sum(events × avg_size × window) per partition. Hot key + длинное окно убивают job.
Левая колонка — источники: три event-потока (clicks 200K/s, impressions 500K/s, trades 50K/s, все partitioned по join-ключу) и три референса (compacted user-table 5GB, versioned FX-rates, Redis lookup-cache). Правая колонка — Flink job parallelism=32 с четырьмя операторами:
FOR SYSTEM_TIME AS OF оператор, держит versioned-state по symbol.Sinks внизу справа: kafka-attributed (matched pairs), kafka-enriched (click + profile), kafka-trades-fx (trade × historical rate).
На flink-join повешены два ADR: state-cost vs precomputed enrichment (выбор stream-stream symmetric hash join) и async-lookup vs co-partitioned KTable vs GlobalKTable (выбор async + Redis).
Канонический stream-stream join: impression@10
буферизуется в right-side hashmap task 7 (hash(user=42) % 32 = 7). Click@10:00 приходит на тот же task (co-partitioned), пробит right-side, находит impression, эмит attribution pair. Дальше показаны три кромки: orphan click без impression (буферизуется в left-side, дропается после window expiry), late impression (event_time < watermark — dropped по watermark policy), и cleanup через RocksDB compactor по TTL =window + allowed_lateness.
Enrichment click'а user-профилем без локального state per task: CDC populate Redis с TTL=60s → Flink async I/O делает non-blocking lookup. Hot user (95% случаев) попадает в локальный LRU-кеш (10K записей per task, ~0.1ms). Cold user — async-call в Redis (~3ms p99), при этом task не блокируется, продолжает обрабатывать другие events. Stale-data window до 60s принят как trade-off для ad targeting.
Versioned join: rates приходят сериями (1.10@10
→ 1.12@11 → 1.15@12). Trade@12 ищетmax(rate.event_time) ≤ 12:00 — выбирает 1.12, не current 1.15. Это корректный исторический rate (отличается от current на ~3%). Anti-pattern: non-versioned table silently вернула бы current rate — financial reporting сломался бы тихо.
Диагностика и фикс: 1-hour window + hot user 100 clicks/s = 360K events × 200 bytes = 72MB per hot user, 1000 hot users = 72GB на task 11 (hot partition). RocksDB write stall, backpressure, OOM risk при compaction, checkpoint restore 5min. Фиксы по убыванию импакта: shorten window 1h → 5min (state /12), filter bot traffic (state /3), composite key (user_id, hour) для spread hot key, parallelism 32 → 128 (state per task /4).
Flink-специфичный interval join: «click BETWEEN imp.time AND imp.time + 30s». Уже чем windowed join (state /10), потому что асимметричный bound — buffer только impressions, click пробит и забыт. Click@10:00
матчится с imp@10 (в bound). Click@10:01 — out of bound, NO MATCH. Используется для attribution с жёсткими time-constraint'ами.ADR-001: Stream-stream symmetric hash join вместо precomputed enrichment. Контекст: ad attribution, нужно match'ить click с impression в окне 30min. State = 252GB raw на 30-min window, ÷32 partitions ≈ 8GB per task. Альтернативы: (1) Kafka Streams precomputed enrichment через KTable — теряет windowed semantics, нельзя retroactively match late impressions; (2) External lookup Redis — миллионы concurrent connections; (3) Managed streaming SQL (Materialize) — vendor lock-in, дорого at scale. Выбрано: Flink symmetric hash join с RocksDB state. Корректная семантика (windowed + watermark + late-event policy), state manageable через TTL и incremental checkpoints в S3. Accepted: ops complexity, recovery ~5min с большим state, ~0.5% late-event drop.
ADR-002: Async lookup с Redis caching вместо co-partitioned KTable / GlobalKTable. Контекст: enrich clickstream user-профилем (5GB, 10M users, ~100 updates/min). Альтернативы: (1) Co-partitioned KTable — re-partition shuffle на 700K events/s expensive, tight coupling source layout; (2) GlobalKTable — full replicate 5GB × 32 instances = 160GB total memory; (3) Async I/O + Redis LRU. Выбрано: вариант 3. State footprint минимален, recovery быстрая, 1000 concurrent lookups per task в Flink AsyncIO, локальный LRU даёт 95% hit rate (Pareto distribution). Accepted: 5% lookups идут network p99 ~3ms, eventual consistency 60s — OK для ad targeting, не для financial.
FOR SYSTEM_TIME, interval join, regular join) + DataStream. Канон для stream-stream и temporal.Production: Uber — Flink для trip+payment+driver join. Netflix — stream-stream для personalization. Lyft — temporal joins для surge pricing. Airbnb — Kafka Streams для booking enrichment.
-/+ records, к которым consumer не готов.