Параллельные и распределённые алгоритмы, консенсус
Все предыдущие статьи трека предполагали одну машину, одно ядро и надёжную память. Реальность другая: у сервера 128 ядер, у GPU — десятки тысяч потоков, а данные лежат в кластере из сотен узлов, которые падают, тормозят и теряют пакеты. Здесь начинаются два разных мира:
- Параллельные алгоритмы — общая память, потоки в одном процессе, отказов нет (падение процесса убивает всех). Вопрос: как разложить работу, чтобы ускорение было близко к p?
- Распределённые алгоритмы — сеть, независимые узлы, частичные отказы, нет общих часов. Вопрос: как договориться, когда часть участников молчит и непонятно — они умерли или тормозят?
- Мост между мирами — консенсус: договорённость о едином порядке при частичных отказах.
Часть I. Параллельные алгоритмы
Модель работа/глубина: правильный способ считать сложность
Обычная асимптотика O(n log n) ничего не говорит о параллелизуемости. Нужны две величины:
работа W(n) — суммарное число операций (обычная последовательная сложность) — и глубина
(span, critical path) D(n) — длина самой длинной цепочки зависимостей, то есть время
на бесконечном числе процессоров. Отношение W/D называют параллелизмом: это верхняя граница разумного числа ядер.
Если W = O(n), а D = O(log n), параллелизм n / log n — можно занять тысячи потоков.
Если D = Θ(n) (обычный цикл-аккумулятор), параллелизма нет вообще.
Теорема Брента: на p процессорах жадный планировщик даёт T_p ≤ W/p + D. Отсюда правило:
пока W/p ≫ D, ускорение почти линейно; как только D сравнима с W/p, ядра бессмысленны.
Модель PRAM из классических учебников — идеализация с мгновенной общей памятью, полезная
для нижних границ. Для инженерии удобнее work-depth: она не зависит от p и композируется —
при последовательной композиции складываются и работы, и глубины; при параллельной работы
складываются, а глубина берётся максимумом.
Закон Амдала и почему он не приговор
Если доля s программы принципиально последовательна, ускорение ограничено
S(p) = 1 / (s + (1−s)/p), что при p → ∞ даёт потолок 1/s. При s = 5% — 20×,
сколько ядер ни дай. Это закон Амдала, и он убийственен для фиксированной задачи.
Но закон Густафсона смотрит иначе: получив больше ядер, мы решаем большую задачу
за то же время, а последовательная часть (инициализация, чтение конфига) не растёт:
S_scaled(p) = s + p(1−s). Вывод для инженера: не спрашивайте «сколько ядер дадут ускорение»,
спрашивайте «растёт ли последовательная часть вместе с n». Если нет — масштабирование будет.
между элементами?"} B -- "нет" --> C["Embarrassingly parallel:
map / фильтрация / рендер тайлов
D = O(1)"] B -- "есть, но ассоциативные" --> D["Reduce / scan / segmented scan
D = O(log n)"] B -- "рекурсивное разбиение" --> E["Divide & conquer:
merge sort, quickhull, FFT
D = O(log^2 n)"] B -- "цепочка состояний" --> F{"Можно переписать как
ассоциативный оператор?"} F -- "да" --> D F -- "нет" --> G["Последовательное ядро.
Параллелить снаружи:
батчи, шардирование, pipeline"] C --> H["Упёрлись в пропускную
способность памяти?"] D --> H E --> H G --> H H --> I["Профилировать: перф-счётчики,
false sharing, NUMA"]
Три кита: map, reduce, scan
Map тривиален: W = O(n), D = O(1). Reduce — свёртка ассоциативным оператором
по дереву: W = O(n), D = O(log n). Требование ассоциативности критично: сложение float
не ассоциативно, поэтому параллельная сумма даёт другой (обычно более точный!) результат,
чем последовательная — типичный источник «нестабильных» тестов в численном коде.
Scan (префиксные суммы) — самый недооценённый примитив. Кажется сугубо последовательным
(out[i] = out[i-1] + a[i]), но параллелизуется на O(log n) глубины. Через scan выражаются
уплотнение массива (stream compaction), radix sort, распределение работы между потоками,
разбор скобок, line-of-sight.
| Алгоритм | Работа W |
Глубина D |
Когда применять |
|---|---|---|---|
| Hillis–Steele (наивный) | O(n log n) |
O(log n) |
внутри warp’а GPU (n ≤ 32), где лишняя работа бесплатна |
| Blelloch (work-efficient) | O(n) |
O(log n) |
блоки и глобальный уровень, большие n |
| Decoupled look-back | O(n) |
O(log n), один проход по памяти |
state-of-the-art на GPU (CUB) |
Hillis–Steele тривиален (out[i] += out[i - step] для всех i параллельно, step удваивается),
но делает лишнюю работу. Blelloch — это ровно алгоритм с картинки:
def scan_blelloch(a):
"""Exclusive scan за O(n) работы и O(log n) глубины.
Внутренние циклы for — параллельные раунды: итерации независимы по i."""
n = len(a)
assert n & (n - 1) == 0, "для простоты — степень двойки"
t = list(a)
d = 1 # --- up-sweep: суммы поддеревьев в правые узлы ---
while d < n:
for i in range(2 * d - 1, n, 2 * d):
t[i] += t[i - d]
d *= 2
total, t[n - 1] = t[n - 1], 0 # корень заменяем нейтральным элементом
d = n // 2 # --- down-sweep: раздаём префиксы вниз ---
while d >= 1:
for i in range(2 * d - 1, n, 2 * d):
left = t[i - d] # сумма левого поддерева (осталась с подъёма)
t[i - d] = t[i] # левому ребёнку — значение родителя
t[i] = t[i] + left # правому — родитель + сумма левого
d //= 2
return t, total
assert scan_blelloch([3, 1, 7, 0, 4, 1, 6, 3])[0] == [0, 3, 4, 11, 11, 15, 16, 22]
Каноническая работа — Guy Blelloch, «Prefix Sums and Their Applications» (1993); инженерный разбор — NVIDIA CUB DeviceScan.
Параллельная сортировка
Наивная параллелизация merge sort (https://courses.digitable.life/post/algorithms/05-recursion-and-divide-conquer/)
даёт W = O(n log n), но глубину D = Θ(n) из-за последовательного слияния — параллелизм
всего O(log n). Лечится параллельным слиянием: чтобы слить A и B, берём медиану A,
бинарным поиском находим её позицию в B, рекурсивно сливаем половины ⇒ D = O(log² n).
На практике для многоядерных CPU и кластеров чаще используют sample sort:
- Каждый из
pпотоков берётsслучайных элементов; сортируемp·sсэмплов, берёмp−1splitter’ов. - Каждый поток бинпоиском раскладывает свой чанк по
pкорзинам; гистограмма + exclusive scan дают точные смещения записи — без единого мьютекса. - Корзина
iсортируется независимо (pdqsort/radix), конкатенация и есть ответ.
Sample sort делает один проход обмена данными вместо log p проходов у merge-based подходов —
поэтому он выигрывает и на CPU, и в Spark/MapReduce, где обмен по сети дорог.
От модели к коду: fork-join и границы зернистости
GRAIN = 1 << 15 # порог: ниже — последовательно, иначе накладные расходы съедят выигрыш
def parallel_sum(a, lo, hi, pool, depth=0):
"""Классический fork-join. W = O(n), D = O(log n) при неограниченном пуле."""
if hi - lo <= GRAIN or depth > 8: # ограничиваем глубину дробления
return sum(a[lo:hi])
mid = (lo + hi) // 2
fut = pool.submit(parallel_sum, a, lo, mid, pool, depth + 1)
right = parallel_sum(a, mid, hi, pool, depth + 1) # правую ветку считаем сами
return fut.result() + right
Два правила, которые нарушают чаще всего. Зернистость: задача дешевле ~10–50 мкс не окупает
постановки в очередь — дробите до порога, а не до элемента. Одну ветку выполняйте сами:
отправлять в пул обе половины значит удваивать число задач и голодать по стеку; именно так
устроены work-stealing планировщики Cilk, Java ForkJoinPool, Rust rayon. В Go тот же паттерн через
errgroup с семафором (https://courses.digitable.life/post/golang/00-overview/), в Elixir — Task.async_stream/3
с max_concurrency (https://courses.digitable.life/post/elixir/00-overview/).
Синхронизация: блокировки, CAS, и где всё ломается
Модель работа/глубина предполагает бесплатную синхронизацию. Реальность:
- Гонка (data race) — два потока обращаются к одной ячейке, хотя бы один пишет, и между ними нет happens-before. В C++/Rust/Go это UB, а не «иногда неправильный ответ».
- False sharing — потоки пишут в разные переменные, попавшие в одну кэш-линию (64 байта). Гонки формально нет, а производительность падает в 5–20 раз из-за пинг-понга линии между ядрами. Лечение: padding до 64 байт либо агрегация в локальной переменной с одной записью в конце (https://courses.digitable.life/post/algorithms/18-practical-optimization/).
- CAS
compare_and_swap(addr, expected, new)универсален: по теореме Херлихи о консенсус-числе он реализует любой lock-free объект, в отличие от atomic read/write (1) и test-and-set (2).
# Псевдокод lock-free стека Трайбера — канонический CAS-цикл
def push(stack, value):
node = Node(value)
while True:
node.next = atomic_load(stack.top) # acquire
if atomic_compare_exchange(stack.top, node.next, node): # release
return # иначе проиграли гонку — повтор
def pop(stack):
while True:
head = atomic_load(stack.top)
if head is None:
return None
nxt = head.next # ОПАСНО: head мог быть освобождён другим потоком
if atomic_compare_exchange(stack.top, head, nxt):
return head.value
Этот pop содержит проблему ABA: поток читает head = A и засыпает; другие снимают A,
снимают B, возвращают A обратно; CAS проходит (указатель тот же!), но nxt указывает
на уже удалённый B. Лечения — теговые указатели, hazard pointers, epoch-based reclamation
(crossbeam в Rust, RCU в ядре Linux).
Главный вывод: lock-free — это не «быстрее», это «без блокировок». Lock-free гарантирует прогресс системы, wait-free — прогресс каждого потока за конечное число шагов. Под низкой конкуренцией обычный мьютекс с адаптивным спином часто быстрее самописного CAS-цикла. Читать: Herlihy & Shavit, The Art of Multiprocessor Programming (2020); модель памяти Go.
Часть II. Распределённые алгоритмы
Здесь меняются аксиомы: нет общей памяти — только сообщения, нет глобальных часов, и главное — частичные отказы: узел может исчезнуть, а система обязана жить дальше.
Модели системы: что вы предполагаете, то и получаете
| Ось | Варианты | Последствие |
|---|---|---|
| Синхронность | синхронная / частично синхронная / асинхронная | в асинхронной невозможно отличить «упал» от «тормозит» |
| Отказы | crash-stop / crash-recovery / omission / византийские | византийские требуют втрое большей избыточности |
| Сеть | надёжная / с потерями / с разбиением | разбиение — не экзотика, а ежедневная реальность |
Практически все продакшн-системы живут в модели частичной синхронности (Dwork, Lynch, Stockmeyer, 1988): система асинхронна произвольно долго, но после неизвестного момента GST (Global Stabilization Time) ведёт себя синхронно — это и даёт живучесть при сохранении безопасности.
Время: почему нельзя доверять часам
Физические часы на узлах расходятся; NTP-коррекция может отмотать время назад. Любая логика «у кого timestamp больше — тот прав» на настенных часах — отложенный баг. Инструмент — логические часы (Lamport, 1978):
- Скалярные часы Лампорта:
Lинкрементируется на каждом событии, кладётся в сообщение, при полученииL = max(L, L_msg) + 1. Свойствоa → b ⇒ L(a) < L(b); обратное неверно — по числам нельзя понять, был ли параллелизм. - Векторные часы: вектор из
nсчётчиков,V(a) < V(b) ⟺ a → b. Несравнимые векторы = конкурентные события = конфликт, который надо разрешать явно. Цена —O(n)метаданных. - Гибридные логические часы (HLC) — комбинация физического и логического времени (CockroachDB, YugabyteDB): близко к настенным значениям с сохранением причинности.
- Spanner TrueTime: атомные часы и GPS дают интервал
[earliest, latest], транзакция ждёт (commit wait) выхода за границу неопределённости — покупка внешней консистентности за железо (Spanner, OSDI 2012).
Задача консенсуса
Набор узлов должен согласовать одно значение. Формально требуются: Agreement (два корректных узла не решают разные значения), Validity (решённое значение было кем-то предложено) и Termination (каждый корректный узел рано или поздно решает).
Первые два — свойства безопасности (safety), третье — живучести (liveness). Это различие — ключевой концепт области: safety нарушать нельзя никогда, liveness можно откладывать.
FLP: почему детерминированного решения не существует
Теорема Фишера — Линча — Патерсона (1985): в полностью асинхронной системе не существует
детерминированного алгоритма консенсуса, устойчивого даже к одному crash-отказу.
Интуиция доказательства: показывается существование «бивалентной» конфигурации, из которой
достижимы оба исхода; затем строится бесконечное расписание, всегда переводящее систему
из одной бивалентной конфигурации в другую — планировщик-противник просто задерживает нужное
сообщение (оригинал).
FLP не означает «консенсус невозможен». Он означает, что нельзя иметь одновременно safety,
liveness и полную асинхронность при детерминизме. Три обхода: частичная синхронность
(Paxos, Raft, PBFT — safety всегда, liveness после GST; выбор 99% продакшна); рандомизация
(Ben-Or 1983, общие монеты, HoneyBadgerBFT — завершение с вероятностью 1, см.
https://courses.digitable.life/post/algorithms/14-randomized-algorithms/); детекторы отказов (Chandra & Toueg 1996: ◇W —
слабейший детектор, достаточный для консенсуса; формализация того, чем на практике является таймаут).
Кворумы: почему работает большинство
Всё держится на элементарном факте: если |A| + |B| > n, множества пересекаются. Взяв кворум
⌊n/2⌋ + 1, мы гарантируем, что любые два кворума имеют общий узел, «видевший» предыдущее
решение. Отсюда n ≥ 2f + 1 для f crash-отказов.
Тонкость (Flexible Paxos, 2016): пересекаться должны кворумы фазы 1 и фазы 2, а не кворумы
одной фазы между собой — при n = 5 можно взять |Q1| = 4, |Q2| = 2 и удешевить горячий путь
записи ценой более дорогих выборов лидера. Кворумы работают и без консенсуса: Dynamo-подобные
системы берут R + W > N для read-your-writes, но без порядка операций — это не линеаризуемость,
а «sloppy quorum» с разрешением конфликтов векторными часами.
Paxos
Paxos (Lamport, 1998/2001) решает консенсус на одно значение. Роли: proposer, acceptor, learner.
ОБЯЗАН предложить "x", а не своё значение! Note over P,A3: Фаза 2 — зафиксировать P->>A1: accept(5, "x") P->>A2: accept(5, "x") A1-->>P: accepted(5, "x") A2-->>P: accepted(5, "x") Note over P: Кворум ⇒ значение "x" выбрано навсегда
Ключевое правило, обеспечивающее safety: proposer не свободен в выборе значения. Если хоть один acceptor в кворуме уже что-то принял, proposer обязан предложить значение с наибольшим номером — это и есть «наследование» решения через пересечение кворумов. Два узких места. Дуэль proposer’ов (livelock): два proposer’а по очереди перебивают номера, ни один не доходит до фазы 2 — ровно то, что предсказывает FLP; лечится рандомизированным backoff и выделенным лидером. Две фазы на решение = 2 RTT; Multi-Paxos выполняет фазу 1 один раз для всего диапазона слотов, дальше каждая команда стоит 1 RTT. Читать: Paxos Made Simple (короткая, но обманчиво трудная) и Paxos Made Live — честный рассказ Google о том, сколько инженерных деталей статья умалчивает.
Raft: тот же консенсус, но понятный
Raft (Ongaro & Ousterhout, 2014) декомпозирует задачу на выборы лидера, репликацию лога и безопасность.
(рандомизированный 150–300 мс) Candidate --> Candidate: split vote —
новый срок, новый таймаут Candidate --> Leader: получено большинство голосов Candidate --> Follower: увидел AppendEntries
с term >= своего Leader --> Follower: увидел сообщение
с большим term Leader --> Leader: heartbeat каждые ~50 мс note right of Candidate Рандомизация таймаутов — обход FLP-livelock: split vote разрешается вероятностно за 1–2 раунда. Запись коммитится, когда реплицирована на большинство И принадлежит текущему сроку лидера. end note
Ключевые инварианты Raft — они объясняют 90% багов самописных реализаций:
- Election Safety: в одном сроке (term) не более одного лидера.
- Leader Append-Only: лидер никогда не перезаписывает и не удаляет записи своего лога.
- Log Matching: если две записи в разных логах имеют одинаковые index и term, логи идентичны
во всех предыдущих записях. Проверяется полями
prevLogIndex/prevLogTerm. - Leader Completeness: если запись закоммичена в сроке
T, она есть в логах всех лидеров сроков> T. Обеспечивается правилом голосования: голос — только за кандидата, чей лог не старее своего (сравнение пары(lastLogTerm, lastLogIndex)). - State Machine Safety: если узел применил запись с индексом
i, никакой другой узел не применит другую запись с индексомi.
Самый коварный нюанс — §5.4.2 статьи: новый лидер не имеет права коммитить запись предыдущего срока только потому, что она реплицирована на большинство — существует сценарий, где такая запись потом перезаписывается. Лечение: запись прошлого срока коммитится косвенно, после коммита хотя бы одной своей записи текущего срока (обычно no-op при старте). Читать: диссертацию Ongaro (там же membership changes через joint consensus и снапшоты) и визуализацию raft.github.io.
Византийские отказы
Если узел может лгать (взлом, баг, повреждение данных), crash-модели мало. Результат
Lamport, Shostak, Pease (1982): согласие возможно тогда и только тогда, когда n ≥ 3f + 1.
Интуиция: корректных n − f, из них до f могли не ответить, и ещё f византийских голосов
надо перевесить. PBFT (Castro & Liskov, 1999) даёт практичный BFT-консенсус за 3 фазы —
pre-prepare / prepare / commit — с O(n²) сообщений: prepare фиксирует порядок внутри
вида, commit — что согласие переживёт смену вида (view change). HotStuff (2018) сводит
смену вида к O(n) через threshold-подписи. Nakamoto-консенсус в блокчейнах — другой класс:
вероятностная финальность и открытое множество участников, а не классические 3f+1 кворумы.
Практический вывод: BFT в корпоративных системах почти никогда не нужен — он стоит кратной
избыточности, а типовые угрозы (битые диски, баги) дешевле закрывать контрольными суммами.
CAP и что он на самом деле утверждает
Формулировка Гильберта–Линч (2002): в асинхронной сети нельзя одновременно обеспечить линеаризуемость (C), доступность всех узлов (A) и устойчивость к разбиению (P). Три частых заблуждения: (1) «можно выбрать CA» — нельзя, разбиения случаются независимо от вашего желания, и реальный выбор делается только во время разбиения, CP или AP; (2) «C в CAP = C в ACID» — нет, в CAP это линеаризуемость (порядок операций), в ACID — сохранение инвариантов схемы; (3) «CAP описывает систему целиком» — обычно нет: платежи CP, лента рекомендаций AP. Полезнее PACELC (Abadi): если Partition, то A или C; иначе (Else) — Latency или Consistency. Вторая половина объясняет, почему даже без отказов кворумная запись через регионы стоит десятки миллисекунд.
Когда консенсус не нужен: CRDT
Если операции коммутативны, порядок согласовывать не нужно — достаточно, чтобы реплики сошлись к одному состоянию. Это CRDT (Shapiro et al., 2011).
class GCounter:
"""Grow-only counter — простейший state-based CRDT:
состояние — вектор счётчиков, merge = поэлементный max."""
def __init__(self, node_id, n):
self.id, self.p = node_id, [0] * n
def inc(self, k=1):
self.p[self.id] += k
def value(self):
return sum(self.p)
def merge(self, other):
# сеть может дублировать, переупорядочивать и терять-переслать — всё равно сойдёмся
self.p = [max(a, b) for a, b in zip(self.p, other.p)]
Требование к merge — быть операцией join в полурешётке: ассоциативность, коммутативность, идемпотентность. Тогда любой порядок доставки даёт одинаковый результат (strong eventual consistency). Известные типы: G-Counter, PN-Counter, LWW-Register, OR-Set, RGA/Yjs для текста. Цена — метаданные (tombstones, версионные векторы) и то, что не все инварианты выразимы: «остаток на счёте ≥ 0» через CRDT не реализуется без резервирования — тут нужен консенсус. Продакшн: Redis Active-Active, Riak, Automerge/Yjs.
Типичные ошибки
Параллельные:
- Считать, что
O(n log n)автоматически параллелизуется. Смотрите на глубинуD. - Дробить до одного элемента без grain size — накладные расходы съедают всё.
- Забывать про пропускную способность памяти: memory-bound алгоритм (arithmetic intensity < 1 flop/byte) на 32 ядрах ускорится в 3 раза, а не в 32.
- Разделяемый счётчик прогресса в горячем цикле — классический false sharing.
- Самописный lock-free без hazard pointers/epoch — почти гарантированный ABA или use-after-free.
- Полагаться на побитовую воспроизводимость суммы float при параллельной редукции.
Распределённые:
- Верить в fallacies of distributed computing (Deutsch/Gosling): сеть надёжна, задержка нулевая, полоса бесконечна, топология неизменна, транспорт бесплатен.
- Ретраи без идемпотентности: если ответ потерян, клиент повторит — нужен idempotency key и дедупликация на сервере. Порядок событий — по логическим часам, не по настенным.
- Распределённая блокировка без fencing token: клиент захватил лок, ушёл в GC-паузу на 20 с, лок истёк, другой клиент взял его — а первый проснулся и пишет. Монотонный токен, проверяемый хранилищем при записи, — единственное надёжное лечение (разбор Kleppmann про Redlock).
- Кластер консенсуса из чётного числа узлов: 4 узла выдерживают тот же 1 отказ, что и 3, но требуют больше подтверждений. Всегда 3, 5, 7.
- Прогонять весь трафик через Raft: пропускная способность ограничена одним лидером и fsync’ом. Данные шардируйте на много raft-групп (CockroachDB ranges, TiKV regions), а через консенсус пускайте метаданные и координацию.
- Не тестировать разбиения. Обязательно: fault injection, Jepsen, детерминированная симуляция (FoundationDB), TLA+ для критичных протоколов.
Как это выглядит в проде
| Система | Что использует | Зачем |
|---|---|---|
| etcd / Kubernetes | Raft | состояние кластера, лидер-элекция контроллеров |
| ZooKeeper | ZAB (близок к Multi-Paxos) | конфигурация, локи, service discovery |
| Kafka (KRaft) | Raft; раньше — ZooKeeper + ISR | метаданные топиков, выбор лидера партиции |
| CockroachDB / TiDB | Raft на каждый range + HLC/2PC | горизонтальное масштабирование ACID |
| Spanner | Paxos-группы + TrueTime | внешняя консистентность на глобальном масштабе |
| Cassandra / DynamoDB | кворумы R+W>N, LWT через Paxos |
AP по умолчанию, CP по требованию |
| Spark / MapReduce | BSP-суперстепы, sample sort в shuffle | пакетная обработка петабайтов |
| CUB / Thrust (GPU) | scan, segmented reduce, radix | базовые примитивы для всего остального |
Отдельно стоит BSP (Bulk Synchronous Parallel, Valiant): вычисления делятся на суперстепы
«локальные вычисления → обмен сообщениями → барьер», время предсказывается как w + g·h + l.
На нём построены Pregel, Giraph, Spark GraphX — параллельные версии обходов из
https://courses.digitable.life/post/algorithms/08-graph-traversal/; Δ-stepping продолжает https://courses.digitable.life/post/algorithms/09-shortest-paths/.
Мини-итог
- Параллельная сложность — пара
(работа W, глубина D);T_p ≤ W/p + D(Брент), параллелизмW/D— потолок полезных ядер. Амдал ограничивает ускорение фиксированной задачи, Густафсон объясняет масштабирование: смотрите, растёт ли последовательная часть вместе сn. - Map / reduce / scan покрывают большинство задач; scan за
O(n)работы иO(log n)глубины — ключ к сортировкам, компактификации и распределению работы. - Синхронизация ломает идеальную модель: false sharing, память как бутылочное горлышко, ABA.
- В распределённой системе всё меняют частичные отказы и отсутствие общих часов: используйте логические/векторные/гибридные часы, не настенные.
- FLP: детерминированный консенсус в асинхронной модели невозможен; практика обходит это
частичной синхронностью (safety всегда, liveness после GST) и рандомизацией таймаутов.
Всё держится на пересечении кворумов:
n ≥ 2f+1для crash,n ≥ 3f+1для византийских. - Paxos и Raft решают одну задачу; Raft выигрывает понятностью. Log Matching, Leader Completeness и правило коммита текущего срока — то, что реально обеспечивает safety.
- Консенсус дорог: коммутативные операции — в CRDT, инварианты — через консенсус на метаданных.
Источники
- Fischer, Lynch, Paterson. Impossibility of Distributed Consensus with One Faulty Process (1985)
- Dwork, Lynch, Stockmeyer. Consensus in the Presence of Partial Synchrony (1988)
- Lamport. Paxos Made Simple (2001); Ongaro, Ousterhout. Raft (2014)
- Castro, Liskov. Practical BFT (1999); Yin et al. HotStuff (2018)
- Gilbert, Lynch. Brewer’s Conjecture… (CAP) (2002)
- Shapiro et al. Convergent and Commutative Replicated Data Types (2011)
- Blelloch. Prefix Sums and Their Applications (1993); Herlihy, Shavit, Luchangco, Spear. The Art of Multiprocessor Programming, 2nd ed. (2020)
- Kleppmann. Designing Data-Intensive Applications (2017) — главы 5, 8, 9; CLRS 4e, глава 26
- Jepsen — разборы реальных нарушений консистентности
Что дальше
Мы разобрали, как распараллелить вычисление и как договориться между машинами. Осталась третья ось масштаба: что делать, когда данные не помещаются — ни в память, ни даже в один проход. Дальше — алгоритмы, которые видят поток однократно и хранят логарифм от него, и те, для которых узкое место не процессор, а диск.