Files
nis2/project-tasks/p14-b-flink-rescaling-recovery.md
T
2026-09-06 19:26:12 +03:00

5.3 KiB
Raw Blame History

P14-B. Flink: rescaling и восстановление

  • Версия и дата проверки: 1.0, 05.09.2026.
  • Статус: готово к назначению.

Статья и исходные материалы

  • Основная статья: Paris Carbone и соавт. — Apache Flink: Stream and Batch Processing in a Single Engine. IEEE Data Engineering Bulletin 2015, 11 страниц.
  • Кратко о статье: Apache Flink объединяет потоковую и пакетную обработку в общем распределённом движке с состоянием операторов. Согласованные контрольные точки и savepoint позволяют восстанавливать вычисление и переносить состояние при изменении параллелизма. В этом проекте исследуется корректность перераспределения состояния и зависимость паузы восстановления от его размера и структуры.
  • Почему результат актуален: исходная статья описывает раннюю архитектуру, но система стала устойчивой отраслевой платформой и продолжает активно развиваться; Flink 2.3.0 выпущен в июне 2026 года. Проект использует современную фиксированную версию и проверяет сохраняющийся механизм распределённых контрольных точек.
  • Артефакты и данные: apache/flink под Apache-2.0; доступны актуальная стабильная ветвь, отдельная LTS-ветвь и официальный локальный режим без обязательного облака. Зафиксированная основная ревизия: apache/flink@81389aca7136 (Apache-2.0).

Обязательный результат

  • Проверяемый вопрос или утверждение: Контрольная точка или savepoint позволяет корректно перераспределить состояние при изменении параллелизма, однако время восстановления и пауза обработки зависят от размера и распределения состояния.
  • Технический результат: Построить key-partitioned stateful-конвейер Flink с повторяемым источником, автоматической проверкой эквивалентности результатов и автоматизированным сохранением состояния. Реализовать остановку и восстановление с прежним и изменённым параллелизмом.
  • Эксперимент: Для не менее чем трёх размеров состояния сравнить обычный restart без изменения параллелизма как baseline и восстановление с rescaling. Измерить паузу, полное время восстановления, размер сохранения, throughput после запуска, пропуски, дубликаты и ошибки агрегатов; выполнить по три серии.
  • Границы выводов: корректность и пауза восстановления измеряются для выбранного stateful-конвейера, размеров состояния и параллелизма на одной машине; результат не характеризует эластичность производственного кластера и удалённого хранилища.
  • Ресурсный профиль: одна машина, CPU, 4–8 ГБ памяти и локальная файловая система; Docker и облако необязательны. При нехватке памяти уменьшаются параллелизм и состояние, но сохраняются отдельные процессы, контрольные точки, отказ и автоматическая проверка корректности.

Содержательные направления

  • нагрузка, состояние и эталон корректности.
  • автоматизация savepoint/restart/rescaling.
  • измерения восстановления, повторные серии и новый режим.

Возможное продолжение

Исследовать перекос ключей, incremental checkpoints, разные state backend, автоматический выбор параллелизма или последовательность нескольких rescaling.