Содержание курса
Модуль 1. Введение в потоковую дата-инженерию и Apache Flink
8 уроков
1.
О курсе: как проходить обучение и канал «Логово Дата-Инженера»
↗
2.
Зачем дата-инженеру Apache Flink и какие задачи он решает
↗
3.
Поток данных против пакетной обработки: первая ментальная модель
↗
4.
Bounded и unbounded streams: конечные и бесконечные данные
↗
5.
Streaming analytics, ETL и event-driven приложения во Flink
↗
6.
Flink, Spark Structured Streaming, Kafka Streams и Beam: где…
↗
7.
Карта курса: от первого job до продакшн-пайплайна
↗
8.
Как выполнять практику и проверять себя без кластера
↗
Модуль 2. Рабочая среда
7 уроков
1.
Версии Flink 2.x, Java 17 и что важно знать про LTS-ветки
↗
2.
Установка JDK, Maven и проверка окружения
↗
3.
Структура Maven-проекта Flink и зависимости provided
↗
4.
Первый DataStream job в IDE: StreamExecutionEnvironment
↗
5.
Локальный mini cluster и запуск из main
↗
6.
Flink CLI: run, list, cancel и что происходит при отправке job
↗
7.
Первый взгляд на Web UI и логи локального запуска
↗
Модуль 3. Архитектура Flink
7 уроков
1.
Client, JobManager и TaskManager: кто за что отвечает
↗
2.
JobGraph, ExecutionGraph и операторный DAG простыми словами
↗
3.
Tasks, subtasks и parallelism: как работа делится на части
↗
4.
Task slots: единица ресурсов и почему слот не равен CPU
↗
5.
Operator chaining: почему несколько операторов могут стать…
↗
6.
Shuffle и network exchange: где начинается дорогое…
↗
7.
Application, job и deployment mode в современной Flink 2.x-линии
↗
Модуль 4. DataStream API: первые преобразования
7 уроков
1.
Что такое DataStream и почему его нельзя просто распечатать…
↗
2.
Source, transformation, sink: анатомия Flink-программы
↗
3.
map и flatMap: преобразование событий
↗
4.
filter и простая очистка входного потока
↗
5.
keyBy: логическое разделение потока по ключу
↗
6.
reduce и простая инкрементальная агрегация
↗
7.
execute и ленивость: когда программа реально стартует
↗
Модуль 5. Типы данных, функции и сериализация
7 уроков
1.
Java POJO, record и Tuple: какие типы Flink понимает лучше
↗
2.
Почему сериализация важна для производительности и state
↗
3.
MapFunction, FlatMapFunction и lambda: что выбрать в учебном…
↗
4.
RichFunction: open, close, RuntimeContext и аккуратные ресурсы
↗
5.
Configuration и параметры job без захардкоженных значений
↗
6.
Ошибки типов и сериализации: как читать первые stack trace
↗
7.
Чистые функции в stream processing: меньше скрытого состояния
↗
Модуль 6. Источники и приёмники данных
7 уроков
1.
Source API: от коллекции и DataGen к реальным источникам
↗
2.
FileSource: чтение файлов как bounded stream
↗
3.
DataGen и Print connector для учебных экспериментов
↗
4.
Sink API: куда уходит результат и почему println не продакшн
↗
5.
FileSink: запись файлов и rolling policy
↗
6.
Параллельные sink-и и порядок записей: чего не обещает Flink
↗
7.
Мини-пайплайн: читаем события, чистим, пишем результат
↗
Модуль 7. Время в потоках
7 уроков
1.
Processing time против event time: две разные правды о времени
↗
2.
Timestamp assigner: откуда Flink берёт время события
↗
3.
Watermark: отметка прогресса во времени события
↗
4.
Out-of-order события и bounded out-of-orderness
↗
5.
Idleness: почему один молчащий partition может остановить время
↗
6.
Поздние события: что значит late data
↗
7.
Отладка времени: как увидеть, почему окно не закрывается
↗
Модуль 8. Окна: tumbling, sliding, session и агрегаты
7 уроков
1.
Зачем нужны окна и почему бесконечный поток нельзя просто…
↗
2.
Tumbling windows: непересекающиеся интервалы
↗
3.
Sliding windows: пересекающиеся интервалы и цена дублирования
↗
4.
Session windows: группировка по активности пользователя
↗
5.
Window assigner, trigger, evictor: кто управляет окном
↗
6.
ReduceFunction, AggregateFunction и ProcessWindowFunction
↗
7.
Allowed lateness и side output для поздних событий
↗
Модуль 9. Keyed state: состояние как сердце Flink
7 уроков
1.
Stateless против stateful operators: когда память становится…
↗
2.
Keyed state и key groups: как Flink распределяет состояние
↗
3.
ValueState: последний статус, счётчик и простые накопления
↗
4.
ListState и MapState: когда одного значения недостаточно
↗
5.
ReducingState и AggregatingState: состояние с агрегирующей…
↗
6.
State TTL: удаление старого состояния без ручной уборки
↗
7.
Типичные ошибки state: общий mutable object, null и рост без…
↗
Модуль 10. ProcessFunction, таймеры и нестандартная логика
7 уроков
1.
ProcessFunction: когда map/filter/window уже не хватает
↗
2.
KeyedProcessFunction и доступ к keyed state
↗
3.
Processing-time timers: отложенные действия по времени обработки
↗
4.
Event-time timers: логика, привязанная к времени события
↗
5.
Side outputs: отдельный поток для ошибок и поздних данных
↗
6.
Broadcast state: правила и справочники, приходящие отдельным…
↗
7.
Async I/O: обогащение событий внешним сервисом без блокировки
↗
Модуль 11. Checkpoints, savepoints и state backends
7 уроков
1.
Fault tolerance во Flink: что именно восстанавливается после…
↗
2.
Checkpoint barrier и согласованное состояние
↗
3.
Checkpointing mode: exactly-once и at-least-once без мифов
↗
4.
Checkpoint config: интервалы, timeout, min pause и…
↗
5.
Savepoint: управляемая точка остановки, миграции и апгрейда
↗
6.
HashMapStateBackend, EmbeddedRocksDBStateBackend и…
↗
7.
Incremental checkpoints и большая state: где выигрыш и где цена
↗
Модуль 12. Kafka и end-to-end гарантии доставки
7 уроков
1.
KafkaSource: bootstrap servers, topics, group id и offsets…
↗
2.
Десериализация событий: строки, JSON и собственные схемы
↗
3.
KafkaSink: сериализация и запись результата
↗
4.
DeliveryGuarantee: none, at-least-once и exactly-once
↗
5.
Транзакционный Kafka sink и transactional id prefix
↗
6.
Offsets, checkpoints и безопасный рестарт job
↗
7.
Практический пайплайн Kafka -> Flink -> Kafka
↗
Модуль 13. Flink SQL: потоковые таблицы и SQL Client
7 уроков
1.
Flink SQL как язык потоковых и пакетных запросов
↗
2.
SQL Client и первый SELECT над DataGen
↗
3.
CREATE TABLE: connectors, formats и WITH options
↗
4.
Dynamic tables: append, retract и upsert простыми словами
↗
5.
SELECT, WHERE и computed columns
↗
6.
INSERT INTO как запуск непрерывного запроса
↗
7.
EXPLAIN: читаем план SQL-запроса
↗
Модуль 14. Время, окна и joins во Flink SQL
7 уроков
1.
Time attributes: rowtime, proctime и WATERMARK в DDL
↗
2.
Window TVF: TUMBLE, HOP и CUMULATE
↗
3.
Window aggregation и группировка по window_start/window_end
↗
4.
Deduplication и Top-N в потоковом SQL
↗
5.
Regular join и проблема бесконечного состояния
↗
6.
Interval join: соединение событий по временному диапазону
↗
7.
Temporal join: справочники, версии и поток фактов
↗
Модуль 15. Table API, UDF и мост к DataStream
7 уроков
1.
TableEnvironment и StreamTableEnvironment: точки входа
↗
2.
Table API против SQL: когда код удобнее строки запроса
↗
3.
Преобразование DataStream в Table и обратно
↗
4.
ScalarFunction: простая пользовательская функция
↗
5.
TableFunction и разворачивание одной строки в несколько
↗
6.
AggregateFunction: пользовательская агрегация
↗
7.
Граница ответственности: где оставить SQL, а где перейти в…
↗
Модуль 16. Форматы, озеро данных и CDC
7 уроков
1.
JSON, CSV, Avro, Protobuf и Raw: как выбирать формат
↗
2.
Schema Registry и Confluent Avro: зачем нужна внешняя схема
↗
3.
Debezium format: changelog из базы данных
↗
4.
Flink CDC: когда нужен отдельный коннектор для изменений
↗
5.
Upsert Kafka: ключи, tombstone и компактированные топики
↗
6.
FileSystem connector и потоковая запись в data lake
↗
7.
Iceberg, Paimon и lakehouse-таблицы: место Flink в озере
↗
Модуль 17. Деплой: standalone, Docker, YARN и Kubernetes
7 уроков
1.
Session mode, application mode и per-job: как выбирать режим…
↗
2.
Упаковка fat jar и зависимости connector-ов
↗
3.
Standalone cluster: JobManager, TaskManager и конфигурация
↗
4.
Docker Compose для локального кластера Flink
↗
5.
YARN deployment: когда Flink живёт рядом с Hadoop
↗
6.
Native Kubernetes и Flink Kubernetes Operator
↗
7.
High availability: ZooKeeper, Kubernetes HA и recovery storage
↗
Модуль 18. Наблюдаемость, отладка и эксплуатация
7 уроков
1.
Web UI: job graph, subtasks, exceptions и accumulators
↗
2.
Логи JobManager и TaskManager: где искать причину падения
↗
3.
Metrics: throughput, busy time, backpressure и checkpoint…
↗
4.
Backpressure: как распознать и где искать узкое место
↗
5.
Debugging windows и event time в реальном job
↗
6.
Restart strategies и failure rate: как Flink перезапускает…
↗
7.
Production readiness checklist для Flink job
↗
Модуль 19. Производительность и тюнинг Flink job
7 уроков
1.
Parallelism, max parallelism и rescaling без сюрпризов
↗
2.
Key skew: перекос ключей и способы борьбы
↗
3.
Network buffers, shuffle и operator chaining
↗
4.
State backend tuning: heap, RocksDB, managed memory
↗
5.
Checkpoint tuning под backpressure и большую state
↗
6.
Async I/O, cache и rate limiting внешних сервисов
↗
7.
Чек-лист ревью производительности Flink-пайплайна
↗
Модуль 20. Итоговый проект и дорога к senior data engineer
7 уроков
1.
Постановка итогового проекта: события заказов в реальном времени
↗
2.
Проектируем схемы событий, ключи и контракты данных
↗
3.
Реализация DataStream pipeline: очистка, keyBy, state и таймеры
↗
4.
Реализация SQL-слоя: витрина агрегатов и upsert sink
↗
5.
Тесты и локальный прогон итогового Flink-проекта
↗
6.
Деплой, checkpointing, savepoint и план безопасного релиза
↗
7.
Дорожная карта дальше: Flink internals, CDC, lakehouse и open…
↗