DR fundamentals для streaming systems
Disaster recovery (DR) в batch-системах относительно прост: данные лежат в S3 или HDFS, реплицируются через cross-region replication, и при падении региона ETL-pipeline просто перезапускается в другом регионе с теми же входными файлами. Streaming-системы устроены принципиально иначе. Поток данных непрерывен, состояние операторов хранится в RocksDB на TaskManager-нодах, exactly-once гарантии зависят от 2PC-протокола между Kafka и Flink — и всё это надо вытащить из упавшего региона в другой за минуты, а не часы.
В этом уроке разберём терминологию (RPO, RTO), почему стандартные DR-подходы не работают для streaming, и какие принципиальные архитектуры существуют — active-passive и active-active.
RPO/RTO и failover для Kafka кластеровRPO и RTO: язык, на котором говорят с бизнесом
Два центральных термина в DR-планировании:
RPO (Recovery Point Objective) — максимально допустимая потеря данных в результате disaster. Если RPO = 5 минут, значит при падении региона мы можем потерять до 5 минут продакшн-событий. Если RPO = 0, значит ни одного события потерять нельзя.
RTO (Recovery Time Objective) — максимально допустимое время восстановления сервиса. Если RTO = 10 минут, значит от момента падения до возобновления полноценной работы должно пройти не более 10 минут.
В streaming эти метрики специфичны. RPO измеряется не в файлах backup, а в отставании savepoint от текущего event time. Если savepoint мы триггерим раз в 5 минут и копируем в DR-регион асинхронно с лагом 1 минута, то реальный RPO составит примерно 6 минут — все события, обработанные между моментом savepoint и моментом disaster, потенциально потеряются (если их нет в Kafka с длинной retention) или будут переобработаны (если Kafka реплицируется в DR-регион и Flink может «перечитать» от смещения savepoint).
RTO в streaming включает: время на детектирование disaster (1-2 минуты обычно), запуск standby Flink-кластера в DR-регионе (если он не был запущен — 3-5 минут на K8s), restore из savepoint (от секунд до десятков минут в зависимости от размера state), warmup state (RocksDB block cache, network shuffle channels) — иногда минуты до устаканивания throughput.
Никогда не путайте RPO/RTO целевые значения с фактическими. Целевые формулируются совместно с бизнесом. Фактические измеряются регулярными DR-учениями (game days). Без учений вы не знаете реального RTO — обычно фактическое в 2-5 раз больше теоретического.
Почему batch-DR не работает для streaming
В batch ETL pipeline DR обычно выглядит так: source data лежит в S3 cross-region replicated bucket, batch job (Spark, Trino) перезапускается в DR-регионе и читает данные с того же бакета. Если job упал на полпути — Spark с idempotent writes начнёт от начала batch’а, никакой проблемы.
Streaming так работать не может по трём причинам:
Stateful operators хранят состояние в локальном RocksDB. При падении региона RocksDB на TaskManager-нодах потеряны. Восстановить можно только из checkpoint/savepoint, который должен лежать в durable storage, доступном из DR-региона. Это значит — обязательная cross-region репликация checkpoint storage (S3 с cross-region replication, или прямая запись в multi-region object store).
Source offsets — тоже state. Kafka source хранит offsets в Flink state. Если перезапустить job в DR-регионе без synchronized state, то Kafka offsets будут несинхронны с реальностью: либо мы перечитаем уже обработанные сообщения (дубликаты), либо пропустим сообщения, которые не успели зафиксироваться в DR Kafka.
Exactly-once требует transactional sink. При failover сначала надо abort незавершённые транзакции в downstream sink (Kafka producer, JDBC), затем restart from savepoint. Без этого consumers увидят либо дубликаты, либо потерянные транзакции.
Active-passive vs active-active: фундаментальный выбор
Есть две принципиально разные архитектуры DR для Flink:
Active-passive (warm standby): job постоянно работает только в одном (primary) регионе. В secondary регионе либо вообще ничего не запущено (cold standby), либо запущен пустой кластер без job (warm standby). При disaster в primary — оператор (или automation) запускает job в secondary, restore из последнего savepoint, переключает downstream consumers. RTO измеряется в минутах. RPO зависит от частоты savepoint replication.
Active-active (dual-region): job одновременно работает в обоих регионах, обрабатывая один и тот же входной поток (или независимые партиции потока). Downstream consumers видят результаты из обоих регионов через какой-то merge layer. При disaster в одном регионе — второй продолжает работать без перерыва. RTO близок к нулю, но архитектурно гораздо сложнее.
Active-passive проще и дешевле — DR-регион в основном простаивает. Active-active требует двух production-grade кластеров, дополнительной логики merge/dedup на consumer-стороне, и обычно жертвует exactly-once в пользу at-least-once с дедупликацией.
Большинство production Flink-инсталляций (Netflix, Uber, Pinterest) используют active-passive — это компромисс между стоимостью и сложностью. Active-active встречается только там, где RTO=0 это бизнес-требование (например, fraud detection в реальном времени).
Что измерять и зачем мониторить DR-метрики
Без мониторинга DR-готовности вы узнаете о проблемах только во время disaster. Минимальный набор метрик:
# Лаг репликации savepoint в DR S3
aws s3 ls s3://flink-savepoints-dr/myJob/ --recursive | \
tail -1 | awk '{print $1, $2}'
# Сравнить с локальным last savepoint timestamp:
aws s3 ls s3://flink-savepoints-primary/myJob/ --recursive | \
tail -1 | awk '{print $1, $2}'
# Лаг MirrorMaker 2 для Kafka cross-region
kafka-consumer-groups.sh \
--bootstrap-server dr-broker:9092 \
--describe --group mm2-myCluster | \
awk '{print $1, $2, $3, $5}' # topic, partition, current, lag
# Время с последнего успешного DR-учения
date -d "$(cat /var/lib/dr/last-game-day)" +%s
echo "$(($(date +%s) - $(date -d "$(cat /var/lib/dr/last-game-day)" +%s)))/86400" | bc
Метрики, которые должны быть в Grafana с алертами:
flink_dr_savepoint_age_seconds— возраст последнего savepoint в DR-регионе. Алерт если > 2x от RPO target.kafka_mirror_lag_ms— лаг MirrorMaker 2. Алерт если > 30s.flink_dr_cluster_health— health-check standby кластера в DR-регионе. Алерт если pods не Ready.dr_drill_days_since_last— дни с последнего DR-учения. Алерт если > 90 дней.
Кроме метрик, нужны runbook’и: документированные процедуры failover для каждого типа disaster. Без runbook’а во время реального инцидента команда теряет 30+ минут на «вспомнить, что делать». В капстон-модуле мы построим автоматизированные runbook’и для типовых сценариев.
Лучший индикатор зрелости DR-стратегии — периодические game days. Раз в квартал намеренно «убивайте» primary регион (в staging environment, естественно) и измеряйте фактическое RTO. Это единственный способ обнаружить regressions: новый Kafka топик, который не реплицируется; забытый Vault secret; неактуальный runbook.
Итоги
DR в streaming — это про synchronization трёх независимых компонентов: checkpoint storage, Kafka topics и downstream sinks. RPO/RTO формулируются совместно с бизнесом и измеряются регулярно. Активно используются две архитектуры: active-passive (проще, дешевле, RTO в минутах) и active-active (RTO ≈ 0, но сложнее на порядок).
В следующих уроках разберём savepoint-based DR (стандартный подход), Kafka mirroring (как реплицировать топики между регионами без offset rot), и active-active паттерн с реальными production examples.