Алгоритмы Параллельные и распределённые алгоритмы, консенсус
0%

Параллельные и распределённые алгоритмы, консенсус

Параллельные и распределённые алгоритмы, консенсус

Все предыдущие статьи трека предполагали одну машину, одно ядро и надёжную память. Реальность другая: у сервера 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». Если нет — масштабирование будет.

Три кита: 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.

Up-sweep и down-sweep в scan Блеллоха

Алгоритм Работа 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:

  1. Каждый из p потоков берёт s случайных элементов; сортируем p·s сэмплов, берём p−1 splitter’ов.
  2. Каждый поток бинпоиском раскладывает свой чанк по p корзинам; гистограмма + exclusive scan дают точные смещения записи — без единого мьютекса.
  3. Корзина 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.

Ключевое правило, обеспечивающее 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) декомпозирует задачу на выборы лидера, репликацию лога и безопасность.

Ключевые инварианты 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, инварианты — через консенсус на метаданных.

Источники


Что дальше

Мы разобрали, как распараллелить вычисление и как договориться между машинами. Осталась третья ось масштаба: что делать, когда данные не помещаются — ни в память, ни даже в один проход. Дальше — алгоритмы, которые видят поток однократно и хранят логарифм от него, и те, для которых узкое место не процессор, а диск.

Потоковые алгоритмы и алгоритмы внешней памяти

Нашли неточность? Выделите фрагмент текста — рядом появится жучок.

Нужен разбор именно вашей ситуации?

Статья описывает общий случай. Если у вас частный — можно разобрать его отдельно, платно. А если не хватает целого материала, предложите тему: её оплачивают вскладчину, и она выходит открытой для всех.

Доска запросов