Deep Engineering

ЗАМЕР

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

Источники

Скрипт

284 строк
"""Почему отказавший кластер не отвечает и на чтение — и что это меняет.

ЗАЧЕМ ЭТОТ ЗАМЕР ПОЯВИЛСЯ. В первой версии статьи стояло утверждение:
документация `cluster-require-full-coverage` говорит про ЗАПИСИ, а замер
показал `CLUSTERDOWN` и на обычный `GET`, — значит, «наблюдаемое поведение
шире написанного». Утверждение оказалось неверным: у Redis есть отдельная
настройка `cluster-allow-reads-when-down`, по умолчанию `no`, и именно она
отвечает за чтения на узле отказавшего кластера. То есть поведение не шире
документации — оно описано, просто в другом месте.

Переписывать вывод словами было бы недостаточно: получилось бы одно
утверждение без замера вместо другого утверждения без замера. Поэтому здесь
тот же расклад, что в блоке 5, прогоняется ДВАЖДЫ — с настройкой `no` и с
настройкой `yes`, — и разница видна на одном и том же `GET`.

ЧТО ЭТО НЕ ОТМЕНЯЕТ. Главный результат статьи — «отказали обе половины» —
держится не на этой настройке, а на кворуме и покрытии слотов. Настройка
меняет только то, что видит клиент на чтении в отказавшей половине. Замер
ниже показывает и это.

ЗАПУСК: python3 bench/stretched-cache/readsdown.py
Нужны redis-server и redis-cli в PATH. Занимает около минуты.
"""

from __future__ import annotations

import os
import socket
import subprocess
import tempfile
import time

import redis

# Какой сервер запускать. По умолчанию тот, что в 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")

NODE_TIMEOUT_MS = 2000
SETTLE_S = 8.0

_PIDS: dict[int, int] = {}


def free_port() -> int:
    """Свободный порт ниже 55535: шинный порт узла на 10000 больше."""
    for _ in range(500):
        s = socket.socket()
        s.bind(("127.0.0.1", 0))
        port = s.getsockname()[1]
        s.close()
        if port < 55535:
            return port
    raise RuntimeError("не нашлось свободного порта ниже 55535")


def start_node(port: int, dirname: str, allow_reads: str) -> subprocess.Popen:
    proc = subprocess.Popen(
        [
            SERVER,
            "--port", str(port),
            "--dir", dirname,
            "--cluster-enabled", "yes",
            "--cluster-config-file", f"nodes-{port}.conf",
            "--cluster-node-timeout", str(NODE_TIMEOUT_MS),
            # Единственное, что различается между двумя прогонами.
            "--cluster-allow-reads-when-down", allow_reads,
            "--save", "",
            "--appendonly", "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()
                _PIDS[port] = proc.pid
                return proc
        except Exception:
            time.sleep(0.05)
    raise RuntimeError(f"узел {port} не поднялся")


def state_of(port: int) -> str:
    try:
        c = redis.Redis(port=port, socket_timeout=1.0, decode_responses=True)
        return c.cluster("INFO")["cluster_state"]
    except Exception:
        return "недоступен"


def probe(port: int, key: str, write: bool) -> str:
    """Что получает клиент: значение, перенаправление, отказ или молчание."""
    try:
        c = redis.Redis(port=port, socket_timeout=1.5, decode_responses=True)
        if write:
            c.set(key, "v")
            return "записал"
        value = c.get(key)
        return "ответил" if value is not None else "ответил (пусто)"
    except redis.exceptions.ClusterDownError:
        return "CLUSTERDOWN"
    except (redis.exceptions.MovedError, redis.exceptions.AskError):
        return "перенаправил"
    except redis.exceptions.ResponseError as exc:
        return str(exc).split()[0]
    except Exception:
        return "таймаут"


# Ключи раскладываются по слотам хешем, поэтому одного ключа мало: он
# случайно попадёт на какой-то один узел, и по нему нельзя сказать, что
# происходит с остальным пространством ключей. Двадцать ключей покрывают
# все три слотовых диапазона.
KEYS = [f"probe:{i}" for i in range(20)]


def probe_all(port: int) -> dict[str, int]:
    """Сколько ключей из KEYS узел отдал, перенаправил и по скольким отказал."""
    counts: dict[str, int] = {}
    for k in KEYS:
        answer = probe(port, k, write=False)
        counts[answer] = counts.get(answer, 0) + 1
    return counts


def topology(ports: list[int]) -> dict[int, str]:
    c = redis.Redis(port=ports[0], socket_timeout=2.0, decode_responses=True)
    raw = c.execute_command("CLUSTER", "NODES")
    by_id: dict[str, int] = {}
    rows = []
    for line in raw.strip().splitlines():
        f = line.split()
        node_id, addr, flags, master_id = f[0], f[1], f[2], f[3]
        port = int(addr.split("@")[0].split(":")[1])
        by_id[node_id] = port
        rows.append((port, flags, master_id))
    return {
        port: ("мастер" if "master" in flags else f"реплика {by_id.get(master_id, '?')}")
        for port, flags, master_id in rows
    }


def roles(ports: list[int]) -> dict[int, str]:
    out = {}
    for p in ports:
        try:
            c = redis.Redis(port=p, socket_timeout=1.0, decode_responses=True)
            out[p] = c.info("replication")["role"]
        except Exception:
            out[p] = "?"
    return out


def run_once(allow_reads: str) -> list[tuple[int, str, str, str]]:
    """Поднять кластер, разложить по расклад B, оборвать канал, опросить ДЦ-1.

    Возвращает строки (порт, роль, состояние, ответ на GET) для той половины,
    у которой большинство мастеров, — той самой, про которую и шёл спор.
    """
    tmp = tempfile.mkdtemp(prefix=f"readsdown-{allow_reads}-")
    ports = [free_port() for _ in range(6)]
    procs = []
    try:
        for p in ports:
            d = os.path.join(tmp, str(p))
            os.makedirs(d, exist_ok=True)
            procs.append(start_node(p, d, allow_reads))

        args = [CLI, "--cluster", "create"]
        args += [f"127.0.0.1:{p}" for p in ports]
        args += ["--cluster-replicas", "1", "--cluster-yes"]
        subprocess.run(args, check=True, capture_output=True, timeout=60)
        time.sleep(3)

        cc = redis.RedisCluster(host="127.0.0.1", port=ports[0], decode_responses=True)
        for k in KEYS:
            cc.set(k, "value")
        cc.close()

        topo = topology(ports)
        masters = [p for p, v in topo.items() if v == "мастер"]
        replicas = [p for p in ports if p not in masters]

        # Расклад B из блока 5: узлы поровну, мастера поделены 2 и 1, и
        # реплика третьего мастера уезжает вместе с ним.
        dc1 = [masters[0], masters[1], replicas[2]]
        dc2 = [masters[2], replicas[0], replicas[1]]

        for p in dc2:
            subprocess.run(["kill", "-STOP", str(_PIDS[p])], check=False)
        time.sleep(SETTLE_S)

        r = roles(dc1)
        rows = []
        for p in dc1:
            counts = probe_all(p)
            summary = ", ".join(f"{v} {k}" for k, v in sorted(counts.items()))
            rows.append((p, r.get(p, "?"), state_of(p), summary))
        # Заодно спросим запись — чтобы было видно, что настройка про чтения
        # и только про них.
        write_answer = probe(dc1[0], KEYS[0], write=True)
        rows.append((0, "запись", "", write_answer))

        for p in dc2:
            subprocess.run(["kill", "-CONT", str(_PIDS[p])], check=False)
        return rows
    finally:
        for proc in procs:
            proc.terminate()
            try:
                proc.wait(timeout=5)
            except subprocess.TimeoutExpired:
                proc.kill()


def main() -> None:
    print("Redis · тот же расклад, что в блоке 5: узлы поровну, мастера 2 и 1")
    print(f"cluster-node-timeout {NODE_TIMEOUT_MS} мс; после разрыва {SETTLE_S:.0f} с на выборы")
    print("опрашивается половина С БОЛЬШИНСТВОМ мастеров — та, про которую спор")
    print()

    print("8. ЧТЕНИЯ В ОТКАЗАВШЕМ КЛАСТЕРЕ РЕШАЕТ ОТДЕЛЬНАЯ НАСТРОЙКА")
    print("-" * 66)
    served: dict[str, int] = {}
    for allow in ("no", "yes"):
        rows = run_once(allow)
        # Сколько ключей из двадцати выжившая половина реально отдала.
        # Считается по выводу, а не пишется руками: какие именно мастера
        # окажутся в ДЦ-1, зависит от того, как кластер расставил роли, и
        # число от прогона к прогону разное.
        served[allow] = sum(
            int(part.split()[0])
            for _, role, _, answer in rows
            if role != "запись"
            for part in answer.split(", ")
            if part.endswith("ответил")
        )
        print(f"  cluster-allow-reads-when-down {allow}"
              f"{'  (по умолчанию)' if allow == 'no' else ''}:")
        for port, role, state, answer in rows:
            if role == "запись":
                print(f"    SET на том же узле                        -> {answer}")
            else:
                print(f"    узел {port}: {role:<8} состояние {state:<8} 20 ключей -> {answer}")
        print()

    print("  Обе строки сняты на одном и том же раскладе и одном и том же")
    print("  разрыве. Различается одна настройка.")
    print()
    print("  Отсюда поправка к первой версии статьи. Утверждение «наблюдаемое")
    print("  поведение шире документации» было неверным: отказ на чтении не")
    print("  побочный эффект cluster-require-full-coverage, а отдельно")
    print("  описанное поведение cluster-allow-reads-when-down со значением")
    print("  no по умолчанию. Цепочка такая: неполное покрытие слотов ->")
    print("  кластер помечен отказавшим -> узел отказавшего кластера не")
    print("  обслуживает чтения, пока не сказано обратное.")
    print()
    print("  И чего настройка НЕ меняет — это важнее самой поправки.")
    print()
    print("  Со значением yes узел отвечает только по ТЕМ ключам, чьи слоты")
    print("  остались за ним: остальные он перенаправляет на мастер, который")
    print("  сейчас за оборванным каналом, и туда клиент не доедет.")
    print()
    print(f"  Ключей отдано из двадцати: с no — {served['no']}, с yes — "
          f"{served['yes']}.")
    print("  Второе число — не двадцать и от прогона к прогону разное: оно")
    print("  зависит от того, сколько слотов досталось двум оставшимся")
    print("  мастерам. Треть пространства ключей — слоты мастера за обрывом —")
    print("  недоступна при любой настройке.")
    print()
    print("  Записи отказаны в обоих случаях. То есть yes не возвращает кэш:")
    print("  он превращает «не работает ничего» в «работает часть чтений на")
    print("  части пространства ключей». Главный вывод статьи — при ровном")
    print("  делении кэша не остаётся ни в одном ДЦ — от настройки не зависит.")


if __name__ == "__main__":
    main()