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 в одном регионе — второй просто продолжает работать без перерыва.
Преимущества:
- 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 = 400 в region A (видит state в A), 400 в region A (всё ещё видит state только в A).
Region A state: {user1: 800}
Region B state: {user1: 400}
Реальная сумма: $1200 (превышает limit, но обоим Flink выглядит OK)
Оба Flink-кластера независимо одобряют все три транзакции, потому что в своём local state они не видят полного. Финансовая катастрофа.
Решения 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.
Полноценный 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 с дедупликацией.
Если ваш 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’и для типовых инцидентов.