Deep Engineering

ЗАМЕР

bench/stretched-cache/replication.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 выборы не начнёт. В этих опытах она не мешала, и это теперь напечатано, а не предполагается.

Источники

Скрипт

250 строк
"""Мастер в одном датацентре, реплика в другом: что теряется при переключении.

ЧТО ЗДЕСЬ МЕРЯЕТСЯ. Не «медленно ли работает репликация» — это и так
очевидно, — а сколько ПОДТВЕРЖДЁННЫХ клиенту записей исчезает, когда мастер
пропал и реплику подняли на его место. Величина важна тем, что её обычно
считают нулевой: клиент получил OK, значит записалось.

ПОЧЕМУ ЭТО НЕ ПРИДИРКА К REDIS. Асинхронная репликация — не дефект
реализации, а объявленный договор: мастер отвечает клиенту, не дожидаясь
реплик. Документация Redis говорит об этом прямо, и цитата стоит в статье.
Замер лишь переводит договор в число для конкретного расстояния.

ЗАПУСК: python3 bench/stretched-cache/replication.py
Нужен redis-server в PATH. Ничего не остаётся после выхода: оба процесса
поднимаются во временных каталогах и убиваются в finally.
"""

from __future__ import annotations

import os
import shutil
import socket
import subprocess
import sys
import tempfile
import time

sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
from link import DelayedLink  # noqa: E402

import redis  # noqa: E402

# Какой сервер запускать. По умолчанию тот, что в PATH; REDIS_SERVER_BIN и
# REDIS_CLI_BIN позволяют прогнать тот же замер на другой сборке — так снята
# таблица регрессии на Redis 8.10.1 и smoke-прогон на Valkey.
SERVER = os.environ.get("REDIS_SERVER_BIN", "redis-server")
CLI = os.environ.get("REDIS_CLI_BIN", "redis-cli")

# ОДНОСТОРОННЯЯ задержка. Круг — вдвое больше. Двадцать миллисекунд в одну
# сторону это типичное расстояние между датацентрами в пределах одной страны:
# Москва — Петербург около того. Величина выбрана не «чтобы было хуже», а
# как обычный случай, о котором и приходят просить.
ONE_WAY_MS = 20.0
WRITES = 200


def free_port() -> int:
    s = socket.socket()
    s.bind(("127.0.0.1", 0))
    port = s.getsockname()[1]
    s.close()
    return port


def start_redis(port: int, dirname: str) -> subprocess.Popen:
    proc = subprocess.Popen(
        [
            SERVER,
            "--port", str(port),
            "--dir", dirname,
            "--save", "",
            "--appendonly", "no",
            "--daemonize", "no",
            "--logfile", os.path.join(dirname, "redis.log"),
        ]
    )
    deadline = time.time() + 10
    while time.time() < deadline:
        try:
            c = redis.Redis(port=port, socket_timeout=0.5)
            if c.ping():
                c.close()
                return proc
        except Exception:
            time.sleep(0.05)
    raise RuntimeError(f"redis на порту {port} не поднялся")


def wait_sync(client: redis.Redis, timeout: float = 15.0) -> None:
    """Дождаться, пока реплика доложит link_status:up и первую синхронизацию."""
    deadline = time.time() + timeout
    while time.time() < deadline:
        info = client.info("replication")
        if info.get("master_link_status") == "up":
            return
        time.sleep(0.1)
    raise RuntimeError("реплика не синхронизировалась")


def main() -> None:
    tmp_a = tempfile.mkdtemp(prefix="dc-a-")
    tmp_b = tempfile.mkdtemp(prefix="dc-b-")
    port_a, port_b = free_port(), free_port()
    proc_a = proc_b = None
    try:
        proc_a = start_redis(port_a, tmp_a)
        proc_b = start_redis(port_b, tmp_b)
        a = redis.Redis(port=port_a, decode_responses=True)
        b = redis.Redis(port=port_b, decode_responses=True)

        version = a.info("server")["redis_version"]
        print(f"Redis {version} · один хост, два процесса")
        print(f"канал между ДЦ: {ONE_WAY_MS:.0f} мс в одну сторону, "
              f"то есть круг {2 * ONE_WAY_MS:.0f} мс")
        print()

        # --- Блок 1: цена расстояния для самой репликации -------------------
        near = DelayedLink(port_a, delay_ms=0.0)
        b.replicaof("127.0.0.1", near.port)
        wait_sync(b)
        a.set("probe", "1")
        t0 = time.perf_counter()
        while b.get("probe") != "1":
            time.sleep(0.0005)
        near_lag = (time.perf_counter() - t0) * 1000

        b.execute_command("REPLICAOF", "NO", "ONE")
        b.flushall()
        far = DelayedLink(port_a, delay_ms=ONE_WAY_MS)
        b.replicaof("127.0.0.1", far.port)
        wait_sync(b)
        a.set("probe2", "1")
        t0 = time.perf_counter()
        while b.get("probe2") != "1":
            time.sleep(0.0005)
        far_lag = (time.perf_counter() - t0) * 1000

        print("1. ЗАПИСЬ ДОХОДИТ ДО РЕПЛИКИ НЕ СРАЗУ, И РАССТОЯНИЕ ЭТО ВИДНО")
        print("-" * 66)
        print(f"  реплика в том же ДЦ: запись видна через, мс   {near_lag:.1f}")
        print(f"  реплика в другом ДЦ: запись видна через, мс   {far_lag:.1f}")
        print(f"  разница, мс                                   {far_lag - near_lag:.1f}")
        print()
        print("  Клиент к этому моменту давно получил OK: мастер отвечает,")
        print("  не дожидаясь реплики. Отставание само по себе не беда —")
        print("  бедой оно становится в следующем блоке.")
        print()

        # --- Блок 2: подтверждённые записи, которых не стало ----------------
        a.flushall()
        time.sleep(0.3)
        acked = 0
        for i in range(WRITES):
            a.set(f"key:{i}", i)
            acked += 1
        # Разрыв ровно в тот момент, когда все записи подтверждены клиенту.
        far.cut()
        # Реплику поднимают на место мастера: это и есть переключение.
        b.execute_command("REPLICAOF", "NO", "ONE")
        time.sleep(0.2)
        survived = sum(1 for i in range(WRITES) if b.get(f"key:{i}") is not None)
        lost = acked - survived

        # Та же связка, но разрыв не сразу после записей, а через секунду.
        # ЗАЧЕМ ВТОРАЯ ПОЛОВИНА: без неё число из первой читается как
        # «теряется почти всё», а это неверно. Теряется то, что не успело
        # улететь, и вопрос только в том, сколько записей помещается в один
        # круг. Двести записей по 0,15 мс помещаются целиком — отсюда и 99,5%.
        far.cut()
        b.execute_command("REPLICAOF", "NO", "ONE")
        b.flushall()
        far3 = DelayedLink(port_a, delay_ms=ONE_WAY_MS)
        b.replicaof("127.0.0.1", far3.port)
        wait_sync(b)
        a.flushall()
        time.sleep(0.3)
        for i in range(WRITES):
            a.set(f"calm:{i}", i)
        time.sleep(1.0)          # тишина: репликация догнала
        far3.cut()
        b.execute_command("REPLICAOF", "NO", "ONE")
        time.sleep(0.2)
        calm_survived = sum(1 for i in range(WRITES) if b.get(f"calm:{i}") is not None)

        print("2. ПЕРЕКЛЮЧЕНИЕ НА РЕПЛИКУ ТЕРЯЕТ ПОДТВЕРЖДЁННЫЕ ЗАПИСИ")
        print("-" * 66)
        print(f"  записей подтверждено клиенту                  {acked}")
        print("  разрыв сразу после записи:")
        print(f"    оказалось на реплике                        {survived}")
        print(f"    потеряно                                    {lost} ({100 * lost / acked:.1f}%)")
        print("  разрыв через секунду после записи:")
        print(f"    оказалось на реплике                        {calm_survived}")
        print(f"    потеряно                                    {acked - calm_survived}"
              f" ({100 * (acked - calm_survived) / acked:.1f}%)")
        print()
        print("  Читать надо обе строки сразу. Теряется не «почти всё» и не")
        print("  «ничего»: теряется то, что в момент разрыва ещё не доехало")
        print("  до реплики. Залп из двухсот записей уходит быстрее, чем")
        print("  летит, — отсюда первая строка. Через секунду тишины")
        print("  репликация догнала — отсюда вторая.")
        print()
        print("  ЧЕМ ИЗМЕРЯЕТСЯ ЭТО ОКНО, ЗДЕСЬ НЕ ВИДНО, и соблазн ответить")
        print("  «кругом» велик. В первой версии статьи так и было написано,")
        print("  и это неверно: круг здесь совпал с потерей случайно. Блок 9")
        print("  меняет один лишь ТЕМП записи на том же канале и получает")
        print("  другие потери — смотреть надо на разницу смещений в")
        print("  INFO replication, а не на круг.")
        print()
        print("  Ни одна из потерянных записей об ошибке не сообщила — на все")
        print("  клиент получил OK. Это не сбой репликации, а её договор:")
        print("  мастер не ждёт реплику.")
        print()

        # --- Блок 3: WAIT покупает не то, за что его принимают --------------
        far2 = DelayedLink(port_a, delay_ms=ONE_WAY_MS)
        b.replicaof("127.0.0.1", far2.port)
        wait_sync(b)
        a.flushall()
        time.sleep(0.3)

        t0 = time.perf_counter()
        for i in range(20):
            a.set(f"plain:{i}", i)
        plain_ms = (time.perf_counter() - t0) * 1000 / 20

        t0 = time.perf_counter()
        acked_by_replica = 0
        for i in range(20):
            a.set(f"waited:{i}", i)
            acked_by_replica += a.execute_command("WAIT", 1, 1000)
        wait_ms = (time.perf_counter() - t0) * 1000 / 20

        print("3. WAIT ДЕЛАЕТ ЗАПИСЬ ДОРОЖЕ, НО НЕ ДЕЛАЕТ ЕЁ ТРАНЗАКЦИЕЙ")
        print("-" * 66)
        print(f"  обычная запись, мс                            {plain_ms:.2f}")
        print(f"  запись с WAIT 1, мс                           {wait_ms:.2f}")
        print(f"  во сколько раз дороже                         {wait_ms / plain_ms:.1f}")
        print(f"  реплик подтвердило (из 20 записей по 1)       {acked_by_replica}")
        print()
        print("  Цена — круг между датацентрами на КАЖДУЮ запись. А купленное")
        print("  свойство слабее ожидаемого: WAIT сообщает, сколько реплик")
        print("  подтвердило, но не отменяет запись, если их меньше. К моменту")
        print("  ответа WAIT запись на мастере уже применена и уже видна")
        print("  читателям — откатить её нечем.")
        print()
    finally:
        for proc in (proc_b, proc_a):
            if proc is not None:
                proc.terminate()
                try:
                    proc.wait(timeout=5)
                except subprocess.TimeoutExpired:
                    proc.kill()
        shutil.rmtree(tmp_a, ignore_errors=True)
        shutil.rmtree(tmp_b, ignore_errors=True)


if __name__ == "__main__":
    main()