Gossip Protocol — concept page. Epidemic dissemination with push/pull/push-pull, periodic random peer selection, O(log N) convergence, SWIM failure detection (suspect -> confirmed via ping/ping-req), delta-based bandwidth control. Used in Cassandra, Consul, Serf, Riak, Akka.
Когда в кластере сотни или тысячи нод, централизованный membership service (ZooKeeper, etcd) превращается в узкое горлышко: каждая нода держит watch, любое изменение состояния порождает шторм нотификаций, а Raft-кворум начинает тонуть в трафике. Нужен способ распространять знание (кто жив, какая версия схемы, какие токены у кого) без единой точки координации, с предсказуемой пропускной способностью и устойчивостью к partition.
Gossip protocol решает это эпидемическим распространением: каждая нода периодически выбирает несколько случайных пиров и обменивается состоянием. Через O(log N) раундов новость доходит до всех — без брокера, без координатора, без watch-стормов. То, что в природе делает вирус, в распределённой системе делает rumor.
Используется для трёх задач: membership (кто в кластере), failure detection (кто умер), anti-entropy (догнать пропущенное). Это инфраструктурный слой под Cassandra, Consul, Serf, Riak, Akka Cluster, Redis Cluster.
«Каждую секунду каждая нода выбирает 2 случайных пира, говорит "вот мой digest, дай свой", обменивается только дельтами. Любая новость через
log2(N)раундов известна всем. Никто никем не управляет — все равноправны.»
Три варианта обмена:
Ключевая интуиция: fanout (сколько пиров на раунд) важнее, чем интервал. Fanout=2 даёт log2(N) раундов, fanout=3 — log3(N). Слишком высокий (>5) — гossip-шторм при каждом изменении; слишком низкий (1) — конвергенция деградирует с логарифма до линейной.
8 нод в полном меш-кластере (n1..n8), все симметричны — нет лидера, нет координатора. Edges нарисованы как репрезентативные TCP-каналы; в реальном кластере связи логические (UDP-датаграммы или эфемерные TCP), а не пред-установленные соединения. Каждая нода может гossipить с любой другой; на каждом раунде выбирается случайный subset.
Три сценария проигрывают три ключевых режима работы protocol:
log2(8)=3 раунда.node-1 узнаёт новый факт (например, новая нода присоединилась, или обновилась версия схемы). В раунде 1 он pickает 2 случайных пира n2, n3 — знают трое. В раунде 2 каждый из троих pickает по 2 — знают шестеро. В раунде 3 — все восемь. Это и есть O(log2 N): log2(8)=3, log2(1024)≈10, log2(1M)≈20. Для кластера в 10 000 нод хватит ~14 раундов, что при интервале 1 секунда даёт 14-секундную convergence latency.
Важная деталь: per-round bandwidth константна — каждая нода говорит ровно с fanout пирами независимо от размера кластера. Total bandwidth растёт как O(N) (сумма по нодам), но per-node — O(1). Это и есть главная причина, почему gossip скейлится туда, куда watch-based системы не доходят.
Наивная схема "не ответил на ping за T секунд = мёртв" ломается на масштабе по двум причинам: (1) head-of-line blocking — нода с GC-паузой ошибочно объявляется мёртвой, триггерится rebalance, остальные тоже начинают тонуть; (2) network partition — minority объявляет majority мёртвой.
SWIM (Das/Gupta/Motivala, 2002) добавляет два слоя гистерезиса:
n1 → n4 не дошёл, n1 просит k witness-нод (n2, n3, n5) проpingить n4 от их имени. Это покрывает случай "плохой link именно у n1 к n4", но n4 жив.n4 помечается как SUSPECT (не DEAD!), эта новость гossipится. У n4 есть T_suspect (~5s) на то, чтобы refute через incremented incarnation number. Если refutation пришло — suspect снимается. Если нет — promote в CONFIRMED DEAD, и об этом гossipится.Cassandra и Akka используют Phi-accrual detector (Hayashibara, 2004) вместо фиксированных таймаутов: continuous suspicion score из распределения межприбытийных интервалов. Адаптируется к GC и сетевому джиттеру автоматически. Hashicorp memberlist добавляет Lifeguard extension — нода учитывает собственное здоровье, чтобы не штамповать false positives во время своей же GC-паузы.
Push-pull в Cassandra Gossiper — это три-фазный протокол:
n1 шлёт n2 digest (version vector по всем известным нодам, ~200 байт). Без payload.n2 сравнивает: {n1@v15 > my v14} → нужен pull; {n5@v3 < my v5} → нужен push. Шлёт обратно: запрос на updates по n1 + полный state n5.n1 шлёт n2 полный state n1@v15.Итого ~1.1 KB на полный обмен между двумя пирами. Если бы каждый раунд гонять полный state — было бы ~10 KB. На кластере в 1000 нод разница: ~1 KB/sec/node vs ~10 KB/sec/node gossip overhead. Это десятикратная экономия bandwidth.
Подводные камни:
ADR-001 — Gossip vs centralized membership. Выбираем gossip, когда: кластер > ~500 нод, нет приемлемой single point of coordination, нужна partition-tolerance (gossip деградирует gracefully — partitioned subset продолжает гossipить внутри себя), eventually consistent membership (~секунды задержки) приемлем. Выбираем централизованный (ZK/etcd), когда: кластер < ~100 нод, нужна strong consistency на membership-решениях (leader election, config rollout где stale view опасен), уже есть ZK для других целей. Hybrid распространён: Consul использует Serf/SWIM gossip для liveness на тысячах agent-ов И Raft-кворум для маленького server-set, который хранит source of truth.
ADR-002 — Push-pull, не push-only/pull-only. Push-only стопорится в конце (большинство уже знает, push тратится впустую). Pull-only стопорится в начале (мало кто знает, pulls возвращают пустоту). Push-pull (Cassandra, memberlist/Serf/Consul) — один round-trip обменивает информацию в обе стороны. Fanout = 2-3 в дефолтах: Cassandra 1s/3, Serf 200ms/3. Fanout > 5 → gossip-шторм при каждом изменении; fanout = 1 → конвергенция деградирует с log(N) до N.
ADR-003 — SWIM трёхстадийный, не наивный таймаут. Direct ping → indirect ping-req via witnesses → SUSPECT (gossipится с шансом refute) → CONFIRMED. Каждая стадия добавляет гистерезиса против ложных срабатываний во время GC-пауз и transient partition. Phi-accrual (Cassandra/Akka) дополнительно адаптируется к сетевым условиям без ручного тюнинга таймаутов.
O(log N) convergence.