Перейти к содержанию
Learning Platform
Глоссарий Troubleshooting
Урок 21.04 · 25 мин
Продвинутый
Active-activeSplit-brainCross-regionDedupMulti-region exactly-once

Active-active architecture — RTO ≈ 0 без compromise

Active-passive DR из предыдущих уроков даёт RTO 15-30 минут. Для многих use cases это приемлемо. Но если у вас fraud detection, который должен работать без секунды простоя, или real-time bidding на ad-exchange — даже 5 минут downtime означают миллионы долларов потерь. Здесь приходит active-active архитектура: оба региона работают одновременно, никакого failover не требуется в момент disaster.

Active-Active и Active-Passive топологии Kafka

В этом уроке разберём, как устроена active-active для Flink, какие сложности она создаёт (split-brain, дубликаты, координация state) и какие production-системы её реально используют.


Базовая идея active-active

Два Flink-кластера в разных регионах обрабатывают один и тот же входной поток. Результаты обоих регионов попадают к downstream consumer’у, который выбирает (или merge’ит) данные. При disaster в одном регионе — второй просто продолжает работать без перерыва.

Active-active basic topology
Event sourceEvents генерируются и попадают в multi-region Kafka или реплицируются bi-directionally через MirrorMaker 2 с loop prevention.
Kafka (A)Kafka в region A. Все события доступны здесь — либо через producer write, либо через MM2 replication из region B.
MM2 bi-directional
Kafka (B)Kafka в region B. Те же события через MM2 replication.
Flink (A)Flink job в region A. Обрабатывает все события из Kafka (A). Свой независимый state, savepoint, checkpoints.
no direct coordination
Flink (B)Flink job в region B. Точная копия job в A: тот же graph, тот же parallelism. Обрабатывает те же события, потому что Kafka реплицируется bi-directionally.
Output AOutput Kafka в обоих регионах. Каждый Flink пишет результаты в свой regional output. Downstream consumer должен читать из обоих и дедуплицировать.
merge / dedup
Consumer (dedup)Downstream consumer. Видит результаты из обоих регионов, дедуплицирует по event-id или idempotent ключу. Может быть DB с unique constraint, ClickHouse с ReplacingMergeTree, или другой Flink job.
Output BOutput Kafka в region B.

Преимущества:

  • RTO ≈ 0: один регион упал — второй уже обрабатывает поток. Никакого manual failover.
  • Independent failure domains: ошибка в одном Flink job не влияет на другой. Если в primary дедлок в RocksDB, DR продолжает работать.
  • Geographic load distribution: downstream consumer’ы могут читать из ближайшего региона для уменьшения latency.

Недостатки (огромные):

  • 2x cost: два production-grade кластера вместо одного + standby.
  • Дубликаты в output: оба региона производят результаты — consumer должен дедуплицировать.
  • Невозможен exactly-once cross-region: разные state machines в двух Flink, нельзя 2PC между ними.
  • Сложность операций: deploy, migration, rollback — всё нужно делать на двух кластерах синхронно.

Split-brain: главный риск

Split-brain — это сценарий, когда два «активных» компонента считают себя единственным правильным владельцем данных и принимают конфликтующие решения. В active-active Flink split-brain выглядит так: оба региона имеют свой view состояния, и при network partition между ними они расходятся.

Пример: stateful aggregation sum(amount) over user. Allowed transaction limit = 1000.Пользовательделаеттритранзакции:1000. Пользователь делает три транзакции: 400 в region A (видит state в A), 400вregionB(видитstateвB),400 в region B (видит state в B), 400 в region A (всё ещё видит state только в A).

Region A state:  {user1: 800}
Region B state:  {user1: 400}
Реальная сумма:  $1200  (превышает limit, но обоим Flink выглядит OK)

Оба Flink-кластера независимо одобряют все три транзакции, потому что в своём local state они не видят полного. Финансовая катастрофа.

Split-brain scenario
t=0: $400 -> At=0: user сделал транзакцию $400 в region A. Flink (A) обработал, state стал 400.
replication to BMM2 копирует событие в B. Flink (B) обрабатывает с лагом, state в B становится 400. Пока всё OK.
t=1: network partitiont=1: network partition между A и B. MM2 не может реплицировать. Оба региона работают независимо.
mirroring brokenBi-directional MM2 broken. Каждый регион видит только свои события.
A: state 400 -> 800t=2: user делает $400 в A. Flink (A) видит state 400, добавляет — становится 800. OK.
conflict
B: state 400 -> 800t=2: user делает $400 в B. Flink (B) видит state 400 (последнее реплицированное), добавляет — становится 800. Тоже OK с её точки зрения.
t=3: heal — divergent statet=3: partition healed, MM2 возобновляет. Но оба state уже разошлись (по 800 в каждом, реальная сумма 1200). Невозможно автоматически смержить — это semantic conflict.
manual resolutionManual resolution required. Бизнес-логика должна решить, как обработать divergence. Это часы работы on-call team.

Решения split-brain (все плохие):

1. Sharding by region. Каждый region отвечает только за свой subset users. Region A: users with even id. Region B: users with odd id. Тогда split-brain невозможен — каждый region владеет своими данными. Но это уже не active-active в полном смысле — это два независимых active кластера, и при падении region A users with even id просто недоступны, RTO ≠ 0.

2. CRDT для state. Использовать Conflict-free Replicated Data Types для аггрегаций — например, G-Counter (grow-only counter). Но это работает только для commutative aggregations (sum, max, min, set union). Для других (financial limits, top-K) CRDT не существует.

3. Conflict resolution в downstream. Принять, что Flink output может содержать дубликаты или конфликты, и разрулить их в downstream consumer (например, в DB с custom conflict resolution). Это переносит проблему вниз по pipeline.

4. Просто at-least-once + dedup. Самый распространённый production-подход: принять, что cross-region exactly-once невозможен, и положиться на idempotent consumers.

WARNING

Полноценный active-active с exactly-once cross-region — это distributed consensus задача, эквивалентная Paxos/Raft через network с latency 100ms+ и периодическими partition. Это очень дорого в performance — обычно неприемлемо для high-throughput streaming.


Паттерн 1: independent computation + downstream dedup

Самый распространённый production-паттерн. Оба региона независимо обрабатывают весь поток, дубликаты дедуплицируются в downstream:

// Один и тот же job в обоих регионах:
DataStream<Transaction> source = env.fromSource(
    kafkaSource,
    WatermarkStrategy.<Transaction>forBoundedOutOfOrderness(Duration.ofSeconds(5)),
    "transactions"
);

DataStream<EnrichedTransaction> processed = source
    .keyBy(t -> t.userId)
    .map(new EnrichmentFunction())  // stateful enrichment
    .filter(t -> t.amount > 10);  // business filter

processed.sinkTo(
    KafkaSink.<EnrichedTransaction>builder()
        .setBootstrapServers("regional-kafka:9092")
        .setRecordSerializer(KafkaRecordSerializationSchema.builder()
            .setTopic("transactions-enriched-region-a")  // регион-специфичный топик
            .setValueSerializationSchema(new TransactionSerializer())
            .build())
        .build()
);

Downstream consumer (например, ClickHouse) использует ReplacingMergeTree:

CREATE TABLE transactions_enriched (
    transaction_id String,
    user_id UInt64,
    amount Decimal(18, 2),
    region String,
    processed_at DateTime64
)
ENGINE = ReplacingMergeTree(processed_at)
ORDER BY transaction_id;

При вставке дубликатов с одним transaction_id ClickHouse оставит только запись с максимальным processed_at. Дубликаты «дедуплицируются» в фоне.

Альтернатива — Postgres с ON CONFLICT:

INSERT INTO transactions_enriched (transaction_id, user_id, amount)
VALUES (?, ?, ?)
ON CONFLICT (transaction_id) DO NOTHING;

Главное требование: каждое событие должно иметь stable unique identifier, который не меняется при reprocessing.


Паттерн 2: region-sharded processing

Каждый регион обрабатывает свой shard данных. Cross-region replication только для DR backup:

// Region A:
Properties kafkaProps = new Properties();
kafkaProps.setProperty("bootstrap.servers", "kafka-a:9092");
kafkaProps.setProperty("group.id", "flink-region-a");

DataStream<Event> source = env.fromSource(
    KafkaSource.<Event>builder()
        .setProperties(kafkaProps)
        .setTopics("events")  // полный топик, но мы фильтруем по shard
        .setStartingOffsets(OffsetsInitializer.latest())
        .build(),
    WatermarkStrategy.noWatermarks(),
    "events-source"
);

DataStream<Event> shardA = source
    .filter(e -> e.userId % 2 == 0);  // только even userIds в region A

shardA
    .keyBy(e -> e.userId)
    .process(new StatefulProcessor())
    .sinkTo(...);
// Region B (deployed independently):
DataStream<Event> shardB = source
    .filter(e -> e.userId % 2 == 1);  // только odd userIds в region B
// ... rest is the same

Преимущества:

  • No split-brain — каждый region владеет своим shard, конфликтов state нет.
  • Linear scalability — 2x throughput с 2 регионами.
  • Geographic locality — users from Europe processed в EU region.

Недостатки:

  • Half-region failure означает half-data downtime — если region A упал, all even userIds не обрабатываются, пока не сделаешь failover (= becomes active-passive).
  • Routing complexity — нужен upstream routing, чтобы события попадали в правильный shard.

Production examples

Netflix: использует active-passive для большинства Flink jobs (savepoint-based, ~5 min RTO). Для критичных fraud detection и personalization — region-sharded active-active по user geography. Cross-region failover тестируется monthly.

Uber: uReplicator (custom MirrorMaker derivative) для Kafka mirroring. Flink jobs в active-active по regions, дедупликация в downstream через storage layer (Pinot, MySQL). RPO < 1 минута.

Pinterest: Flink на Kubernetes с active-passive в двух регионах. Savepoint каждые 5 минут, replication через AWS S3 CRR with RTC (15-minute RTO guarantee).

Stripe: financial transactions — strict exactly-once requirement. Решение: active-passive, но с очень частым checkpoint (10s) и в основном sync replication через Aurora Global Database (not Flink-level).

Pattern: полное active-active с exactly-once в production почти не встречается. Все известные мне реализации либо region-sharded (не настоящий active-active), либо принимают at-least-once с дедупликацией.

TIP

Если ваш use case реально требует RTO=0 и exactly-once cross-region — критически пересмотрите requirements с бизнесом. Чаще всего «нужен RTO=0» означает «нужен RTO < 30 секунд», что достижимо в active-passive с warm standby + быстрая automation. Реальный RTO=0 требует архитектуры, которая стоит в 5-10x дороже.


Hybrid: warm active с быстрым take-over

Компромисс между active-passive и active-active: standby cluster в DR region реально обрабатывает события (читает Kafka, держит state up-to-date), но не пишет в output. При failover — просто разрешает writes и take over downstream.

boolean isPrimary = Config.getBoolean("flink.cluster.is-primary");

DataStream<EnrichedTransaction> processed = source
    .keyBy(t -> t.userId)
    .map(new EnrichmentFunction());

if (isPrimary) {
    processed.sinkTo(kafkaSink);  // только primary пишет
} else {
    processed.addSink(new DiscardingSink<>());  // standby обрабатывает но discard
}

Преимущества:

  • State в DR-кластере всегда warm — не нужен RocksDB warmup.
  • Lag реплицируется на consume-стороне (Kafka offsets), не на flink-savepoint-стороне.
  • Failover = переключение config is-primary=true в DR + DNS switch для consumers.

RTO достижим 30 секунд (детектирование 10с + DNS update 10с + config reload 10с). Стоимость — почти такая же как полный active-active.

Это паттерн, который реально используется в production когда нужен низкий RTO без сложности true active-active.


Итоги

Active-active даёт RTO ≈ 0, но за счёт значительной сложности: split-brain risks, отсутствие cross-region exactly-once, дублирование стоимости. Чистый active-active в production почти не встречается — большинство известных реализаций это либо region-sharded (квази active-active), либо warm-active с быстрым take-over.

Для большинства бизнес-задач active-passive с warm standby и хорошей automation даёт достаточный RTO (5-15 минут) при разумной сложности.

В следующем модуле (capstone) мы построим production-grade Flink platform, которая объединяет всё: multi-tenancy, autoscaling, savepoint automation, observability и runbook’и для типовых инцидентов.

Проверка знанийKnowledge check
Команда хочет построить active-active Flink для fraud detection с требованием exactly-once cross-region. Архитект предлагает использовать distributed Raft consensus между двумя Flink clusters для coordination state. Какие проблемы возникнут и какова правильная рекомендация?
ОтветAnswer
Распределённый consensus (Raft, Paxos) требует quorum write, что означает synchronous round-trip на каждое state mutation. Cross-region latency 100ms+ делает это catastrophically slow: throughput high-frequency streaming jobs (100K+ events/s) падает на порядки. Это фундаментальное ограничение physics — нельзя обойти. В production это всегда решается компромиссами: region-sharding (нет cross-region coordination нужна), CRDTs (для commutative operations), at-least-once + downstream dedup (большинство production), или warm active-passive (не настоящий active-active, но RTO менее 1 минуты достижим). Никакая известная production система не делает true active-active с strict exactly-once cross-region для high-throughput streaming.

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

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

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

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