Deep Engineering

ЗАМЕР

bench/stretched-cache/link.py

Скрипт, которым получены числа в статье, и запись прогона. Файл читается на сборке из репозитория — это тот самый код, который запускали, а не его копия.

Цитируется в статье
/ru/system-design/caching/stretched-cache
Прогон
Redis 7.0.15, PostgreSQL 16.13, Python 3.11.15, redis-py 8.1.0
Как запустить
python3 bench/stretched-cache/replication.py   # блоки 1-3
python3 bench/stretched-cache/cluster.py       # блоки 4-6
python3 bench/stretched-cache/latency.py       # блок 7
python3 bench/stretched-cache/readsdown.py     # блок 8
python3 bench/stretched-cache/offsets.py       # блок 9
python3 bench/stretched-cache/waitcost.py      # блок 10
python3 bench/stretched-cache/relaycheck.py    # блоки 11-14
python3 bench/stretched-cache/versions.py      # блок 15

Запись прогона

Замеры для статьи «Растянутый кэш на два датацентра»

Скрипт Что делает
link.py канал между датацентрами: ретранслятор в пользовательском коде, который задерживает байты и умеет обрываться. Два варианта: конвейерный DelayedLink и сериализующий SerializingLink — разница между ними измерена в блоке 13. Общий модуль, сам ничего не печатает
replication.py мастер в одном ДЦ, реплика в другом: отставание, потеря подтверждённых записей при переключении, цена WAIT — блоки 1–3
cluster.py кластер из шести узлов, три раскладки по датацентрам, разрыв канала — блоки 4–6
latency.py кэш в чужом ДЦ против базы в своём — блок 7
readsdown.py что решает cluster-allow-reads-when-down — блок 8
offsets.py чем измеряется окно потерь: отставание против круга — блок 9
waitcost.py из чего складывается цена WAIT, на трёх расстояниях — блок 10
relaycheck.py проверка самого стенда: во что он обходится при нулевой задержке, сколько раз задержка применяется к операции, сериализующий канал против конвейерного, калибровка — блоки 11–14
versions.py те же утверждения на Redis 7.0.15, Redis 8.10.1 и Valkey 9.1.2 — блок 15
python3 bench/stretched-cache/replication.py   # блоки 1-3
python3 bench/stretched-cache/cluster.py       # блоки 4-6
python3 bench/stretched-cache/latency.py       # блок 7
python3 bench/stretched-cache/readsdown.py     # блок 8
python3 bench/stretched-cache/offsets.py       # блок 9
python3 bench/stretched-cache/waitcost.py      # блок 10
python3 bench/stretched-cache/relaycheck.py    # блоки 11-14
python3 bench/stretched-cache/versions.py      # блок 15

Нужен redis-server в PATH, для кластера ещё и redis-cli, для latency.py — работающий PostgreSQL (DSN в DE_TEST_DATABASE_URL). Другую сборку можно подсунуть через REDIS_SERVER_BIN и REDIS_CLI_BIN; versions.py берёт список из DE_REDIS_BUILDS вида «имя=каталог/bin» через запятую. Узлы поднимаются во временных каталогах на случайных портах и убиваются в finally; после себя скрипты не оставляют ни процессов, ни файлов.

Чем изображается разрыв, и почему двумя разными способами

В replication.py, latency.py, offsets.py, waitcost.py и relaycheck.py — ретранслятором link.py. Он слушает порт, соединяется с настоящим Redis и переносит байты, задерживая каждую порцию на заданное время. tc netem был бы честнее, но в ядре этой среды его попросту нет: CONFIG_NET_SCH_NETEM не собран, и tc qdisc add ... netem отвечает Error: Specified qdisc kind is unknown. Права тут ни при чём.

И у этого ретранслятора нашлось два дефекта — оба измерены и устранены. Первый: на его двух соединениях не был выставлен TCP_NODELAY, и пара «алгоритм Нейгла — отложенное подтверждение» добавляла сорок миллисекунд к каждой операции, где по каналу шло несколько мелких сообщений подряд. Ровно эти сорок миллисекунд первая версия статьи опубликовала как «слагаемое неизвестной природы» в цене WAIT. Второй: задержка ставилась между recv и sendall в одном потоке, поэтому порции шли по очереди, а не параллельно, и операция стоила на один лишний перелёт дороже. Пока первый дефект был на месте, второй был невидим — он вчетверо меньше. Подробности печатает relaycheck.py.

У ретранслятора есть cut(): он закрывает все соединения и перестаёт принимать новые. Это важнее задержки. Разрыв между датацентрами — это не «стало медленно», а «перестало ходить вовсе», и изображать его увеличенной задержкой нельзя: репликация с задержкой в десять секунд всё равно догонит, а оборванная — нет.

В cluster.py — сигналом STOP дальней половине. Здесь ретранслятор не годится: узлы кластера общаются друг с другом по шине напрямую, а не через клиентский порт, и ставить ретранслятор пришлось бы на каждую пару. Процесс, остановленный SIGSTOP, жив и держит порт открытым, но не отвечает и не шлёт heartbeat — для оставшейся половины он неотличим от узла за оборванным каналом. Убить узлы было нельзя: вторую половину надо было опросить тоже, и а каждая половина наблюдается на СВОЁМ свежем кластере с той же, заданной руками топологией. Это исправление: раньше обе половины проверялись по очереди на одном кластере, и повышение реплики в первой проверке могло изменить топологию для второй. Две последовательные проверки — не две стороны одного разрыва.

Что здесь важно прочитать правильно

Абсолютные миллисекунды не переносятся никуда — ни на другую машину, ни на другую версию Redis. Содержательны отношение внутри блока, знак разницы и то, какая половина кластера ответила, а какая нет. Двадцать миллисекунд в одну сторону выбраны не «чтобы было хуже», а как обычное расстояние между датацентрами в пределах одной страны.

Номера портов в блоках 4–6 случайны и в каждом прогоне свои. Читаются не они, а роли узлов и столбец состояния. Именно поэтому в блоке 5 стоит проследить пальцем, чья реплика где: в этом весь результат.

Блоки 5 и 6 различаются одной перестановкой реплики. Деление узлов по датацентрам в них одинаковое; меняется только то, чью реплику оставили дома. Разница в исходе — между «отказали обе половины» и «одна работает».

Замер задержки сравнивает самую дешёвую операцию базы. Пять тысяч строк, всё в памяти, чтение по первичному ключу. Из того, что первые две строки блока 7 почти совпали, не следует, что кэш не нужен, — следует только правило: поход в кэш обязан быть дешевле той операции источника, которую он заменяет, на этом расстоянии. Для дорогого агрегата ответ может быть и другим, и здесь он не мерился.

Опубликованная основа — Redis 7.0.15, но она не единственная проверка. Блок 15 (versions.py) прогоняет те же утверждения на Redis 8.10.1 и Valkey 9.1.2: ни одно не разошлось. Числа блоков 1–10 при этом на новые версии не переносятся и не заменяются ими — вопрос блока 15 другой: не разъехалось ли поведение.

Что получилось (Redis 7.0.15, PostgreSQL 16.13, Python 3.11.15, redis-py 8.1.0)

Мастер и реплика, канал 20 мс в одну сторону:

что значение
реплика в том же ДЦ: запись видна через 1,1 мс
реплика в другом ДЦ: запись видна через 20,1 мс

В этом опыте при низкой нагрузке дополнительная задержка видимости почти совпала с добавленным перелётом в одну сторону. Совпадение не есть закон: обработка, буферизация и нагрузка добавляют своё.

Переключение на реплику, 200 подтверждённых клиенту записей:

когда оборвали оказалось на реплике потеряно
сразу после записи 107 93 (46,5 %)
через секунду после записи 200 0

Обе строки — один и тот же опыт с единственной разницей: успела ли репликация догнать до обрыва. Проценты — не свойство Redis; переносится граница: под угрозой то, что не доехало до реплики. Чем эта граница измеряется, см. ниже — здесь первая версия ошибалась.

Цена WAIT 1:

что значение
обычная запись 0,24 мс
запись с WAIT 1 41,73 мс
во сколько раз дороже ×171,6

Это ровно один круг до реплики — разбор на трёх расстояниях ниже. Купленное свойство в любом случае слабее ожидаемого: WAIT сообщает число подтвердивших реплик, но не отменяет запись, если их меньше.

Кластер из шести узлов, разрыв канала между датацентрами:

раскладка ДЦ-1 ДЦ-2
A: все мастера в ДЦ-1, все реплики в ДЦ-2 работает CLUSTERDOWN
B: узлы поровну, мастера 2 и 1 CLUSTERDOWN CLUSTERDOWN
C: то же деление, реплика чужого мастера дома работает CLUSTERDOWN

Главный результат — строка B. У ДЦ-1 большинство мастеров, и голосов на повышение реплики ему хватает, а кэша всё равно нет: реплика мастера, оставшегося в ДЦ-2, тоже осталась в ДЦ-2, и слоты этого мастера не покрыты никем.

Задержка, медиана из 200 замеров:

откуда читаем медиана p95
кэш в своём ДЦ 0,10 мс 0,14 мс
база в своём ДЦ 0,09 мс 0,13 мс
кэш в чужом ДЦ 41,09 мс 41,36 мс

Выводы, которые пришлось исправить

Разборов снаружи было два. Первый указал на три вывода, сделанных из одной точки каждый; второй — на сам измерительный стенд. Все указанные места пришлось не переписать словами, а перемерить, и четыре вывода из пяти оказались неверны.

«Наблюдаемое поведение шире документации» — неверно. Было сказано: раз cluster-require-full-coverage описан через записи, а GET тоже получил CLUSTERDOWN, значит поведение шире написанного. На самом деле за чтения отвечает отдельная описанная настройка cluster-allow-reads-when-down, по умолчанию no. readsdown.py меряет обе: с no двадцать ключей из двадцати получают CLUSTERDOWN; с yes узел отдаёт только ключи своих слотов — четырнадцать из двадцати, — остальные перенаправляет за оборванный канал, а записи отказаны при обоих значениях.

«Под угрозой всё, что записано за последний круг» — неверно. Круг совпал с потерей случайно. offsets.py меняет один только темп записи на том же канале и получает 104, 3 и 0 потерянных записей при отставании 3548, 99 и 0 байт. Мерой служит разница смещений потока репликации, а не круг.

«82 мс — это два круга» — неверно, и заменившее его объяснение тоже. waitcost.py показал, что фиксированным числом кругов цена не описывается, и статья честно написала, что сверх круга остаётся слагаемое в 42–44 мс неизвестной природы. Природа оказалась в приборе: relaycheck.py включает и выключает TCP_NODELAY на ретрансляторе при НУЛЕВОЙ задержке и получает 44,00 и 0,70 мс. После исправления цена WAIT — ровно один круг: 10,90, 41,52 и 81,40 мс при отношении 1,09, 1,04 и 1,02.

«Две половины одного разрыва» — так называть две последовательные проверки было нельзя. После первой проверки топология не обязана вернуться к исходной: реплика могла повыситься, слоты переехать. cluster.py теперь поднимает свежий кластер на каждую сторону и печатает топологию до и после, а также смену ролей.

«Большинство плюс реплика каждого потерянного мастера» — условие неполное. Третья часть — свежесть реплики: отставшая сверх cluster-replica-validity-factor выборы не начнёт. В этих опытах она не мешала, и это теперь напечатано, а не предполагается.

Источники

Скрипт

240 строк
"""Канал между датацентрами: ретранслятор с задержкой и рубильником.

ПОЧЕМУ РЕТРАНСЛЯТОР, А НЕ `tc netem`. Настоящий эмулятор сети ставится
дисциплиной очереди `netem` на сетевом интерфейсе. Она требует прав и, что
важнее здесь, поддержки в ядре: `CONFIG_NET_SCH_NETEM`. В песочнице, где
снимались эти числа, ядро собрано без неё — `tc qdisc add ... netem` отвечает
`Error: Specified qdisc kind is unknown`, и никакими правами это не лечится.
Ретранслятор в пользовательском коде даёт наблюдаемую величину того же рода
(известная задержка между двумя портами) и не трогает ничего за пределами
процесса.

ЧТО ЭТО ЗНАЧИТ ДЛЯ ВЫВОДОВ. Ретранслятор — приближение задержки, а не
эмулятор на уровне пакетов, и разница между ними измерима. Она измерена:
`relaycheck.py` считает, во что обходится сам стенд при нулевой задержке,
сколько раз задержка применяется к одной операции и что меняет замена
сериализующего ретранслятора на конвейерный. Читать выводы, снятые через этот
канал, стоит вместе с блоками 11-14.

ДВА РЕТРАНСЛЯТОРА, И РАЗНИЦА МЕЖДУ НИМИ СОДЕРЖАТЕЛЬНА.

`SerializingLink` — первая версия стенда. Она спит между `recv` и `sendall` в
одном потоке:

    data = src.recv(65536); time.sleep(delay); dst.sendall(data)

Пока поток спит, он не читает следующую порцию. Поэтому k порций, пришедших
подряд, стоят k задержек ПОДРЯД, а не одну: канал не только удлиняется, но и
перестаёт пропускать больше одной порции за время полёта. Для протокола
«запрос — ответ» разницы почти нет; для потока репликации и для `WAIT`,
где по каналу за одну операцию проходит несколько порций, разница есть — и
блок 13 её печатает.

`DelayedLink` — то, чем стенд стал после разбора. Чтение и отправка разнесены:
читающий поток кладёт порцию в очередь со сроком доставки `сейчас + delay` и
сразу читает дальше, отправляющий поток ждёт срок и отправляет. Порции,
пришедшие одновременно, уходят одновременно — это и есть постоянная задержка
распространения. Границ сообщений TCP по-прежнему не хранит, так что и это
приближение; но приближение, у которого задержка не зависит от того, на
сколько порций разбился поток.

ЧТО ЗДЕСЬ ЕЩЁ ЕСТЬ — рубильник `cut()`. Разрыв между датацентрами это не
«стало медленно», это «перестало ходить вовсе», и изображать его увеличенной
задержкой нельзя: узел, ждущий ответа, и узел, получивший сброс соединения,
ведут себя по-разному. `cut()` закрывает все установленные соединения и
перестаёт принимать новые — с точки зрения обеих сторон канал оборван.

СЧЁТЧИКИ. Оба ретранслятора считают порции и байты в каждую сторону
(`stats()`). Это не отладка, а измерительный прибор: число порций и есть то,
сколько раз была применена задержка.
"""

from __future__ import annotations

import queue
import socket
import threading
import time
from dataclasses import dataclass, field


@dataclass
class LinkStats:
    """Сколько раз задержка была применена и к скольким байтам."""

    chunks_up: int = 0
    chunks_down: int = 0
    bytes_up: int = 0
    bytes_down: int = 0
    _lock: threading.Lock = field(default_factory=threading.Lock, repr=False)

    def add(self, upstream: bool, size: int) -> None:
        with self._lock:
            if upstream:
                self.chunks_up += 1
                self.bytes_up += size
            else:
                self.chunks_down += 1
                self.bytes_down += size

    @property
    def chunks(self) -> int:
        return self.chunks_up + self.chunks_down


class _BaseLink:
    """Общая часть: слушающий порт, приём соединений, рубильник."""

    def __init__(self, upstream_port: int, delay_ms: float, nodelay: bool = True) -> None:
        self.upstream_port = upstream_port
        self.delay_ms = delay_ms
        # nodelay=False воспроизводит прежний стенд. Оставлено не для
        # совместимости, а затем, чтобы блок 11 мог включить и выключить ровно
        # одну вещь и показать, чего она стоила.
        self.nodelay = nodelay
        self.alive = True
        self.stats = LinkStats()
        self._conns: list[socket.socket] = []
        self._lock = threading.Lock()
        self.sock = socket.socket()
        self.sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
        self.sock.bind(("127.0.0.1", 0))
        self.sock.listen(128)
        self.port = self.sock.getsockname()[1]
        threading.Thread(target=self._accept_loop, daemon=True).start()

    def reset_stats(self) -> None:
        self.stats = LinkStats()

    def cut(self) -> None:
        """Оборвать канал: закрыть установленное и не принимать новое."""
        self.alive = False
        with self._lock:
            conns, self._conns = self._conns, []
        for c in conns:
            try:
                c.close()
            except OSError:
                pass

    def _accept_loop(self) -> None:
        while True:
            try:
                client, _ = self.sock.accept()
            except OSError:
                return
            if not self.alive:
                client.close()
                continue
            threading.Thread(target=self._serve, args=(client,), daemon=True).start()

    def _serve(self, client: socket.socket) -> None:
        try:
            server = socket.create_connection(("127.0.0.1", self.upstream_port))
        except OSError:
            client.close()
            return
        # TCP_NODELAY НА ОБА СОКЕТА, И ЭТО НЕ МИКРООПТИМИЗАЦИЯ.
        #
        # Redis ставит этот флаг на своих соединениях сам. Ретранслятор создаёт
        # ДВА НОВЫХ соединения, и на них флага не было. Отсюда пара «алгоритм
        # Нейгла на отправителе — отложенное подтверждение на получателе»:
        # маленькая порция ждёт подтверждения предыдущей, а подтверждение
        # придерживается до 40 мс. Ровно эти 40 мс и оседали в каждом замере,
        # где по каналу шло несколько мелких сообщений подряд, — и именно они
        # стояли в блоке 10 как «слагаемое, от расстояния не зависящее, чем
        # вызвано — замер не устанавливает». Установил блок 11.
        if self.nodelay:
            for s in (client, server):
                s.setsockopt(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1)
        with self._lock:
            self._conns += [client, server]
        # upstream=True — порция идёт от клиента к серверу за каналом.
        self._start_pipe(client, server, upstream=True)
        self._start_pipe(server, client, upstream=False)

    def _start_pipe(self, src: socket.socket, dst: socket.socket, upstream: bool) -> None:
        raise NotImplementedError

    @staticmethod
    def _shutdown(*socks: socket.socket) -> None:
        for s in socks:
            try:
                s.close()
            except OSError:
                pass


class SerializingLink(_BaseLink):
    """Первая версия стенда: сон между чтением и отправкой, в одном потоке.

    Оставлена НЕ для совместимости, а как предмет измерения: блок 13 сравнивает
    её с конвейерной версией на одних и тех же операциях, и без неё сравнивать
    было бы не с чем.
    """

    def _start_pipe(self, src: socket.socket, dst: socket.socket, upstream: bool) -> None:
        threading.Thread(target=self._pump, args=(src, dst, upstream), daemon=True).start()

    def _pump(self, src: socket.socket, dst: socket.socket, upstream: bool) -> None:
        try:
            while True:
                data = src.recv(65536)
                if not data:
                    break
                self.stats.add(upstream, len(data))
                # Задержка ДО отправки и В ТОМ ЖЕ потоке: пока спим, следующая
                # порция не читается. Отсюда и вся разница с конвейерной версией.
                time.sleep(self.delay_ms / 1000)
                dst.sendall(data)
        except OSError:
            pass
        finally:
            self._shutdown(src, dst)


class DelayedLink(_BaseLink):
    """Конвейерная версия: чтение не останавливается на время полёта.

    Читающий поток проставляет порции срок доставки и сразу читает дальше;
    отправляющий поток ждёт срок и отправляет. Порции, пришедшие в один момент,
    уходят в один момент, позже на `delay_ms`, — это и есть постоянная задержка
    распространения, а не задержка на порцию.
    """

    _EOF = object()

    def _start_pipe(self, src: socket.socket, dst: socket.socket, upstream: bool) -> None:
        q: queue.Queue = queue.Queue()
        threading.Thread(target=self._reader, args=(src, q, upstream), daemon=True).start()
        threading.Thread(target=self._sender, args=(q, src, dst), daemon=True).start()

    def _reader(self, src: socket.socket, q: queue.Queue, upstream: bool) -> None:
        delay = self.delay_ms / 1000
        try:
            while True:
                data = src.recv(65536)
                if not data:
                    break
                self.stats.add(upstream, len(data))
                q.put((time.monotonic() + delay, data))
        except OSError:
            pass
        finally:
            q.put((time.monotonic() + delay, self._EOF))

    def _sender(self, q: queue.Queue, src: socket.socket, dst: socket.socket) -> None:
        try:
            while True:
                due, data = q.get()
                left = due - time.monotonic()
                if left > 0:
                    time.sleep(left)
                if data is self._EOF:
                    break
                dst.sendall(data)
        except OSError:
            pass
        finally:
            self._shutdown(src, dst)