Перейти к содержанию
Learning Platform
Глоссарий Troubleshooting
Урок 21.01 · 25 мин
Продвинутый
Disaster RecoveryRPORTOMulti-regionActive-passiveActive-active

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 минут.

RPO vs RTO timeline
Last backup pointПоследний успешный backup или snapshot до disaster. Всё, что произошло между этим моментом и disaster, потенциально потеряно.
RPO window
Disaster momentМомент disaster. Регион упал, primary cluster недоступен. Данные после last backup пока не реплицированы.
RTO window
Service restoredСервис восстановлен в secondary регионе. Данные обрабатываются, downstream consumers получают результаты.

В 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.

WARNING

Никогда не путайте 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 увидят либо дубликаты, либо потерянные транзакции.

Что должно быть синхронизировано между регионами
Checkpoint storageCheckpoints и savepoints должны лежать в multi-region object store или asynchronously replicated в DR-регион. RPO зависит от частоты репликации.
Source KafkaKafka топики реплицируются через MirrorMaker 2 или Confluent Cluster Linking. Offset translation — нетривиальная задача, см. урок 03.
Sink storageDownstream Kafka (для sink) или JDBC реплицируется отдельно. После failover надо abort незавершённые транзакции.
Job configurationКонфигурация job (job graph, parallelism, savepoint location). Хранится в Git и доступна в DR-регионе для запуска standby кластера.
Schema RegistrySchema Registry — критичен для Avro/Protobuf. Должен быть реплицирован, иначе DR-регион не сможет десериализовать сообщения.
SecretsSecrets и credentials (Kafka SASL, S3 keys, DB passwords). Vault или AWS Secrets Manager должны быть доступны cross-region.

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 vs active-active
Active-passive: primaryPrimary регион: единственный источник правды. Все события обрабатываются здесь. Savepoint периодически копируется в S3 DR региона.
savepoint copy
Standby clusterSecondary регион: либо пустой K8s namespace (cold), либо запущенный Flink cluster без job (warm). При disaster — запуск job из последнего savepoint.
Active-active: region AActive-active: оба региона активны. Каждый обрабатывает либо весь поток (с дедупликацией на consumer-стороне), либо непересекающиеся партиции (с merge на consumer-стороне).
coordination?
Region BRegion B работает параллельно. Сложность: как координировать exactly-once между регионами? Обычно — нельзя, переходят на at-least-once с дедупликацией.

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’и для типовых сценариев.

TIP

Лучший индикатор зрелости 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.

Проверка знанийKnowledge check
Job обрабатывает 50K events/s, savepoint триггерится каждые 5 минут, копирование savepoint в DR S3 занимает 2 минуты. Какой минимально достижимый RPO при этой конфигурации, и что произойдёт с событиями, обработанными после savepoint, но до disaster?
ОтветAnswer
RPO = время между последним savepoint и disaster + лаг репликации savepoint в DR. Если savepoint раз в 5 минут и копирование 2 минуты, то худший случай: disaster наступил через 5 минут после savepoint, копирование ещё идёт — теряем 7 минут state. Но это не значит, что события навсегда потеряны: savepoint содержит Kafka offsets, и если Kafka реплицируется в DR через MirrorMaker, Flink в DR-регионе перечитает все сообщения от offset savepoint. Это превращает state loss в event reprocessing — что обычно приемлемо для idempotent pipeline.

Закончили урок?

Отметьте его как пройденный, чтобы отслеживать свой прогресс

Войдите чтобы оценить урок

Прогресс модуля
0 из 4