# P14-B. Flink: rescaling и восстановление - **Версия и дата проверки:** 1.1, 07.09.2026. - **Статус:** готово к назначению. ## Статья и исходные материалы - **Основная статья:** Paris Carbone и соавт. — [Apache Flink: Stream and Batch Processing in a Single Engine](https://asterios.katsifodimos.com/assets/publications/flink-deb.pdf). IEEE Data Engineering Bulletin 2015, 11 страниц. - **Кратко о статье:** Apache Flink объединяет потоковую и пакетную обработку в общем распределённом движке с состоянием операторов. Согласованные контрольные точки и savepoint позволяют восстанавливать вычисление и переносить состояние при изменении параллелизма. В этом проекте исследуется корректность перераспределения состояния и зависимость паузы восстановления от его размера и структуры. - **Почему результат актуален:** исходная статья описывает раннюю архитектуру, но система стала устойчивой отраслевой платформой и продолжает активно развиваться; [Flink 2.3.0](https://flink.apache.org/2026/06/25/apache-flink-2.3.0-release-announcement/) выпущен в июне 2026 года. Проект использует современную фиксированную версию и проверяет сохраняющийся механизм распределённых контрольных точек. - **Артефакты и данные:** [apache/flink](https://github.com/apache/flink) под Apache-2.0; доступны актуальная стабильная ветвь, отдельная LTS-ветвь и [официальный локальный режим](https://nightlies.apache.org/flink/flink-docs-stable/docs/getting-started/local_installation/) без обязательного облака. Зафиксированная основная ревизия: `apache/flink@81389aca7136` (Apache-2.0). - **Что уже предоставляет артефакт:** Flink предоставляет сохранение состояния, restart, rescaling и средства наблюдения. Их разрешено использовать для выполнения и восстановления конвейера. ## Обязательный результат - **Проверяемый вопрос или утверждение:** Контрольная точка или savepoint позволяет корректно перераспределить состояние при изменении параллелизма, однако время восстановления и пауза обработки зависят от размера и распределения состояния. - **Технический результат:** Построить key-partitioned stateful-конвейер Flink с повторяемым источником, автоматической проверкой эквивалентности результатов и автоматизированным сохранением состояния. Реализовать остановку и восстановление с прежним и изменённым параллелизмом. - **Обязательное приращение команды:** Создать описанный key-partitioned-конвейер, повторяемый источник и независимую проверку результатов; автоматизировать сохранение и восстановление с прежним и изменённым параллелизмом. Сравнить предусмотренные три размера состояния и разобрать паузу и корректность восстановления. - **Эксперимент:** Для не менее чем трёх размеров состояния сравнить обычный restart без изменения параллелизма как baseline и восстановление с rescaling. Измерить паузу, полное время восстановления, размер сохранения, throughput после запуска, пропуски, дубликаты и ошибки агрегатов; выполнить по три серии. - **Границы выводов:** корректность и пауза восстановления измеряются для выбранного stateful-конвейера, размеров состояния и параллелизма на одной машине; результат не характеризует эластичность производственного кластера и удалённого хранилища. - **Ресурсный профиль:** одна машина, CPU, 4–8 ГБ памяти и локальная файловая система; Docker и облако необязательны. При нехватке памяти уменьшаются параллелизм и состояние, но сохраняются отдельные процессы, контрольные точки, отказ и автоматическая проверка корректности. ## Содержательные направления - нагрузка, состояние и эталон корректности. - автоматизация savepoint/restart/rescaling. - измерения восстановления, повторные серии и новый режим. ## Возможное продолжение Исследовать перекос ключей, incremental checkpoints, разные state backend, автоматический выбор параллелизма или последовательность нескольких rescaling.