1. What is Flink?1. Что такое Flink?
Flink is a tool that computes over data as it flows, continuously, in real time. Instead of waiting for all the data and then crunching it (batch), Flink reacts to each event the moment it arrives and keeps results up to date.
Flink — инструмент, который вычисляет по данным прямо в потоке, непрерывно, в реальном времени. Вместо того чтобы ждать все данные и потом считать (батч), Flink реагирует на каждое событие в момент прихода и держит результаты свежими.
2. Kafka vs Flink (and batch vs streaming)2. Kafka vs Flink (и батч vs стриминг)
| Kafka | Flink | |
|---|---|---|
| Job | Move & store the stream | Compute over the stream |
| Think of it as | The pipe / the log | The calculator on the pipe |
| Typical role | Source and destination | Reads from Kafka, writes back to Kafka or a DB |
| Kafka | Flink | |
|---|---|---|
| Задача | Перемещать и хранить поток | Вычислять по потоку |
| Как думать | Труба / лог | Калькулятор на трубе |
| Типичная роль | Источник и приёмник | Читает из Kafka, пишет обратно в Kafka или БД |
Batch vs streaming: batch processing waits for a complete chunk (e.g. "all of yesterday") and computes it once — simple and accurate. Streaming computes continuously on data that never stops — fast and live, but you must handle data that's incomplete or out of order. Flink is built for the streaming case.
Батч vs стриминг: батч ждёт целый кусок (например «весь вчерашний день») и считает его один раз — просто и точно. Стриминг считает непрерывно по данным, которые не заканчиваются — быстро и «вживую», но надо уметь работать с неполными или неупорядоченными данными. Flink создан как раз для стриминга.
3. Event time vs processing time3. Время события vs время обработки
Two different clocks matter:
Важны двое разных «часов»:
- Event time = when the thing actually happened (stamped on the event).
- Processing time = when Flink happened to see it.
- Event time = когда событие реально произошло (отметка на самом событии).
- Processing time = когда Flink его увидел.
They differ because events get delayed: a phone goes offline, the network lags, a retry happens. For correct results you almost always want to group by event time, not by when you received it.
Они расходятся, потому что события задерживаются: телефон ушёл в офлайн, сеть тормозит, случился ретрай. Для корректных результатов почти всегда нужно группировать по event time, а не по моменту получения.
4. Windows: chopping a stream into buckets4. Окна: режем поток на корзины
A stream never ends, so "count everything" never finishes. Instead you compute over windows — time buckets like "every 5 minutes".
Поток не кончается, поэтому «посчитать всё» никогда не завершится. Вместо этого считают по окнам — временны́м корзинам вроде «каждые 5 минут».
- Tumbling windows: fixed, non-overlapping (00:00–00:05, 00:05–00:10, …). Each event in exactly one.
- Sliding windows: overlapping (a 5-min window every 1 min) → smoother "last 5 minutes" metric.
- Session windows: grouped by activity with gaps (a user's burst of clicks = one session).
- Tumbling (кувыркающиеся) окна: фиксированные, без нахлёста (00:00–00:05, 00:05–00:10, …). Каждое событие ровно в одном.
- Sliding (скользящие) окна: с нахлёстом (5-мин окно каждую 1 мин) → более плавная метрика «за последние 5 минут».
- Session (сессионные) окна: группируют по активности с паузами (всплеск кликов пользователя = одна сессия).
5. Watermarks: "I've seen enough to close this window"5. Watermarks: «я видел достаточно, чтобы закрыть окно»
Since events can arrive late, how does Flink know a 10:00–10:05 window is done and can be totalled? It uses a watermark — a moving marker that says "I believe I've now seen all events up to time X." When the watermark passes 10:05, that window fires.
Раз события могут опаздывать, как Flink понимает, что окно 10:00–10:05 завершено и можно подбить итог? Он использует watermark — движущийся маркер, говорящий «кажется, я уже видел все события до момента X». Когда watermark проходит 10:05, окно срабатывает.
Very-late events (after the watermark) can be dropped, sent to a side channel, or used to correct the number later.
Сильно опоздавшие события (после watermark) можно отбросить, отправить в отдельный канал или позже поправить ими число.
6. State & checkpoints: Flink remembers6. Состояние и чекпоинты: Flink помнит
To keep a running count or join two streams, Flink must remember things between events — that memory is called state (e.g. "current total per user"). Flink manages this state for you, even when it's huge.
Чтобы вести бегущий счёт или джойнить два потока, Flink должен помнить вещи между событиями — эта память называется состояние (state) (например «текущая сумма по пользователю»). Flink управляет этим состоянием за тебя, даже когда оно огромное.
Periodically Flink saves a snapshot of all its state — a checkpoint. If a machine crashes, it restarts from the last checkpoint and continues as if nothing happened. No lost progress.
Периодически Flink сохраняет снимок всего своего состояния — checkpoint. Если машина падает, он рестартует с последнего чекпоинта и продолжает, будто ничего не случилось. Прогресс не теряется.
7. When to use Flink7. Когда использовать Flink
- You need answers now, not tomorrow: live dashboards, fraud/abuse detection, alerting.
- You need windowed metrics over event time (per-minute counts, sessionization).
- You need to join or enrich streams on the fly (e.g. add user info to each click).
- If you only need daily/hourly aggregates and can wait, plain batch (e.g. Spark) is simpler and cheaper — don't reach for streaming unless latency matters.
- Ответы нужны сейчас, а не завтра: живые дашборды, детекция фрода/абьюза, алертинг.
- Нужны оконные метрики по event time (поминутный счёт, сессионизация).
- Нужно джойнить или обогащать потоки на лету (например, добавлять инфу о пользователе к каждому клику).
- Если нужны только дневные/часовые агрегаты и можно подождать, обычный батч (например Spark) проще и дешевле — не берись за стриминг, если задержка не критична.
8. Mini-glossary8. Мини-словарь
| Word | Plain meaning |
|---|---|
| Stream processing | Computing continuously on data as it flows. |
| Event time | When the event actually happened. |
| Processing time | When the system saw the event. |
| Window | A time bucket you compute over (e.g. every 5 min). |
| Watermark | "I've seen all events up to here" — when to close a window. |
| State | What Flink remembers between events (running totals, etc.). |
| Checkpoint | A saved snapshot of state for crash recovery. |
| Exactly-once | Each event affects the result once, even after a crash. |
| Слово | Простой смысл |
|---|---|
| Потоковая обработка | Непрерывные вычисления по данным в потоке. |
| Event time | Когда событие реально произошло. |
| Processing time | Когда система увидела событие. |
| Окно (window) | Временна́я корзина, по которой считаешь (например, каждые 5 мин). |
| Watermark | «Я видел все события до этого момента» — когда закрывать окно. |
| Состояние (state) | Что Flink помнит между событиями (бегущие суммы и т.п.). |
| Checkpoint | Сохранённый снимок состояния для восстановления после сбоя. |
| Exactly-once | Каждое событие влияет на результат один раз, даже после сбоя. |
9. Quick self-check9. Быстрая самопроверка
Answer in your head, then tap to flip.Ответь про себя, потом нажми, чтобы перевернуть.