Потоковые алгоритмы и алгоритмы внешней памяти
Почти вся классическая теория алгоритмов молча предполагает одну вещь: данные лежат в оперативной памяти, и обращение к любому элементу стоит одинаково. Эта модель называется RAM-моделью, и она отлично работала последние полвека — ровно до тех пор, пока не появились задачи, где данных больше, чем памяти. Иногда — на порядок. Иногда — данных вообще нет как объекта: есть бесконечный поток событий, который проходит мимо вас один раз и исчезает.
Тут RAM-модель ломается не в деталях, а в основании. Алгоритм за O(n log n) бесполезен,
если он требует O(n) памяти, а n — это 40 миллиардов кликов за сутки. Алгоритм с
«всего лишь» случайным доступом к массиву становится в 10 000 раз медленнее, если массив
лежит на SSD, а не в RAM.
Поэтому появились две другие модели вычислений, и обе — про одно и то же ограничение, только с разных сторон:
- Потоковая модель (streaming). Данные проходят мимо один раз (реже — константное
число раз), памяти у нас
O(polylog n)илиO(√n), и мы согласны на приближённый ответ с вероятностной гарантией. - Модель внешней памяти (external memory / I/O-модель). Данные лежат на диске, память ограничена, но данные никуда не убегают — можно читать сколько угодно раз. Считаем не операции процессора, а блоки ввода-вывода.
Между ними есть глубокая связь: и там и там мы платим за то, что «дёшево» и «дорого» поменялись местами. И обе модели — не академическая экзотика, а буквально то, на чём построены ClickHouse, Redis, Kafka Streams, RocksDB, Spark и любой биллинг, который считает уникальных пользователей.
Зачем это нужно на практике
Прежде чем нырять в математику — список задач, где без этих техник вы физически не уложитесь в бюджет:
- Уникальные пользователи. «Сколько уникальных
user_idбыло в этом сегменте за 30 дней, по любому срезу измерений» — точный ответ требует хранить множество всех id на каждый срез. Комбинаторный взрыв. HyperLogLog решает это в 12 КБ на срез. - Топ-запросы, топ-ошибки, топ-IP. Найти 100 самых частых ключей в потоке из миллиардов — Misra–Gries или Count–Min Sketch, память константная.
- Латентности и SLO. p99 времени ответа по 10 млрд запросов. Хранить все значения нельзя, а среднее бесполезно — нужны потоковые квантили: t-digest, KLL, DDSketch.
- Антифрод и rate limiting. «Этот IP уже видел 5000 разных карт» — distinct count в скользящем окне, в памяти прокси.
- Дедупликация и защита от лишних чтений. Bloom-фильтр перед походом в базу.
- Сортировка терабайта. ETL, построение индексов, join больших таблиц — это внешняя
сортировка, и её асимптотика не
O(n log n), аO((N/B)·log_{M/B}(N/B)). - Базы данных. B-дерево и LSM-дерево — это не «структуры данных вообще», а именно структуры данных внешней памяти, оптимизированные под число блочных операций.
Часть I. Потоковая модель
Формализация: что именно нам разрешено
Поток — это последовательность обновлений a₁, a₂, …, a_m. Классификация Мутукришнана
различает три режима, и путать их опасно:
| Модель | Обновления | Пример | Что можно |
|---|---|---|---|
| Time series | элемент i приходит ровно один раз, по порядку |
показания датчика по минутам | почти всё |
| Cash register | f[i] += c, где c > 0 |
счётчики событий, клики | Count–Min с min, HLL, Misra–Gries |
| Turnstile | f[i] += c, c любого знака |
вход/выход из системы, отмены заказов | Count Sketch, AMS; Count–Min ломается |
Ресурсы: память S, время обработки одного элемента t, число проходов p. Хороший
потоковый алгоритм — это S = O(polylog n) или O(ε⁻¹ polylog n), t = O(1)
амортизированно, p = 1.
Ответ почти всегда (ε, δ)-приближённый: с вероятностью не менее 1 − δ относительная
(или аддитивная) ошибка не превышает ε. Это не «алгоритм иногда врёт» — это
контролируемая ошибка, и вы сами задаёте цену в памяти за нужную точность.
Почему точный ответ невозможен: нижние границы
Прежде чем изучать хитрые приближённые схемы, полезно понять, что они не от лени. Точное решение доказуемо требует линейной памяти.
Классический аргумент — через коммуникационную сложность. Пусть Алиса имеет множество
A ⊆ {1..n}, Боб — множество B. Задача DISJOINTNESS (пересекаются ли они)
требует Ω(n) бит коммуникации даже для рандомизированных протоколов. Теперь
предположим, что есть потоковый алгоритм, который точно считает число различных
элементов F₀ в памяти S. Алиса пропускает через него свои элементы, шлёт Бобу
состояние (S бит), Боб дописывает свои. Если F₀(A ∪ B) = |A| + |B|, множества
не пересекаются. Значит S = Ω(n).
Тот же приём даёт нижние границы для точных частых элементов, точной медианы за один
проход (Ω(n) для точного, а для приближённой хватает O(ε⁻¹ log(εn))) и для F_k
при k > 2. Вывод простой и важный: если кто-то предлагает точно посчитать уникальных
за константную память — он ошибается или считает не то.
Reservoir sampling: равномерная выборка из потока неизвестной длины
Самый базовый приём. Нужно получить k элементов, каждый из которых выбран равномерно
из всего потока, при этом длина потока n заранее неизвестна и не помещается в память.
import random
from typing import Iterable, List, TypeVar
T = TypeVar("T")
def reservoir_sample(stream: Iterable[T], k: int) -> List[T]:
"""Алгоритм R (Vitter, 1985). Память O(k), один проход, время O(n)."""
reservoir: List[T] = []
for i, item in enumerate(stream):
if i < k:
reservoir.append(item) # первые k берём безусловно
else:
# i-й элемент (0-индексация) попадает в выборку с вероятностью k/(i+1)
j = random.randrange(i + 1)
if j < k:
reservoir[j] = item # вытесняет случайного жильца
return reservoir
Доказательство корректности — индукция по n. Утверждение: после обработки n ≥ k
элементов каждый из них лежит в резервуаре с вероятностью k/n. База очевидна.
Шаг: элемент n+1 попадает с вероятностью k/(n+1) — верно по построению. Старый
элемент выживает, если новый либо не попал (1 − k/(n+1)), либо попал, но вытеснил
не его (k/(n+1) · (k−1)/k). Суммарно (k/n)·(1 − 1/(n+1)) = k/(n+1). ∎
Сложность: время O(n), память O(k), ровно один проход. Для взвешенной выборки
(вероятность пропорциональна весу) используют A-Res: для каждого элемента считают
ключ u^(1/w) где u ~ U(0,1), и держат k наибольших ключей в min-heap — O(n log k).
На практике reservoir sampling — это основа сэмплирования трейсов (tail-based sampling в OpenTelemetry), «случайные 1000 строк из огромного файла» и обучающих подвыборок в ML.
Частые элементы: Misra–Gries
Задача Heavy Hitters: найти все элементы с частотой > N/k. Идея Misra–Gries — обобщение
алгоритма голосования Бойера–Мура на k кандидатов. Держим не более k−1 счётчиков;
если пришёл новый ключ, а места нет — уменьшаем все счётчики на единицу и выкидываем
обнулившиеся. Интуиция: это как попарное «взаимное уничтожение» голосов. Элемент,
встречающийся чаще N/k, не может быть уничтожен полностью.
from collections import Counter
from typing import Dict, Hashable, Iterable
def misra_gries(stream: Iterable[Hashable], k: int) -> Dict[Hashable, int]:
"""Не более k-1 счётчиков. Гарантия: f(x) - N/k <= counters[x] <= f(x).
Время O(1) амортизированно на элемент, память O(k).
Все элементы с частотой > N/k гарантированно есть в результате
(но в результате могут быть и лишние — нужен второй проход для проверки).
"""
counters: Dict[Hashable, int] = {}
for x in stream:
if x in counters:
counters[x] += 1
elif len(counters) < k - 1:
counters[x] = 1
else:
# декремент всех: за всю работу таких «шагов» не больше N/k,
# поэтому амортизированная стоимость остаётся O(1)
for key in list(counters):
counters[key] -= 1
if counters[key] == 0:
del counters[key]
return counters
Почему оценка занижена не более чем на N/k. Каждый шаг с декрементом уничтожает
k единиц веса: одну от нового элемента и k−1 из счётчиков. Всего веса N, значит
таких шагов не более N/k. Счётчик конкретного x уменьшается не чаще одного раза за
шаг, следовательно недосчёт ≤ N/k. ∎
Ключевое свойство для распределённых систем — мержабельность: два résumé Misra–Gries
можно сложить поэлементно, оставить k−1 наибольших и вычесть k-е по величине значение;
гарантия сохранится. Именно поэтому такие структуры называют mergeable summaries.
Родственники: Space-Saving (Metwally et al.) — вместо декремента переиспользует
счётчик с минимальным значением; на практике точнее на скошенных распределениях и лежит
в основе topK в ClickHouse и Redis.
Count–Min Sketch: частоты в константной памяти
Misra–Gries отвечает на вопрос «кто чаще всех». Count–Min Sketch отвечает на другой: «какова частота вот этого конкретного ключа» — для любого ключа, в том числе редкого.
Структура — матрица счётчиков d × w и d попарно независимых хеш-функций. Обновление
(x, c) прибавляет c в ячейку CM[i][h_i(x)] каждой строки. Запрос возвращает
минимум по строкам.
import math
from hashlib import blake2b
from typing import Hashable, List
class CountMinSketch:
"""Cormode & Muthukrishnan, 2005.
Память O(w*d). Обновление и запрос — O(d).
Гарантия (только для неотрицательных обновлений, cash register):
f(x) <= estimate(x) всегда,
estimate(x) <= f(x) + eps * N с вероятностью >= 1 - delta.
"""
def __init__(self, eps: float = 0.001, delta: float = 1e-4) -> None:
self.w = int(math.ceil(math.e / eps)) # ширина: e/eps
self.d = int(math.ceil(math.log(1.0 / delta))) # глубина: ln(1/delta)
self.table: List[List[int]] = [[0] * self.w for _ in range(self.d)]
self.total = 0
def _hashes(self, item: Hashable):
# одна 32-байтная хеш-выжимка нарезается на d независимых индексов
digest = blake2b(repr(item).encode("utf-8"), digest_size=4 * self.d).digest()
for i in range(self.d):
chunk = digest[4 * i: 4 * i + 4]
yield i, int.from_bytes(chunk, "big") % self.w
def add(self, item: Hashable, count: int = 1) -> None:
self.total += count
for row, col in self._hashes(item):
self.table[row][col] += count
def estimate(self, item: Hashable) -> int:
return min(self.table[row][col] for row, col in self._hashes(item))
def merge(self, other: "CountMinSketch") -> "CountMinSketch":
"""Мержабельность: поэлементная сумма при одинаковых w, d и хешах."""
assert (self.w, self.d) == (other.w, other.d)
for r in range(self.d):
for c in range(self.w):
self.table[r][c] += other.table[r][c]
self.total += other.total
return self
Разбор гарантии. Зафиксируем строку i. Оценка CM[i][h_i(x)] = f(x) + X, где
X — суммарный вес чужих ключей, попавших в ту же ячейку. По попарной независимости
E[X] = (N − f(x))/w ≤ N/w = εN/e. По неравенству Маркова P(X ≥ εN) ≤ 1/e.
Строки независимы, поэтому все d строк ошибаются одновременно с вероятностью
≤ e^{−d} = δ. Минимум ошибается только если ошиблись все. ∎
Заметьте, где именно используется односторонность: коллизии могут только добавить
вес, никогда не отнять. Поэтому min — правильный агрегат, и поэтому CMS корректен
только в cash-register модели. Если у вас бывают отрицательные обновления
(отмены, удаления), берите Count Sketch (Charikar–Chen–Farach-Colton): там каждая
ячейка обновляется как + s_i(x)·c со случайным знаком s_i(x) ∈ {−1, +1}, оценка —
медиана s_i(x)·CS[i][h_i(x)], и ошибка становится двусторонней, но несмещённой.
Практический размер: ε = 0.001, δ = 10⁻⁴ даёт w ≈ 2719, d = 10, то есть ~27 000
счётчиков — около 110 КБ на 4-байтных int. Это независимо от того, миллион у вас
ключей или триллион.
Типичная ошибка новичка: спрашивать у CMS частоты редких ключей и верить ответу.
Аддитивная ошибка εN пропорциональна всему объёму потока. Если N = 10⁹ и
ε = 0.001, то ошибка ±10⁶ — для ключа с настоящей частотой 50 это мусор. CMS хорош
для «тяжёлых» ключей, для хвоста — нет.
HyperLogLog: сколько различных элементов
Задача F₀ (distinct count). Точное решение — Ω(n) памяти, как мы видели. Приближённое
решение — одна из самых красивых конструкций в алгоритмах.
Интуиция. Возьмём хорошую хеш-функцию и посмотрим на двоичное представление хешей.
Примерно у половины хешей первый бит нулевой, у четверти — два нуля подряд, у одного
из 2^k — k нулей. Если в потоке максимальный наблюдённый «ранг» (позиция первой
единицы) оказался R, то разумно предположить, что различных элементов было около 2^R.
Проблема: одна такая оценка чудовищно шумная — дисперсия огромна. Решение (Flajolet et al.)
— стохастическое усреднение: первые p бит хеша выбирают один из m = 2^p регистров,
остальные биты дают ранг, и в каждом регистре хранится максимум. Итоговая оценка —
гармоническое среднее 2^{R_j} с поправочным коэффициентом. Гармоническое среднее
подавляет влияние выбросов, чем и объясняется выигрыш относительно LogLog.
Стандартная ошибка: σ ≈ 1.04 / √m. При p = 14 (16 384 регистра по 6 бит ≈ 12 КБ)
получаем ~0.81 % — и это на любом количестве уникальных значений вплоть до миллиардов.
import math
from hashlib import blake2b
from typing import Hashable
class HyperLogLog:
"""Flajolet, Fusy, Gandouet, Meunier (2007) + линейная коррекция на малых мощностях.
Память: m регистров по 6 бит (в этой реализации — по байту).
add: O(1). count: O(m). Стандартная ошибка ~ 1.04/sqrt(m).
"""
def __init__(self, p: int = 14) -> None:
assert 4 <= p <= 18
self.p = p
self.m = 1 << p
self.registers = bytearray(self.m)
if self.m == 16:
self.alpha = 0.673
elif self.m == 32:
self.alpha = 0.697
elif self.m == 64:
self.alpha = 0.709
else:
self.alpha = 0.7213 / (1.0 + 1.079 / self.m)
def _hash64(self, item: Hashable) -> int:
return int.from_bytes(blake2b(repr(item).encode("utf-8"), digest_size=8).digest(), "big")
def add(self, item: Hashable) -> None:
x = self._hash64(item)
idx = x >> (64 - self.p) # старшие p бит — номер регистра
bits = 64 - self.p
w = x & ((1 << bits) - 1) # остаток — «хвост» для ранга
# ранг = позиция первой единицы слева в bits-битном слове, 1-индексация
rank = bits - w.bit_length() + 1 if w else bits + 1
if rank > self.registers[idx]:
self.registers[idx] = rank
def count(self) -> float:
z = sum(2.0 ** -r for r in self.registers) # обратное гармоническое среднее
estimate = self.alpha * self.m * self.m / z
if estimate <= 2.5 * self.m:
zeros = self.registers.count(0)
if zeros:
# на малых мощностях точнее обычное linear counting
return self.m * math.log(self.m / zeros)
return estimate
def merge(self, other: "HyperLogLog") -> "HyperLogLog":
"""Объединение множеств = поэлементный максимум регистров. Идемпотентно."""
assert self.p == other.p
for i in range(self.m):
if other.registers[i] > self.registers[i]:
self.registers[i] = other.registers[i]
return self
Свойства, которые делают HLL незаменимым в проде:
- Мержабельность через
max. Можно считать HLL на 500 воркерах и слить в один. Операция идемпотентна и коммутативна — то есть HLL это, по сути, CRDT, и повторная доставка события ничего не портит. - Фиксированный размер. 12 КБ — и точка. Бюджет памяти известен заранее.
- Иерархичность. HLL за час можно слить в HLL за сутки. Точные множества так не агрегируются без хранения всех id.
Чего HLL не умеет — и это ошибка №1 в проде: пересечения. |A ∩ B| = |A| + |B| − |A ∪ B|
формально верно, но при вычитании двух величин с относительной ошибкой 1 % ошибка
результата может превысить сам результат, если пересечение мало. Для пересечений нужны
другие структуры (например, Theta Sketch из Apache DataSketches, где поддерживается
корректная алгебра множеств).
Инженерные улучшения из статьи Google «HyperLogLog in Practice» — 64-битный хеш
(убирает поправку на большие мощности), sparse-представление на малых мощностях
(экономит память, когда уникальных мало) и эмпирическая коррекция смещения. Именно
этот вариант — HLL++ — реализован в BigQuery, Redis (PFADD/PFCOUNT),
ClickHouse (uniqHLL12, uniqCombined) и Presto.
Квантили: p50, p95, p99 из потока
Медиана — известный камень преткновения: она не разложима (не мержабельна как сумма), и точное вычисление за один проход требует линейной памяти. На практике используют три семейства:
- Greenwald–Khanna (GK) — детерминированный,
ε-приближение ранга, памятьO((1/ε)·log(εN)). Хранит отсортированный набор кортежей(значение, g, Δ), гдеg— вклад в ранг,Δ— неопределённость. Слияние возможно, но неудобно. - KLL sketch (Karnin–Lang–Liberty, 2016) — рандомизированный, оптимальный:
O((1/ε)·log log(1/δ)). Идея — иерархия «компакторов», каждый уровень выбрасывает половину элементов, случайно оставляя чётные или нечётные. Полностью мержабельный. Реализован в Apache DataSketches. - t-digest (Dunning) — кластеры точек с переменным весом, гарантия относительной
ошибки на хвостах: чем ближе к p99.9, тем точнее. Именно то, что нужно для SLO.
Используется в Elasticsearch (
percentiles), Spark, Prometheus-совместимых стеках. - DDSketch (Datadog, 2019) — relative-error гарантия по значению, не по рангу; логарифмические бакеты, мержабельность тривиальная.
Практическое правило: для мониторинга латентностей берите t-digest или DDSketch (вам важен хвост), для аналитики общего назначения — KLL. Наивные «гистограммы с фиксированными бакетами» (как в классических Prometheus histogram) тоже работают, но требуют угадать границы бакетов заранее, а p99, посчитанный линейной интерполяцией внутри широкого бакета, — это оценка с неизвестной ошибкой.
Скользящее окно: DGIM
Отдельная сложность — «за последний час», а не «за всё время». Проблема: мы не можем хранить элементы, чтобы вычитать их при выходе из окна.
Алгоритм DGIM (Datar, Gionis, Indyk, Motwani) считает число единиц в скользящем окне
длины N за O(log² N) бит с относительной ошибкой ≤ 50 % (и ≤ ε при увеличении
числа бакетов на уровень). Идея — экспоненциальные бакеты: хранятся временные метки
конца групп из 1, 1, 2, 2, 4, 4, 8, 8, … единиц. Чем старше данные, тем грубее
разрешение — что логично, ведь их вклад скоро исчезнет.
Более практичный компромисс: tumbling/sliding windows поверх мержабельных скетчей.
Считаем по одному HLL на минуту, окно «за час» = слияние 60 скетчей. Память O(60·12КБ),
точность прежняя, и это ровно то, что делают Flink и Kafka Streams.
Как всё это живёт в распределённой системе
Главное свойство, из-за которого скетчи победили: они образуют коммутативную полугруппу относительно слияния. Порядок агрегации не важен, промежуточные результаты можно кэшировать и переиспользовать.
Сравните с точным COUNT(DISTINCT ...): там координатору пришлось бы получить все
уникальные id со всех шардов — трафик, пропорциональный мощности множества. Это и есть
причина, по которой в аналитических базах uniq() дешёвый, а uniqExact() — дорогой.
Про модели согласованности, коммутативные слияния и CRDT подробнее — в статье https://courses.digitable.life/post/algorithms/16-parallel-and-distributed/. Вероятностный аппарат за гарантиями скетчей (неравенства Маркова, Чебышёва, Чернова, попарная независимость) разобран в https://courses.digitable.life/post/algorithms/14-randomized-algorithms/, а свойства хеш-функций — в https://courses.digitable.life/post/algorithms/13-number-theory/.
Как выбрать структуру под задачу
~10 бит на элемент"] C -->|Нет| C2["Точное множество или
индекс на диске"] B -->|"Сколько различных?"| D{"Нужны пересечения множеств?"} D -->|Нет| D1["HyperLogLog++
12 КБ, ошибка 0.8%"] D -->|Да| D2["Theta Sketch / KMV
алгебра множеств"] B -->|"Частота ключа"| E{"Бывают отрицательные
обновления?"} E -->|Нет| E1["Count-Min Sketch
оценка сверху, min"] E -->|Да| E2["Count Sketch
несмещённая, median"] B -->|"Кто самый частый"| F["Misra-Gries /
Space-Saving"] B -->|"p95 и p99"| G{"Важен хвост
распределения?"} G -->|Да| G1["t-digest / DDSketch
относительная ошибка"] G -->|Нет| G2["KLL sketch
оптимальная память"] B -->|"Репрезентативная выборка"| H["Reservoir sampling
алгоритм R или A-Res"] style D1 fill:#5fa87a,fill-opacity:0.25 style E1 fill:#6a7fb5,fill-opacity:0.25 style G1 fill:#b5842c,fill-opacity:0.25
Часть II. Модель внешней памяти
Модель DAM: считаем блоки, а не операции
Модель Disk Access Machine (Aggarwal & Vitter, 1988) описывается тремя параметрами:
N— число записей во входных данных;M— число записей, помещающихся во внутреннюю память;B— число записей в одном блоке передачи (страница диска, строка кэша).
Стоимость алгоритма — число операций ввода-вывода, то есть переносов блока размера
B между внешней и внутренней памятью. Вычисления внутри памяти считаются бесплатными.
Это грубое приближение, но оно поразительно точно предсказывает реальную
производительность, потому что разрыв в скорости между уровнями памяти составляет
2–5 порядков.
Три базовые величины, которые надо помнить наизусть:
| Операция | I/O-сложность | Смысл |
|---|---|---|
| Сканирование | scan(N) = Θ(N/B) |
прочитать всё подряд |
| Сортировка | sort(N) = Θ((N/B)·log_{M/B}(N/B)) |
и это нижняя граница |
| Поиск | Θ(log_B N) |
B-дерево, и это тоже оптимум |
Обратите внимание на две вещи. Во-первых, log в сортировке имеет основание M/B,
а не 2 — фан-ин слияния равен числу буферов, влезающих в память. Во-вторых, поиск стоит
log_B N, а не log₂ N: за одно чтение мы получаем B ключей сразу. Именно поэтому
B-дерево с фактором ветвления 100+ бьёт красно-чёрное дерево на диске в разы, хотя в
RAM-модели у них одинаковая асимптотика.
Числовой пример. N = 10¹⁰ записей, M = 10⁹, B = 10⁴. Тогда
N/B = 10⁶, M/B = 10⁵, и log_{M/B}(N/B) = log(10⁶)/log(10⁵) = 1.2. То есть
один-два прохода слияния. Практический вывод, который стоит запомнить: на реальных
размерах памяти внешняя сортировка почти всегда двухфазная, и sort(N) мало отличается
от 2·scan(N). Отсюда же берётся эвристика «сортировка терабайта стоит примерно
как три чтения терабайта».
Внешняя сортировка слиянием
Фаза 1 — формирование ранов: читаем куски по M записей, сортируем в памяти любым
быстрым внутренним алгоритмом (см. https://courses.digitable.life/post/algorithms/02-sorting/), пишем обратно.
Получаем N/M отсортированных ранов, стоимость 2N/B.
Фаза 2 — многопутевое слияние: k = M/B − 1 входных буферов по B записей плюс
один выходной. Берём минимум из голов через min-heap за O(log k), пишем в выходной
буфер, при заполнении сбрасываем на диск. Каждый проход уменьшает число ранов в k раз.
import heapq
import itertools
import os
import tempfile
from typing import Iterator, List
def external_sort(src_path: str, dst_path: str, chunk_lines: int = 1_000_000) -> None:
"""Двухфазная внешняя сортировка строк.
Фаза 1: N/M ранов, стоимость 2N/B I/O.
Фаза 2: k-путевое слияние, k = число открытых файлов.
Итого при N/M <= k — ровно 4N/B операций ввода-вывода.
"""
run_paths: List[str] = []
# --- фаза 1: раны ---
with open(src_path, "r", encoding="utf-8") as src:
while True:
# islice читает ровно chunk_lines строк — это наш «размер памяти» M
chunk = list(itertools.islice(src, chunk_lines))
if not chunk:
break
chunk.sort() # внутренняя сортировка, I/O не тратим
fd, path = tempfile.mkstemp(suffix=".run")
with os.fdopen(fd, "w", encoding="utf-8") as run:
run.writelines(chunk)
run_paths.append(path)
# --- фаза 2: k-путевое слияние ---
files = [open(p, "r", encoding="utf-8") for p in run_paths]
try:
with open(dst_path, "w", encoding="utf-8") as dst:
# heapq.merge — ленивое слияние через кучу, память O(k), а не O(N)
dst.writelines(heapq.merge(*files))
finally:
for f in files:
f.close()
for p in run_paths:
os.unlink(p)
Практические тонкости, которые решают всё:
- Размер буфера
B. Слишком маленькие буферы → случайный доступ и деградация на HDD. Слишком большие → фан-ин падает, появляется лишний проход. Оптимум обычно 1–8 МБ. - Replacement selection вместо простой сортировки чанков даёт раны средней длины
2M(за счёт «просачивания» подходящих записей) — вдвое меньше ранов. На современных SSD выигрыш редко окупает сложность, но в классических СУБД приём применялся десятилетиями. - Сортировка ключей, а не записей. Если запись весит 4 КБ, сортируйте пары
(ключ, смещение), а потом делайте один проход материализации. - Кодирование: сортировка по префиксу ключа (poor man’s normalized keys) резко сокращает число полных сравнений.
Postgres делает ровно это: при превышении work_mem планировщик переключается на
external merge Disk: NNNkB — увидите в EXPLAIN ANALYZE. Spark в shuffle-стадии
пишет отсортированные spill-файлы и сливает их. Утилита GNU sort управляется
флагами -S (это M) и --batch-size (это k).
Нижняя граница sort(N) — откуда берётся log_{M/B}
Интуиция аргумента Aggarwal–Vitter: каждая операция ввода-вывода приносит в память B
элементов, и максимум информации, который мы можем «узнать» об их взаимном порядке
относительно уже находящихся в памяти M элементов, ограничен log(C(M, B) · B!) бит.
Всего для сортировки нужно log(N!) ≈ N log N бит. Деление одного на другое и даёт
Ω((N/B)·log_{M/B}(N/B)). Подробнее — в оригинальной статье и в обзоре Виттера
(ссылки в конце).
Отсюда важное следствие для проектирования: если ваша задача сводится к сортировке
или к перестановке, вы не сможете уложиться в scan(N). Задача «переставить N
элементов по заданной перестановке» стоит min(N, sort(N)) — то есть на диске бывает
дешевле обрабатывать элементы по одному, чем сортировать. Это неинтуитивно и регулярно
удивляет инженеров, оптимизирующих ETL.
B-деревья, LSM-деревья и амплификация
Два канонических индекса внешней памяти решают одну задачу с разными приоритетами.
B-дерево оптимизирует чтения: O(log_B N) I/O на точечный поиск, узлы размером с
страницу. Но запись случайного ключа стоит как минимум одно чтение и одну запись
страницы — вставка 100-байтной записи трогает 8 КБ. Это write amplification ≈ 80.
LSM-дерево (O’Neil et al., 1996) переворачивает компромисс: все записи идут в memtable в памяти, сбрасываются на диск последовательными файлами (SSTable), и лишь фоновая компакция приводит их в порядок. Запись становится последовательной и дешёвой, зато чтение может потребовать заглянуть в несколько уровней.
Три вида «амплификации», которыми меряют такие движки:
| Метрика | B-дерево | LSM (leveled) | LSM (tiered) |
|---|---|---|---|
| Write amplification | высокая (страница на запись) | средняя (O(L) перезаписей) |
низкая |
| Read amplification | низкая (log_B N) |
средняя (уровни + bloom) | высокая |
| Space amplification | ~1.5× (фрагментация) | ~1.1× | высокая (дубликаты) |
Bloom-фильтр на каждый SSTable — обязательный элемент: он отсекает походы на уровни,
где ключа нет, превращая read amplification из O(L) в «почти 1». Это, кстати,
идеальная иллюстрация связи двух частей статьи: потоковая структура (Bloom) используется
для экономии операций внешней памяти.
Отдельно стоит Bε-дерево (Bender et al.) — узлы содержат не только ключи, но и
буферы отложенных обновлений. Оно достигает вставки за O((1/(B^{1−ε}))·log_B N) при
поиске O(log_B N), то есть даёт настраиваемый компромисс между B-деревом и LSM.
На нём построены TokuDB и файловая система BetrFS.
Cache-oblivious алгоритмы: оптимальность без знания B и M
Проблема I/O-модели: чтобы настроить алгоритм, надо знать B и M. Но иерархия памяти
многоуровневая (L1/L2/L3/RAM/SSD), у каждого уровня свои параметры, и в облаке они ещё
и меняются.
Cache-oblivious модель (Frigo, Leiserson, Prokop, Ramachandran, 1999) требует
асимптотической оптимальности при любых B и M, о которых алгоритм не знает.
Достигается это рекурсивным делением задачи: рано или поздно подзадача становится
размером с кэш, и с этого момента всё идеально.
Каноничные примеры:
- Раскладка van Emde Boas для статического бинарного дерева: рекурсивно кладём
верхнюю половину высоты, затем поддеревья. Поиск стоит
O(log_B N)без знанияB. На практике даёт ускорение бинарного поиска в разы — сравните с обсуждением Eytzinger-раскладки в https://courses.digitable.life/post/algorithms/03-searching-and-binary-search/. - Funnelsort — cache-oblivious аналог merge sort с
k-мержерами, достигающий оптимальныхO((N/B)·log_{M/B}(N/B)). - Рекурсивное умножение матриц делением на 8 подзадач:
O(n³/(B√M))вместоO(n³/B)у наивного тройного цикла — то самое, почему блочные алгоритмы быстрее. - Рекурсивный обход при перемножении и transposition-friendly раскладки — практическая база для BLAS.
Практический вывод для повседневного кода: рекурсивное деление задачи пополам почти бесплатно даёт кэш-эффективность. Если у вас есть выбор между итеративным проходом с большим шагом и рекурсивным делением — второе часто быстрее, даже если операций больше. Детали про кэш-линии, префетч и измерения — в https://courses.digitable.life/post/algorithms/18-practical-optimization/.
Графы во внешней памяти
Отдельная боль: обход графа по своей природе — случайный доступ. Наивный BFS
(см. https://courses.digitable.life/post/algorithms/08-graph-traversal/) на диске стоит O(V + E) операций ввода-вывода,
то есть по одной на ребро — катастрофа.
Приёмы, которые применяют:
- Semi-external: если массив меток вершин (
Vбит или байт) влезает в память, а рёбра — нет, всё резко упрощается: рёбра читаем последовательно, метки трогаем в RAM. Для миллиарда вершин это ~1 ГБ — вполне реально. - Munagala–Ranade BFS: строим уровни через сортировку и слияние —
O(V + sort(E))вместоO(V + E). - Edge-centric обработка (модель X-Stream, GraphChi): вместо «для вершины пройти
её рёбра» делаем «прочитать все рёбра последовательно, сгенерировать обновления,
отсортировать, применить». Меняем случайный доступ на сортировку — и выигрываем,
потому что
sort(E) ≪ Eв I/O-метрике. - Разбиение на shard’ы по диапазонам вершин так, чтобы «окно» помещалось в память.
Это тот редкий случай, когда добавление сортировки в алгоритм ускоряет его.
Типичные ошибки
- Верить CMS на редких ключах. Ошибка
εNабсолютная и зависит от объёма всего потока. Проверяйте: еслиεNсравнимо с интересующей вас частотой — структура не та. - Использовать Count–Min при удалениях. Гарантия
minдержится только на неотрицательных обновлениях. С удалениями нужен Count Sketch с медианой. - Считать пересечения HLL через включение-исключение. На маленьких пересечениях результат — шум. Берите Theta Sketch или храните точные множества для этих срезов.
- Смешивать скетчи с разными параметрами. HLL с
p=12иp=14не сливаются напрямую (нужен downgrade), CMS с разнымиw/dили разными сидами хешей — тем более. В проде это выражается в «после апгрейда сервиса метрика уников сошла с ума». - Слабая хеш-функция. Все гарантии скетчей опираются на попарную (а иногда 4-wise)
независимость.
hash()из Python рандомизирован по процессам, CRC32 плохо перемешивает, аString.hashCode()в Java даёт коллизии на структурированных ключах. Берите MurmurHash3, xxHash или BLAKE2 — и фиксируйте сид, иначе перезапуск сервиса обнулит совместимость скетчей. - Считать, что
O(N/B)— это «линейно, значит быстро». КонстантаB— это разница между 2 минутами и 6 часами. Всегда проверяйте, действительно ли доступ последовательный:iostat,blktrace, счётчики страниц Postgres. - Забыть про фан-ин при внешней сортировке. Если
kмал (мало памяти или мало файловых дескрипторов —ulimit -n!), число проходов растёт, и время удваивается. - Тестировать на данных, влезающих в память. Классика: всё летает на 10 ГБ и падает
на 500 ГБ, потому что где-то в коде притаился
list(...)илиGROUP BYбез spill. - Игнорировать skew. Zipf-распределение ключей ломает и равномерность бакетов, и балансировку шардов. Скетчи как раз к skew устойчивы — а вот наивные хеш-партиции нет.
- Путать «мержабельно» и «идемпотентно». HLL идемпотентен (повторное добавление элемента ничего не меняет), CMS — нет: повторная доставка события удвоит счётчик. В системах at-least-once это принципиальная разница.
Как это применяют в проде
- ClickHouse.
uniq()— адаптивная схема,uniqCombined()— комбинация хеш-таблицы, массива и HLL с переключением по мощности,uniqHLL12— чистый HLL,quantileTDigest— t-digest,topK— Space-Saving. ТипAggregateFunction(uniq, UInt64)вAggregatingMergeTreeпозволяет хранить сам скетч в колонке и досливать при merge — это ровно та мержабельность, о которой шла речь. - Redis.
PFADD/PFCOUNT/PFMERGE— HyperLogLog в 12 КБ на ключ, со sparse- представлением для малых мощностей. - Apache DataSketches (Yahoo/Verizon → Apache) — эталонная библиотека: Theta, HLL, KLL, Frequent Items, Tuple sketches. Интегрирована в Druid, Hive, Pig, Spark.
- Apache Flink / Kafka Streams. Оконная агрегация поверх мержабельных состояний, RocksDB как state backend (то есть LSM-дерево прямо в стриминговом движке).
- RocksDB / LevelDB / Cassandra / ScyllaDB / HBase. LSM-деревья, bloom-фильтры на SSTable, настраиваемые стратегии компакции (leveled vs universal/tiered).
- PostgreSQL. B-деревья для индексов, external merge sort при нехватке
work_mem,HashAggregateсо spill на диск начиная с 13-й версии, расширениеpostgresql-hll. - Google BigQuery.
APPROX_COUNT_DISTINCT(HLL++),APPROX_QUANTILES,APPROX_TOP_COUNT— все три из этой статьи. - Prometheus / VictoriaMetrics / Datadog. Гистограммы и DDSketch для латентностей.
- Сетевое оборудование. Sketch-based телеметрия (например, семейство
UnivMon,SketchVisor) — на линейной скорости 100 Гбит/с ничего, кроме скетчей, не успевает. - STXXL и TPIE — библиотеки C++ для алгоритмов внешней памяти: контейнеры и сортировки, работающие поверх нескольких дисков.
Мини-итог
- Когда данных больше памяти, меняется не константа, а модель стоимости. Оптимизируйте то, что дорого: число проходов по потоку или число блоков ввода-вывода.
- Точные ответы в потоке доказуемо требуют линейной памяти (аргумент через коммуникационную
сложность). Приближённость — не компромисс качества, а осознанная покупка памяти за
контролируемую ошибку
(ε, δ). - Скетчи: Bloom (принадлежность), HyperLogLog (мощность), Count–Min и Count Sketch (частоты), Misra–Gries и Space-Saving (тяжёлые элементы), KLL, t-digest и DDSketch (квантили), reservoir sampling (выборка). Все они мержабельны — вот почему они победили в распределённых системах.
- Внешняя память:
scan(N) = Θ(N/B),sort(N) = Θ((N/B) log_{M/B}(N/B)), поискΘ(log_B N). Основание логарифма — это фан-ин, и на реальном железе он порядка тысяч, так что сортировка почти всегда двухфазная. - B-дерево против LSM — это выбор компромисса между write-, read- и space-амплификацией, а не «что лучше».
- Cache-oblivious: рекурсивное деление даёт оптимальность на всей иерархии памяти без знания её параметров.
Источники
Потоковые алгоритмы:
- S. Muthukrishnan. Data Streams: Algorithms and Applications — каноничный обзор моделей.
- N. Alon, Y. Matias, M. Szegedy. The Space Complexity of Approximating the Frequency Moments, 1996 — статья, с которой началась область (премия Гёделя).
- G. Cormode, S. Muthukrishnan. An Improved Data Stream Summary: The Count-Min Sketch and its Applications, 2005.
- P. Flajolet, É. Fusy, O. Gandouet, F. Meunier. HyperLogLog: the analysis of a near-optimal cardinality estimation algorithm, 2007.
- S. Heule, M. Nunkesser, A. Hall. HyperLogLog in Practice, EDBT 2013 — инженерные поправки Google.
- M. Greenwald, S. Khanna. Space-Efficient Online Computation of Quantile Summaries, 2001.
- Z. Karnin, K. Lang, E. Liberty. Optimal Quantile Approximation in Streams, 2016 — KLL.
- T. Dunning. The t-digest: Efficient estimates of distributions.
- M. Datar, A. Gionis, P. Indyk, R. Motwani. Maintaining Stream Statistics over Sliding Windows, 2002 — DGIM.
- J. Leskovec, A. Rajaraman, J. Ullman. Mining of Massive Datasets, глава 4 — лучшее учебное изложение.
- Apache DataSketches — документация с разбором точности каждой структуры.
Внешняя память и кэш:
- A. Aggarwal, J. Vitter. The Input/Output Complexity of Sorting and Related Problems, CACM 1988 — определение модели и нижние границы.
- J. Vitter. External Memory Algorithms and Data Structures — большой обзор.
- M. Frigo, C. Leiserson, H. Prokop, S. Ramachandran. Cache-Oblivious Algorithms, FOCS 1999.
- P. O’Neil, E. Cheng, D. Gawlick, E. O’Neil. The Log-Structured Merge-Tree, 1996.
- M. Bender et al. An Introduction to Bε-trees and Write-Optimization, ;login: 2015.
- D. Knuth. The Art of Computer Programming, том 3, раздел 5.4 — внешняя сортировка во всех подробностях, включая replacement selection.
- RocksDB Tuning Guide и Compaction — как это выглядит в реальном движке.
- Документация: ClickHouse — uniqCombined, Redis — PFADD, PostgreSQL — work_mem.
Смежные статьи трека: базовые оценки и доказательства — https://courses.digitable.life/post/algorithms/01-analysis-and-proofs/; сортировки, из которых вырастает внешнее слияние — https://courses.digitable.life/post/algorithms/02-sorting/; кэш-эффективные раскладки для поиска — https://courses.digitable.life/post/algorithms/03-searching-and-binary-search/; скользящее окно как приём — https://courses.digitable.life/post/algorithms/04-two-pointers-sliding-window/; рекурсия и разделяй-и-властвуй как источник cache-oblivious схем — https://courses.digitable.life/post/algorithms/05-recursion-and-divide-conquer/; обходы графов, которые ломаются на диске — https://courses.digitable.life/post/algorithms/08-graph-traversal/; хеширование строк — https://courses.digitable.life/post/algorithms/11-string-algorithms/; вероятностный анализ гарантий скетчей — https://courses.digitable.life/post/algorithms/14-randomized-algorithms/; распределённые вычисления и слияние состояний — https://courses.digitable.life/post/algorithms/16-parallel-and-distributed/.
Что дальше
Мы весь текст считали блоки ввода-вывода и делали вид, что вычисления внутри памяти бесплатны. Пора снять это допущение и спуститься на уровень ниже: кэш-линии и промахи, предсказание ветвлений, векторизация, профилирование и то, почему «оптимальный по асимптотике» код иногда проигрывает наивному в десять раз: Практическая оптимизация: кэш, ветвления, профилирование, SIMD.