Apache Flink — the friendly introApache Flink — дружелюбное введение

What stream processing is, in plain words. Read Kafka first if you haven't — Flink usually sits on top of it. When comfortable, switch to Advanced. Что такое потоковая обработка — простыми словами. Сначала прочитай Kafka, если ещё нет — Flink обычно стоит поверх неё. Когда станет комфортно, переключайся на Advanced.
🟢 BeginnerНовичок 🔴 AdvancedПродвинутый

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 реагирует на каждое событие в момент прихода и держит результаты свежими.

AnalogyАналогия If Kafka is the conveyor belt carrying items, Flink is the worker standing at the belt who inspects, counts, and combines items as they pass — keeping a live running tally on a whiteboard. Если Kafka — это конвейерная лента, везущая предметы, то Flink — работник у ленты, который осматривает, считает и комбинирует предметы по мере их проезда, ведя живой счёт на доске.
One lineОдной строкой Flink is a stream-processing engine: ongoing computation (counts, sums, windows, joins) over never-ending streams of events. Flink — это движок потоковой обработки: непрерывные вычисления (счёт, суммы, окна, джоины) над бесконечными потоками событий.

2. Kafka vs Flink (and batch vs streaming)2. Kafka vs Flink (и батч vs стриминг)

KafkaFlink
JobMove & store the streamCompute over the stream
Think of it asThe pipe / the logThe calculator on the pipe
Typical roleSource and destinationReads from Kafka, writes back to Kafka or a DB
KafkaFlink
ЗадачаПеремещать и хранить потокВычислять по потоку
Как думатьТруба / логКалькулятор на трубе
Типичная рольИсточник и приёмникЧитает из 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, а не по моменту получения.

AnalogyАналогия A letter is written on the 1st (event time) but arrives on the 5th (processing time). To count "letters written in January" correctly, you use the postmark date, not the delivery date. Письмо написано 1-го (event time), но пришло 5-го (processing 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 (сессионные) окна: группируют по активности с паузами (всплеск кликов пользователя = одна сессия).
events: • • • • • • • • • • tumbling: [ 0–5 ][ 5–10 ][ 10–15 ] ← each event lands in exactly one bucket

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, окно срабатывает.

It's a trade-off dialЭто «ручка» компромисса Wait longer (looser watermark) → catch more late events, but results are slower. Wait less (tighter) → faster results, but you might drop some latecomers. You tune it to your data's real lateness. Ждать дольше (свободнее watermark) → поймаешь больше опоздавших, но результат медленнее. Ждать меньше (туже) → быстрее результат, но можешь отбросить опоздавших. Настраиваешь под реальную «опоздалость» данных.

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. Если машина падает, он рестартует с последнего чекпоинта и продолжает, будто ничего не случилось. Прогресс не теряется.

AnalogyАналогия Checkpoints are like save points in a video game: if you die, you reload the last save and keep going instead of starting over. Чекпоинты — как точки сохранения в видеоигре: если умер, грузишь последнее сохранение и продолжаешь, а не начинаешь сначала.
Exactly-once, simplyExactly-once, просто Checkpoints + a destination that can accept updates safely (idempotent sink) let Flink give exactly-once results: even after a crash and replay, nothing is counted twice. Чекпоинты + приёмник, который умеет безопасно принимать обновления (идемпотентный sink), позволяют Flink давать exactly-once результат: даже после падения и перечитки ничто не посчитается дважды.

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) проще и дешевле — не берись за стриминг, если задержка не критична.
Honest noteЧестно Flink is powerful but operationally heavier (always-on, stateful). Many teams do streaming only where it's truly needed and keep everything else in batch. Flink мощный, но операционно тяжелее (всегда включён, stateful). Многие команды делают стриминг только там, где он реально нужен, а остальное держат в батче.

8. Mini-glossary8. Мини-словарь

WordPlain meaning
Stream processingComputing continuously on data as it flows.
Event timeWhen the event actually happened.
Processing timeWhen the system saw the event.
WindowA time bucket you compute over (e.g. every 5 min).
Watermark"I've seen all events up to here" — when to close a window.
StateWhat Flink remembers between events (running totals, etc.).
CheckpointA saved snapshot of state for crash recovery.
Exactly-onceEach 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.Ответь про себя, потом нажми, чтобы перевернуть.

Kafka vs Flink in one line?Kafka vs Flink одной строкой?
tapнажми
Kafka moves & stores the stream; Flink computes over it. Flink usually reads from Kafka and writes back.Kafka перемещает и хранит поток; Flink вычисляет по нему. Flink обычно читает из Kafka и пишет обратно.
Why prefer event time over processing time?Почему event time лучше processing time?
tapнажми
Because events arrive late/out of order. Grouping by when it happened (postmark) gives correct results; grouping by when you saw it does not.Потому что события приходят с опозданием/не по порядку. Группировка по тому, когда произошло (штемпель), даёт верный результат; по тому, когда увидел — нет.
What is a window?Что такое окно?
tapнажми
A time bucket to compute over (e.g. every 5 min), because an endless stream can't be "summed all at once." Tumbling / sliding / session.Временна́я корзина для вычислений (например, каждые 5 мин), потому что бесконечный поток нельзя «просуммировать разом». Tumbling / sliding / session.
What does a watermark decide?Что решает watermark?
tapнажми
When a window is "complete enough" to fire. Looser = catch more late data but slower; tighter = faster but may drop latecomers.Когда окно «достаточно полно», чтобы сработать. Свободнее = ловишь больше опоздавших, но медленнее; туже = быстрее, но можешь отбросить опоздавших.
What are checkpoints for?Зачем нужны чекпоинты?
tapнажми
Saved snapshots of Flink's state (its memory). After a crash it reloads the last checkpoint and continues — like a video-game save point. Enables exactly-once.Сохранённые снимки состояния Flink (его памяти). После сбоя он грузит последний чекпоинт и продолжает — как точка сохранения в игре. Даёт exactly-once.
When NOT to use Flink?Когда НЕ брать Flink?
tapнажми
When you only need daily/hourly aggregates and can wait — plain batch (Spark) is simpler and cheaper. Use streaming only when latency truly matters.Когда нужны только дневные/часовые агрегаты и можно подождать — обычный батч (Spark) проще и дешевле. Стриминг — только когда задержка реально важна.