#57: Apache Cassandra, часть 2: как работает запись
Вторая часть разбора Apache Cassandra: Александр Пахомов и Дмитрий Константинов (коммитер Apache Cassandra) прослеживают путь одной операции записи от серверного сокета до диска. По дороге — серверный `Netty` с нативным `epoll`, flow control через TCP backpressure, аутентификация и авторизация, consistent hashing с виртуальными нодами, consistency level и кворумы, gossip-протокол и hinted handoff, а затем локальная запись: `commit log`, трёхуровневый memtable (`ConcurrentSkipListMap`/trie над `B-tree`), «last write wins» по таймстемпам, сброс в `SSTable`, компакция и, наконец, удаления через tombstones.
Главное
- Серверная сторона Cassandra на `Netty` симметрична клиентской и умеет включать нативный транспорт (`epoll` на Linux, `BoringSSL` для TLS) простой подкладкой JAR в classpath — прирост порядка 10–15% «за бесплатно».
- Flow control на входе Cassandra ограничивает не число in-flight запросов, а их суммарный размер в байтах (по умолчанию ~10% от хипа); при перегрузке чаще всего просто перестаёт читать из сокета, опираясь на TCP backpressure.
- Cassandra шардирует данные через consistent hashing: `Murmur3`-хэш ключа кладётся на кольцо, реплики находятся бинарным поиском по кольцу, а виртуальные ноды (`vnodes`) выравнивают распределение при добавлении/удалении нод.
- `consistency level` (`ONE`, `QUORUM`, `LOCAL_QUORUM`, `ALL`, `ANY`) — параметр запроса, задающий число ответивших реплик; он двигает систему по CAP-треугольнику между доступностью и консистентностью.
- Живость нод Cassandra отслеживает gossip — эпидемический протокол: нода шлёт heartbeat в seed-ноду и несколько случайных, и информация о `membership` расходится по кластеру, как инфекция.
- Запись на реплике сначала идёт в `commit log`, но по умолчанию `fsync` не синхронный (раз в ~10 секунд) — Cassandra меняет надёжность одной ноды на скорость, полагаясь на репликацию; включение синхронного `fsync` роняет производительность примерно на порядок.
- Memtable — это трёхуровневая структура «мапа мап»: `ConcurrentSkipListMap` или trie по partition-ключу, персистентное `B-tree` по clustering-ключу и `B-tree` по колонкам; слияние идёт по правилу «побеждает больший таймстемп», поэтому в Cassandra критична синхронизация часов клиентов.
- Удаление в LSM — это вставка `tombstone` (могильной плиты) с таймстемпом; выкинуть его при компакции безопасно, только доказав отсутствие «затенённых» записей — по min/max таймстемпам и ключам SSTable и через `Bloom filter`, иначе старые данные воскресают.
В выпуске
- Дмитрий Константинов — Системный архитектор и Java-разработчик, коммитер Apache Cassandra; специализируется на распределённых системах, производительности и отказоустойчивости. Регулярный спикер JPoint/Joker, на Хабре — @netudima. Habr ↗ jpoint.ru ↗
Ссылки
Расшифровка
[00:00] Александр: Здарова! Это 57-й выпуск подкаста «Тысяча фичей». Вы слушаете вторую часть про Cassandra. Ранее мы до косточек разобрали клиента, а сегодня погружаемся в сервер. Поговорим про epoll, TCP backpressure, авторизацию, consistent hashing, gossip-протокол и в конце доберёмся до самого интересного — LSM-дерева. В гостях Дима Константинов, коммитер в Apache Cassandra. Поехали!
[00:39] Александр: В предыдущем выпуске, в части первой, мы очень хорошо поговорили про то, что происходит на клиенте, на стороне драйвера. Если вдруг ты, слушатель, наткнулся на вторую часть, не послушав первую, я крайне рекомендую сначала прослушать первую часть и потом сразу переходить ко второй — как минимум чтобы быть в контексте. Итак, в предыдущей части мы остановились на том, что на клиенте сериализовали данные, засунули их в TCP-соединение и отправили куда-то на сервер. И вот вопрос: что с ними происходит, на каком сервере, что вообще дальше по сюжету? Мы сформировали какой-то набор байтов, записали в TCP-соединение, они как-то долетели до нашего сервера, попали в один из серверов Cassandra.
[01:29] Дмитрий: Мы говорили ранее, что серверы Cassandra равноправные — в принципе неважно, в какой из них попал запрос, любой его обработает. Дальше у нас начинается логика обратного декодирования. И тут вспоминаем, что Cassandra написана на Java, и — вот чудо — чтобы написать логику сетевого взаимодействия на стороне сервера, взяли тоже библиотеку Netty. Так что то, что мы обсуждали для клиента, в каком-то смысле симметрично применимо и для серверной стороны.
[02:01] Александр: Вообще говоря, это же не обязательно было? Могла быть Netty на сервере, а на клиенте что-то другое. То есть тебя не стопроцентно Netty заставляет использовать и там, и там, но это, наверное, было логично.
[02:15] Дмитрий: Да, Netty — одна из самых распространённых и протестированных в продакшене сетевых библиотек. Поэтому это логично. Если мы говорим про Java, это достаточно удобно: мы знаем, что на ней можно написать высокопроизводительное приложение, и там есть полезные плюшки. Во-первых, в ней есть нативная реализация. То есть Netty позволяет работать не только с Java-сокетами — это дефолтный NIO-режим, когда она работает с неблокирующими сокетами из JDK. Но в ней есть ещё и нативный режим, причём этот режим ты можешь получить, практически просто изменив конфигурацию, без переписывания кода. Твой прикладной код, который работает с Netty, не меняется, но под капотом вместо JDK-сокета будет нативная реализация. Это C-библиотека, упакованная в саму Netty, которая использует нативные epoll-вызовы. Она подтюнена по сравнению с JDK-реализацией: там как минимум явно меньше аллокаций Java-объектов, и в целом она работает быстрее. Поэтому в сервере Cassandra…
[03:30] Александр: Включаем помолчание.
[03:31] Дмитрий: Речь, разумеется, в данном контексте идёт про Linux, потому что epoll — это Linux-специфичная вещь. Аналогичный нативный код в Netty есть и для macOS, он на основе, по-моему, технологии kqueue. Это всё можно подключить, но я предполагаю, что мы говорим всё-таки про Linux-системные вызовы. То же самое, кстати, применимо и к клиенту: ты можешь увеличить производительность своего клиента и немного уменьшить аллокации памяти, добавив нативную библиотеку в classpath. Она называется что-то вроде netty-transport-native-epoll — какие-то такие координаты. И если ты добавишь её в classpath, плюс твоя операционная система позволяет, то Netty автоматически её подхватит. Ты увидишь — там, по-моему, пишут специальные сообщения в логе о том, какой вид транспорта используется. То есть это не надо конфигурировать, как стартер в Spring Boot; это конфигурация на уровне Netty. Она сделана в коде драйвера Cassandra и в коде сервера Cassandra: там просто написана логика — если класс в classpath доступен, попытайся включить этот режим, а если нет — включи дефолтный.
[04:47] Александр: Понятно. То есть с точки зрения Netty там надо некий код написать, но этот код уже написан и в сервере, и в клиенте, и тебе нужно только подложить нужные джарники. Слушай, интересный момент. А чем они отличаются концептуально? Понятно, что одна — нативная реализация, которая делает системные вызовы силами операционной системы, и то делает системные вызовы, и то делает epoll. Просто в случае JDK, в моём понимании, мы частично ограничены тем, что есть некий Java API, который мы должны реализовать, и этот API уже не поменяешь.
[05:31] Дмитрий: Да, стандартный Java API поменять очень сложно. Поэтому есть вот эти сущности — все эти selector’ы и так далее, — которые мы выставляем в неблокирующих сокетах наружу, и уже от них приходится плясать. А в Netty что захотели, то и сделали — как говорится, для самих себя делаем. И вот этот интерфейс взаимодействия между нативным кодом и Java-кодом в Netty просто оптимизирован за счёт того, что мы полностью его контролируем и не обязаны поддерживать тут какую-то обратную совместимость. Плюс, возможно, ребята из Netty просто больше вложились в этот интерфейс по сравнению с ребятами из JDK. Но это мои предположения.
[06:11] Александр: Но оно реально быстрее работает, да? На бенчмарках.
[06:13] Дмитрий: Да, у меня получалось где-то процентов 10, 10–15.
[06:18] Александр: Очень приятно. Разница видна.
[06:20] Дмитрий: Да, за бесплатно. Похожая история, даже в более сильном варианте, применима к TLS. Мы уже упоминали, что можно включить TLS. Есть реализация TLS через JDK, когда мы используем SSL-сокеты. И в Netty есть тоже нативная реализация на базе библиотеки, по-моему, BoringSSL. То есть опять-таки мы напрямую используем C-код, который реализует SSL. Точных замеров я там не делал, но по моим ощущениям разница ещё сильнее. Поэтому, особенно если ты беспокоишься о производительности, эта логика была бы тебе полезна. В целом, если пишешь какие-то нативные коммуникации через Netty, держи в уме, что там есть такие оптимизированные реализации, которые достаточно легко подключить, и твой код за бесплатно получит буст.
[07:22] Александр: Практикующим домохозяйкам на заметку. Записывайте рецепт.
[07:28] Дмитрий: Так. Получили эти байты со стороны Netty. Там всё те же event loop’ы в Netty, которые разгребают наши входящие реквесты. То есть побежал некий поток внутри самой Netty, который прочитал байты из сокета. Вспоминаем всю нашу историю про то, что не факт, что ты получишь все байты сразу: может прилететь половина запроса, а может — сразу два запроса в этот сокет. То есть дальше мы должны заниматься декодированием этих запросов. У нас есть некие объекты-декодеры, которые подключены в архитектуре Netty — она называется pipeline. Там есть набор хендлеров, которые работают в цепочке и разгребают этот входящий поток байтов. В какой-то момент мы накапливаем эти байты в буфере. Когда накопили достаточно, чтобы получилось полное сообщение, мы вытаскиваем его из буфера и декодируем в high-level объект Cassandra, который называется frame. Внутри этого фрейма лежит один или несколько запросов. Давай пока для простоты будем говорить, что один фрейм — один запрос.
[08:38] Дмитрий: Итак, мы декодировали наружную часть запроса, для начала вытащили хедер-поля и поняли, что это вообще за тип запроса. Аутентификацию, как мы уже проговорили, мы сделали на этапе установления соединения — то есть мы знаем, с кем общаемся, знали имя пользователя. И дальше первая задача, которая перед нами стоит… Вернее, сломать, на самом деле. Потому что если кто-то будет посылать очень много запросов, а мы в Cassandra всё это будем обрабатывать асинхронно — много потоков и так далее, — то вполне легко получить ситуацию, когда мы разгребаем эти запросы из сетевого сокета, не успеваем их обрабатывать, всё это копится в памяти, и рано или поздно мы падаем с out of memory. Вот эта история про rate limiting, а точнее про контролирование входящего потока — скорее это можно назвать flow control, — первая вещь, которую нужно сделать на входе.
[09:34] Александр: Да, на входе в сервер.
[09:37] Дмитрий: Она появилась в Cassandra не сразу. Народ собрал эти шишки уже в процессе. Я думаю, где-то в районе третьей версии это появилось в полноценном виде. И там применяется следующая идея. Мы говорим, что есть некий суммарный объём памяти для in-flight сообщений. То есть нам пришёл запрос, мы его обрабатываем, и до тех пор, пока мы не ответили на него, мы считаем его in-flight — то есть он в процессе выполнения. И что мы делаем? Мы ограничиваем суммарное количество байтов для таких in-flight запросов. То есть говорим: в системе должно быть не больше стольких-то — не сообщений, а суммарного размера сообщений. Потому что сообщения могут быть разные: короткое или большое, память они потребляют по-разному. Поэтому идея в том, что мы ограничиваем не количество сообщений, а суммарный объём памяти.
[10:33] Александр: Это, кстати, хорошая идея, очень здравая. Когда у тебя проблема на сервере в том, что ты можешь переполниться по памяти и получить out of memory, то и ограничивать тебе нужно память — именно тот параметр, который у тебя является проблемным. А некоторые действительно могли бы подумать: ну что там, тысячу запросов максимум будем держать, и как-то ограничим. Но тысяча запросов — это непонятно сколько памяти, а у тебя проблема в памяти. Поэтому ограничивай память. Это просто, но понимание этой мысли, мне кажется, помогает.
[11:10] Дмитрий: На самом деле понимание того, сколько памяти реально авансируется при обработке запроса, — это очень сложная задача. Поэтому здесь сделан упрощённый вариант. Мы просто говорим, что, наверное, этот объём памяти пропорционален размеру входящего сообщения — того, которое мы только что прочитали из сокета. Мы знаем его размер, он был написан в хедере.
[11:33] Александр: Да-да-да. То есть это число мы бесплатно получили.
[11:35] Дмитрий: Вот этот размер мы и используем в качестве той метрики, того ресурса, который мы потребляем. У нас есть некий глобальный лимит в Cassandra, который по дефолту выставлен как 10% от хипа. Такая цифра там выбрана эмпирически. И мы говорим, что из этого лимита вычитаем память, когда запрос к нам пришёл, а когда ответим в сокет обратно — эту память в кучку возвращаем. И если мы подошли к нулю, то есть не можем для следующего запроса взять память из лимита, то мы находимся на этапе перегрузки. И тут два варианта. Первый вариант контролируется, на самом деле, клиентом. Клиент, когда устанавливает соединение, флагами передаёт, какое поведение он хочет в этом случае. Есть поведение «вернуть клиенту ошибку» и сказать: извини, чувак, слишком много посылаешь запросов, server too busy, приди позже. Но для этого надо, чтобы клиент всё это реализовывал — перепосылки, дожидания. Это сложно, обычно никто ничего не реализует или ленится.
[12:50] Александр: Да и вообще говоря, дедосить нас могут не обязательно из нашего же клиента. Могут просто написать альтернативный клиент, который не следует этому, и фигачить. Нам как серверу, наверное, было бы слишком оптимистично закладываться, что мы работаем исключительно с собственными клиентами. Или такое есть в контракте?
[13:09] Дмитрий: Нет, есть просто протокол, который описан в виде спецификации, и дальше ты его реализуешь. Был период, когда этих драйверов было много, но на самом деле сейчас всё сошлось к стандартным драйверам. Я не видел, чтобы кто-то прямо активно писал собственную реализацию, хотя, по-моему, вариации этого были. Скорее всего, товарищи в ScyllaDB написали, возможно, драйвер для Rust — кажется, что-то такое они делали. Также эта спецификация полезна ещё чем? В другом случае, когда декодировать протокол интересно не только клиентам. Очень полезная вещь — уметь декодировать протокол по TCP-дампу. Когда ты промышленник, включаешь что-то в продакшен, один из вариантов разобраться в каких-нибудь сложных проблемах — это снять TCP-дамп, а потом использовать какую-нибудь тулзу типа Wireshark, чтобы посмотреть, что там реально происходит. Было бы удобно, чтобы в Wireshark у тебя был декодер. Кстати, для Cassandra CQL такой декодер существует. Можно найти плагин для Wireshark, который умеет декодировать Cassandra request и response. Он, по-моему, не поддерживает последней версии протокола — когда я год назад глядел на него внимательно, — но в целом базовые вещи умеет декодировать. И удобно, когда есть спецификация, по которой можно его написать, а не пытаться выковырять это тайное знание из внутренностей клиентов.
[13:37] Александр: Да, круто.
[13:40] Дмитрий: Значит, один из вариантов, что делать, если нас перегрузили, — вернуться к клиенту с ошибкой. А второй вариант, который чаще всего используется, — это использовать свойство flow control, встроенное в TCP-стек. Мы можем, на самом деле, просто перестать читать из сокета как сервер. Не пытаться вычитывать из сокета и складировать в локальную память, пока не сдохнем, а просто перестать читать. Что в этом случае происходит? Мы перестаём читать из сокета на стороне сервера. Есть некий буфер, входящая очередь, где накапливаются входящие пакеты. В какой-то момент этот буфер заполнится. А дальше на уровне самого TCP есть — прямо в спецификации TCP — раздел flow control. Там есть разные реализации, но идея в том, что есть некое окно передачи байтов между клиентом и сервером, по которому клиент и сервер динамически договариваются. И в какой-то момент сервер просто говорит клиенту: мне больше посылать нельзя, у меня входящий буфер забит, приложение его ещё не разобрало. Ты, конечно, можешь мне посылать эти байты, но я буду ругаться и говорить, что у меня места нет. Соответственно, дальше происходит накопление байтов на стороне клиента — там аналогично буфер для отправки, это называется send queue для нашего конкретного сокета. Понятно, что есть ещё всякие буферы на уровне сетевой карты и так далее, но мы сейчас говорим именно про буферы на точке взаимодействия между операционной системой и приложением, для нашего конкретного сокета.
[16:11] Дмитрий: В общем, этот буфер отправителя на стороне клиента тоже в какой-то момент забивается, поскольку он не может проталкивать байты в сеть. И в конце концов это доходит до клиента: если у тебя выбран блокирующий сокет, ты будешь вызывать send, и его метод будет висеть, пока не удастся отправить. А в случае NIO тебе NIO будет говорить: извини, канал не writable, — и тоже тем или иным способом заставит ждать. Тем самым мы как бы замедлим клиента, и ему придётся что-то делать, чтобы искусственно замедляться, реагировать на такую ситуацию.
[16:52] Александр: Ага. То есть получается два способа, и в Cassandra используется вот такой?
[16:55] Дмитрий: А в протоколе — оба, как я сказал: клиент может выбрать, как сервер должен реагировать для его соединения. Это тоже очень интересный момент. По умолчанию именно вариант TCP backpressure используется чаще всего. TCP backpressure — это тот, который по сути реализован средствами TCP-протокола: то, что я описал, — забивается очередь, приостанавливается передача, забивается очередь отправителя, и в итоге клиенту приходится что-то делать, потому что в сокет не пихается.
[17:29] Александр: Опять же, это такая, знаешь, простая штука, но когда ты думаешь: а вот если мне такую задачу поставит мой продакт, что мне нужно будет это реализовать, — как я это буду делать? Ты сразу думаешь: ага, ну это надо backpressure-алгоритм, то есть я должен вычитать что-то, посчитать, как-то дать знать… Это на самом деле очень непростая задача.
[17:50] Дмитрий: Да, это непростая — ты сразу грузиться начинаешь. А точнее, архисложная задача, которую народ до сих пор, как бы, не решил окончательно. В общем случае это задача flow control: как нам максимально утилизировать пропускную способность канала. Там куча алгоритмов со всякими вот этими словами — Reno и прочее. Народ диссертации целые пишет на тему того, как сделать вот эти протоколы flow control, потому что это ещё и динамическая система: ты не просто поставил фиксированный rate; у тебя в случае TCP ещё и канал связи может сейчас лучше работать, потом хуже. И что интересно, эта же задача применима и для приложений. Я видел статью Netflix, где как раз использовалась идея из TCP на уровне прикладного стека: как нам сделать динамический rate limiting — не просто поставить rate в константу, а сделать так, чтобы клиент прощупал границу, когда всё ещё передаётся хорошо, а когда он заходит за некий предел, где сервер начинает отвечать плохо, — он немного снижает свой rate. Он всё время пытается прощупать максимальную границу, чтобы использовать сервер по максимуму, но при этом не перегрузить его. И вот эти алгоритмы — что в TCP, что в задаче прикладного уровня — весьма похожи, идею из TCP можно перетащить в прикладной слой. Народ такое делал: concurrency limits, по-моему, называется статья, в техническом блоге Netflix можно её найти.
[19:38] Александр: А мы возвращаемся к нашим сокетам.
[19:41] Дмитрий: Да, возвращаемся к сокетам. То, что можно было бы сделать… Ну, существует целое семейство алгоритмов и инженерных практик — backpressure, flow control и так далее. То есть можно было бы, например… [здесь у Cassandra всё устроено проще] — если все потоки заняты, читать из сокетов некому, никакие внутренние очереди в собственной памяти приложения не копятся, и всё прекрасно работает.
[20:53] Дмитрий: Как только ты начинаешь какую-нибудь синхронщину делать, актёров приносить и так далее — возьмём библиотеку Akka, — там сразу появляются все эти алгоритмы rate limiting и backpressure, которые тебе надо реализовывать с помощью каких-то готовых классов, утилит, но самому. Потому что ты перешёл в синхронную логику и всю эту доступную из коробки вещь потерял.
[21:29] Александр: Давай вернёмся к нашему запросу. Допустим, мы прошли этот этап rate limiting.
[21:36] Дмитрий: Есть у нас лимиты на сервере, мы можем этот запрос обработать. Дальше он вытаскивается из NIO event loop и передаётся в процессинг в потоке прикладного приложения. То есть мы не хотим слишком долго заниматься прикладной процессинговой логикой в потоках Netty — это нерекомендуемый подход. Мы эту задачу перекладываем в отдельный thread pool в Cassandra и дальше уже работаем в нём.
[22:08] Александр: То есть мы берём фрейм, да? Или, правильно, давай — фрейм или реквест? Это не одно и то же?
[22:15] Дмитрий: В текущей версии протокола различаются понятия фрейма и месседжа. Один реквест или респонс — это месседж, они упакованы во фреймы. Тут можно провести аналогию: как у нас есть TCP-пакеты и IP-пакеты, что-нибудь такое. То есть одно сообщение может быть размазано на несколько фреймов в Cassandra, когда оно очень большое. И может быть наоборот — в один фрейм, по-моему, ты можешь запихать несколько сообщений, если я не ошибаюсь.
[22:53] Александр: Окей, но мы работаем на уровне сообщений, да?
[22:53] Дмитрий: Да, давай пока забудем про это. Это уже довольно низкоуровневая логика, плюс она появилась недавно, чаще всего мы её видеть не будем. Поэтому считаем, что прилетел реквест. Но это реквест как бы по сути… мы знаем только хедер, его как-то определили, но body мы ещё не распарсили.
[23:12] Александр: Только хедер разобрали, да. То есть мы вычитали это как некий набор байтиков, знаем длину, знаем какие-то базовые поля — вот там correlation id и прочее.
[23:24] Дмитрий: А дальше нам надо его обработать. И первое, что мы делаем, — вытаскиваем информацию о том, что это за тип запроса. Это SELECT, это INSERT, это какой-нибудь ALTER TABLE, может быть. И по этому типу запроса мы выбираем некий хендлер-процессор, который его будет обрабатывать. У нас есть разные реализации для разной логики в Cassandra. Допустим, мы сейчас разбираем случай записи, то есть у нас передался какой-нибудь INSERT, DELETE или UPDATE. Всё это попадает в хендлер, который отвечает за соответствующую операцию. А первое, что мы должны сделать, когда нам прилетел такой запрос, — вот что бы ты сделал, ещё до того, как начать непосредственно обработку?
[24:19] Александр: Я хендлер, да? То есть мне уже пришли и говорят: давай, делай логику, вот тебе request, вставляй давай.
[24:27] Дмитрий: Вставляй.
[24:27] Александр: Ну, я бы сразу вставлять не стал. Я бы посмотрел. Может быть, я уже это сделал?
[24:33] Дмитрий: Нет, ещё до этого этапа, на самом деле. Вспоминаем: security, всё такое — надо проверить, а можно ли тебе вставлять?
[24:41] Александр: Ага, то есть аутентификацию пройти.
[24:44] Дмитрий: Ну, аутентификацию-то ты прошёл при подключении соединения. А дальше идёт авторизация: есть ли у твоего пользователя права на запись в эту конкретную таблицу. Для этого нам надо сначала немного дальше разобрать запрос.
[24:59] Александр: Да-да-да. То есть мне уже нужна информация — недостаточно username и пароля, который нужен, чтобы сделать аутентификацию. Для авторизации мне, например, нужно ещё понять, с каким объектом мы сейчас хотим взаимодействовать. Я знаю, кто; теперь мне надо знать, с чем, и какие у этого кого-то есть права. Я вот, например, как человек, который реализовывал для такой системы role-based access control, сейчас удивлён, почему сразу не ответил — потому что это именно то место, где это надо делать.
[25:28] Дмитрий: Слишком очевидно. Так вот, да, мы, соответственно, должны сначала разобрать запрос, понять, с какими ресурсами — это так в терминах Cassandra называется — мы работаем.
[25:40] Александр: Они уже в теле находятся, да? Где-то там, в начале.
[25:43] Дмитрий: Да, у нас в теле находится либо сам запрос — это тот самый simple statement, который мы обсуждали в прошлый раз, там прямо текстом написано: INSERT в такую-то таблицу, keyspace, таблица. Либо там находится prepared statement, и в этом случае лежит просто id — MD5-хэш от этого запроса. Естественно, нам нужно сначала этот MD5-хэш зарезолвить и получить разобранное представление. Поэтому либо мы на предыдущем этапе, когда обрабатывали prepared statement, распарсили запрос, разложили, вытащили информацию о том, с какой таблицей и с какими полями он работает, в некую внутреннюю структуру, подготовили и положили в кэш, — либо ты прислал этот запрос текстом, и этим мы прямо сейчас занимаемся на начальном этапе. То есть мы парсим этот запрос и раскладываем его на сущности. По факту там происходит лексический и синтаксический анализ. Внутри Cassandra используется ANTLR: мы разбиваем текст запроса на лексемы и дальше строим абстрактное дерево согласно нашей грамматике, которая для Cassandra Query Language написана и скомпилирована в некий код парсера.
[27:09] Дмитрий: Значит, мы разобрали, знаем таблицу, с которой ты работаешь. Дальше проверяем, а можно тебе это делать или нет. В Cassandra эта информация хранится в служебных таблицах. Там написано, что такой-то пользователь, например, имеет права на работу с такой-то таблицей для такого рода операций: ты можешь из этой таблицы читать либо в эту таблицу писать. Но эти правила более гибкие. Ты можешь, например, написать, что имеешь права не на каждую таблицу по отдельности, а на весь keyspace. То есть эти ресурсы с точки зрения security образуют некую иерархию: у тебя есть возможность задать гранты на права на конкретную таблицу, а может быть — на вышележащий контейнер, keyspace, в котором они лежат. Либо, может, ты вообще суперадминистратор, и у тебя есть права сразу на все keyspace’ы. И всё это записано в неких системных таблицах. Чтобы каждый раз эти таблицы с диска не поднимать, не читать, есть кэш. На каждой ноде Cassandra, когда мы эти данные поднимаем, они кэшируются, попадают в локальный кэш, который написан, опять-таки, на той же библиотеке, что мы уже обсуждали, — Caffeine. То есть Cassandra не реализовала кэш с нуля с хитрыми умными стратегиями, а просто взяла готовую Java-библиотеку. Благо, есть неплохие варианты в Java. Естественно, мы сходили в кэш, и если повезло — там эти security-сущности лежат.
[28:48] Дмитрий: А тут есть особый момент. Этот кэш имеет ограниченный период жизни по времени, и, по-моему, текущая реализация не сбрасывается, если ты что-то меняешь на уровне системных таблиц. То есть если ты, например, дал пользователю дополнительные гранты или отнял их, это изменение не немедленно пропагируется на ноды — как раз из-за кэшей. И здесь наш классический трейд-офф: мы хотим побыстрее работать и всё кэшируем, но возникает задача консистентности кэшей. И здесь мы платим за быстроту доступа отложенностью применения. Можно было бы попытаться сделать какие-нибудь форсированные сбросы кэшей, но, по-моему, сейчас он ещё не сделан, если только ты сам руками не дёрнешь эту операцию.
[29:39] Дмитрий: Вот. Допустим, нам разрешили сделать запись в эту таблицу. Если нет — понятно, мы просто ошибку вернём. А разрешили — идём дальше.
[29:52] Дмитрий: И мы сейчас находимся в той логике Cassandra, которая называется координация. Вспоминаем, что Cassandra — не просто локальная база данных, а распределённая система. И нам, чтобы выполнить этот запрос на запись, нужно пообщаться с несколькими другими нодами. Вот эта логика общения, собирания результатов с этих нод и ответа назад называется координацией. И, как я сказал, любая нода Cassandra может выступать в роли координатора. Мы попадаем в эту координаторную логику, и для начала нам нужно понять, с какими репликами будем общаться. Нам сказали: запиши в такую-то таблицу вот такую-то строчку. У нас кластер, может быть, из тысячи нод. Какие из этих нод вообще отвечают за эти данные? Когда мы создавали keyspace, мы указали для него некий replication factor, например 3. Наиболее типичное значение — 3 в каждом дата-центре. Пока рассмотрим один дата-центр для простоты: из тысячи нод есть 3 ноды, которые для этой конкретной строчки хранят данные. Нам нужно найти эти ноды, и по сути дела это задача шардирования. Данные в кластере Cassandra нарезаны на кусочки — шарды, или в терминах Cassandra партиции, — которые разбросаны по большому набору нод, и есть несколько нод, называемых реплики, которые отвечают за хранение наших данных.
[31:29] Дмитрий: Можно разные алгоритмы делать того, как разбросать эти данные. Cassandra использует механизм, основанный на хэшировании ключей. Что мы делаем? Мы вспоминаем, что в таблице Cassandra обязательно должен быть partition-ключ как часть primary-ключа. У нас есть составной primary-ключ, который состоит из partition-ключа и clustering-ключа, и partition-ключ — это то, что есть всегда, и мы на него сейчас смотрим. Мы берём значение этого ключа — это несколько байтов, возможно, это несколько ключей, тогда мы их просто конкатенируем.
[32:09] Александр: Композитный ключ такой.
[32:10] Дмитрий: Да. И для начала считаем от них хэш. Сейчас это Murmur3. То есть это не обязательно должен быть криптографический хэш — просто он должен быть быстрым и давать хорошее распределение, статистически быть просто хорошим.
[32:27] Александр: Я бы даже удивился, если бы он был криптографическим. Это же время.
[32:32] Дмитрий: Ну вот видишь, для prepared statement когда-то взяли криптографический, хотя там тоже не то чтобы он прям обязателен.
[32:38] Александр: Хотя там всё-таки мы хотим более низкую вероятность коллизий, чем здесь.
[32:43] Дмитрий: Ну в общем, Murmur3. Сейчас, может быть, существует что-то и получше — появились всякие хэши побыстрее, которые можно вычислять, там всякие CityHash и прочее. Основная проблема в том, что мы хотим считать этот хэш ещё и быстро. Но это Murmur3, и поскольку он используется при взаимодействии между различными нодами, то даже из соображений обратной совместимости нам не хотелось бы его менять. А дальше нам надо вот этот хэш, который по сути дела… Сколько это было? 128?
[33:16] Александр: 128 битов, по-моему.
[33:17] Дмитрий: Нам нужно как-то ассоциировать с какими-то тремя нодами из нашего кластера. И мы хотели бы сделать это таким образом…
[33:27] Александр: То есть можно было бы взять, например, остаток от деления. Чего сложные вещи городить? Вот число у нас, поделил на количество нод, получил остаток от деления — и вот номер ноды, который отвечает за нашу строчку.
[33:45] Дмитрий: Ну да, самый простой такой.
[33:47] Александр: Следующие реплики — как бы плюс один, плюс два к этому числу добавить, например. Цепочкой делать — и всё. Что тут дальше городить?
[33:54] Дмитрий: Ну да, берёшь остаток от деления и погнал. В хэш-таблице уже так работает, почему бы тут не работало?
[34:02] Александр: Проблема в том, что у нас кластер динамический.
[34:06] Дмитрий: То есть мы можем менять его состав. Мы можем добавлять ноды, можем удалять ноды из этого кластера. И для этого алгоритма есть проблема переноса данных. Допустим, у нас было 10 нод, а мы добавили ещё одну, стало 11. И когда мы берём разные числа на вход и берём остаток от деления на 10, а потом на 11, они будут давать разные результаты. Это означает, что когда мы добавили в кластер ещё одну ноду, маппинг между нашими ключами и нодами полностью поменялся. Вот это очень важный момент.
[34:51] Александр: Давай попробуем сейчас для слушателей, которые идут где-то по лесу или сидят за компом, что-то программируя в фоне… Ребят, подключаем нейроны. Нам сейчас надо нарисовать в голове картинку, чтобы понять, как это работает. На каком этапе мы находимся? У нас есть 10 нод — это 10 таких коробочек. И есть ключ, который к нам пришёл. Нам надо понять, в какую коробочку этот ключ положить. Мы берём от него хэш, 128 бит, берём остаток от деления на 10, получаем число в диапазоне от 0 до 9, выбираем нашу коробочку по номеру и в неё кладём. И каждый раз, когда этот ключ теперь приходит на чтение, мы такие: хэш один и тот же, остаток от деления на 10 — пятая коробка, сюда. Так разложили огромное количество данных. И теперь, поскольку кластер динамический, к нам приходит ещё одна коробочка. Коробочек становится 11. И теперь остаток от деления на 11 у тех же самых ключей приводит нас в другие коробки. И всё ломается — мы, по сути, ослепли на один глаз. И теперь нам надо что-то с этим делать, нужен какой-то ребаланс. Короче, у нас проблемы. И вот сейчас мы будем эту проблему решать, правильно?
[36:09] Дмитрий: Да. И проблема в том, что это ребаланс глобальный. Понятно, что какой-то ребаланс будет: мы добавили одну ноду, мы ожидаем, что то, что было разложено по 10 нодам, теперь должно быть разложено по 11 равномерно. Мы же не хотим, чтобы добавили ноду, а она была пустая, — мы хотим, чтобы она участвовала в хранении данных. Но мы хотим, чтобы количество данных, которое мы перемещаем между нодами после добавления новых, было минимальным. По факту это количество данных, которое в итоге окажется на этой новой ноде, — вот это наша целевая задача, переместить только эти данные, а всё остальное вообще лучше не менять никак. И эту задачу в математике называют по-разному, я не уверен, что есть устоявшийся термин. Наверное, лучший вариант — называть это в литературе stable hashing. Это хэш-функция, которая, во-первых, хэш-функция, но у неё есть и дополнительные свойства: когда мы добавляем и удаляем количество нод, по которым она хэширует, она держит объём ребалансинга ровно на том уровне, который необходим, чтобы перераспределить данные на новую ноду. То есть она не перебалансирует больше, чем нужно, — только то, что реально нужно.
[37:27] Александр: Вот смотри, тут надо аккуратненько. Вот хэш-функция как математическая функция: f от x, x — это ключ. А что y?
[37:35] Дмитрий: От двух параметров. У нас есть входной ключ и количество нод. То есть это не самая классическая хэш-функция.
[37:45] Александр: Ну, следующий уровень, да.
[37:47] Дмитрий: Да, следующий уровень. Тоже функция, внутри неё есть классическая хэш-функция, но одна из. А эта функция следующего порядка отвечает нам на вопрос не какой хэш у ключа, а на какой ноде этот ключ.
[37:59] Александр: Типа циферку поменьше даёт. Окей.
[38:02] Дмитрий: Да, на самом деле у тебя такие же функции внутри хэш-таблицы даже есть. Они просто заинлайнены, спрятаны, потому что хэш-таблица — это какой-то набор бакетов, и ты также делаешь остаток от деления, например. А может, и не делаешь — разные алгоритмы есть, но сейчас всё-таки чаще остаток от деления там бывает. И там тоже бывают задачи ребалансировки, потому что в какой-то момент ты понимаешь, что таблица переполнена, и тебе нужно количество бакетов увеличить. И тут начинаются всякие весёлые истории, как это сделать, желательно не сильно долго блокируя таблицу. Так что даже там такие задачи стоят. Но вернёмся к нашему распределённому случаю.
[38:45] Дмитрий: Есть разные алгоритмы, которые решают эту задачу. Она применима не только для распределённых баз данных, но и, например, для балансировщиков. У нас может быть какой-нибудь load balancer, который балансирует между серверами с каким-нибудь статическим контентом — CDN-сети. Прилетает куча запросов, данные расшардированы по куче серверов, и мы хотим балансировать запросы так, чтобы направлять их в одни и те же серверы, потому что там кэши будут прогреты, и мы будем отвечать быстрее и меньше грузить диски. Поэтому такая же задача стоит перед авторами многих балансировщиков. Если посмотреть всякие балансеры — даже open-source, на Envoy какой-нибудь, — или описания того, как реализованы балансировщики у Google, у Amazon, там это отдельная задача, которую люди решают. Есть разные алгоритмы. Исторически, наверное, два первых алгоритма появились примерно в одно время, где-то в начале 90-х, в 93-м, наверное. Один был изобретён в MIT, другой не помню где. Это consistent hashing, о котором мы сейчас поговорим, поскольку он сделан в Cassandra. И другой интересный алгоритм называется rendezvous hashing. Я не знаю, в текущей версии Ignite он используется или нет, но во второй версии Ignite его точно использовали. И есть алгоритмы, которые всякие балансировщики — типа Maglev в Google — делали для той же задачи: там уже более хитрые логики, с решением неких линейных уравнений. В общем, весёлая математическая задача, которую можно решать разными интересными способами.
[39:33] Александр: Каким способом её решила Cassandra?
[39:37] Дмитрий: Так. Вот у нас есть наш хэш. Это число в каком-то диапазоне. Для простоты возьмём от минус maxint до плюс maxint — некий отрезок. Мы этот отрезок как прямую линию закрутим в кольцо. Представляем себе кольцо: взяли за два конца нашу прямую, от минус maxint до плюс maxint, и закольцевали. Получился ring — собственно, этот термин можно встретить в коде Cassandra. И наш хэш — это точка на этом кольце. Мы посчитали хэш от нашего ключа, и он куда-то попал на кольцо.
[40:23] Александр: Так, пока понятно. То есть мы бросили эту точку туда.
[40:27] Дмитрий: Теперь мы каждый наш сервер тоже ассоциируем с какой-то точкой. Для начала с одной. Мы, например, случайным образом — а может, и не случайным, каким-то способом — выберем точку для каждого сервера. И у нас, допустим, 10 серверов. Мы на это кольцо поместим 10 точек. Можно представить, что есть точки другого цвета: нам прилетел запрос — это красная точка, попавшая на кольцо, а вот есть серверы, которые один раз вычислили и поставили свои точки, синие.
[41:03] Александр: Слушай, мне кажется, тут может подойти неплохая аналогия из нашего прошлого выпуска — с циферблатом. Вот ноды — они как бы там всегда, ну пока что. И это циферки на нашем кольце, от 1 до 12: мы их расположили, они там пока есть. А конкретный хэш — это вот мы дротик кинули, куда попали по циферблату, там и есть. Это конкретный хэш конкретного ключа.
[41:29] Дмитрий: Да, в секторе, рядом с пятёрочкой.
[41:31] Александр: Да.
[41:31] Дмитрий: А дальше мы говорим: окей, нам нужно для нашей точки, соответствующей хэшу, понять, к какому серверу она относится. И что мы сделаем? Мы пойдём по этому кольцу, например, по часовой стрелке, до тех пор, пока не встретим точку, связанную с сервером. Бежим-бежим. И вот попали мы, например, где-нибудь на полвторого и побежали, добежали до двух часов, а там наша точка стоит. Соответственно, это сервер, который отвечает за нашу запись. Это реплика, которая будет владеть этой записью.
[43:12] Александр: Ну, по сути, мы получили ответ на вопрос, какая нода. Мы проделали действие с хэшом и расположением на кольце и пришли в двоечку, в два часа. Это наша нода, возвращаем два. Всё.
[43:27] Дмитрий: Да. Если нам нужно несколько реплик, мы можем пробежать чуть дальше и собрать побольше точек. Например, нам нужно три реплики: мы встретили первую, чуть пробежали, там где-то в районе трёх или четырёх часов оказалась вторая серверная точка, и потом ещё где-нибудь третью нашли. Собрали. Эти точки говорят, что вот это три реплики, отвечающие за нашу строчку.
[44:01] Александр: Тут суть, наверное, в том, что мы именно, знаешь, бежим по направлению часовой стрелки. То есть мы не вычисляем как-то хитро по делению и остаток не берём — сразу получаем ответ. А мы реально бежим по кольцу: вот прямо в коде там цикл есть, и мы перебираем эти ноды, правильно?
[44:06] Дмитрий: Ну, не совсем. Потому что если ты будешь перебирать, это алгоритм, пропорциональный количеству этих точек, от n. В худшем случае ты можешь вообще всё кольцо пробежать и только в конце понять, что надо было сразу в конец идти. Похитрее. Эти точки у тебя упорядочены, это числа, которые имеют отношение «больше — меньше». У тебя задача — найти среди набора чисел число, которое максимально близко к твоему, равно ему или меньше его. Не напоминает ли это тебе какую-то задачу?
[44:46] Александр: Поиск ближайшего соседа или поиск по дереву?
[44:51] Дмитрий: Проще. На самом деле это поиск числа в упорядоченном массиве, который мы решаем бинарным поиском.
[44:57] Александр: Всё так. Как бы тут только просто не условие «равно», а «не больше чем».
[45:04] Дмитрий: Да, небольшая модификация этой задачи, но алгоритм по сути тот же самый. Мы можем делать, например, бинарный поиск среди вот этих отрезочков. У нас есть набор упорядоченных отрезков, и мы алгоритмом спускаемся и находим нужный сегмент. Какие есть проблемы у этого алгоритма? За счёт чего у нас получается вот это перераспределение? Как мы задачу ребалансировки здесь решили? Допустим, у нас появилась ещё одна нода. Мы её тоже ассоциировали с какой-то новой точкой, добросили ещё одну точку в это кольцо. Эта точка попала на какой-то отрезок между двумя соседними нодами. Допустим, было 10 нод, и она попала между седьмой и восьмой, например. И этот отрезок она своим попаданием разбила на два. То, что раньше попадало между седьмой и восьмой точкой, принадлежало восьмой точке — потому что мы идём по часовой стрелке, и всё, что перед ней до предыдущей точки, это восьмая. Мы туда, в середину, попали новой нодой, и она откусила часть этого сектора под себя — забрала, например, половину того, что принадлежало восьмой ноде. При этом все остальные точки никак не заимпактили: как было, так и осталось, вот этот маппинг. То есть задачу вида «пожалуйста, оставь нетронутым то, что не двигается» мы решили.
[46:44] Александр: Да, но мы приобрели новую проблемку.
[46:47] Дмитрий: Да, мы нарезали не очень равномерно. Мы отняли половину от восьмой ноды, захватили её в свою новую ноду, а все остальные ноды остались того же размера. Равномерности мы не добились. Смотри: добавив новую ноду в кластер из 10, где каждая хранит по 100 гигабайт, я разгрузил только одну — они стали по 50 гигабайт. Наша восьмая нода хранила 100 гигабайт, а теперь хранит 50, и 50 ушло. И то это если я в середину попал. А если я попал очень-очень близко к восьмёрке, буквально в 59 минут восьмого, то я почти 99 гигабайт буду у себя хранить, а ей останется один. То есть равномерного распределения данных между 11 нодами не произошло, глобальную задачу я не решил.
[47:47] Александр: Ну да, не все условия ты соблюдаешь при этом.
[47:51] Дмитрий: Поэтому как можно этот алгоритм обобщить? А давайте мы для каждой серверной ноды в нашем кластере будем ставить на этом кольце не одну точку, а несколько. Это называется в терминах Cassandra vnode, virtual nodes, виртуальные ноды. То есть одной ноде у нас соответствует, например, 16 точек на этом кольце. Каждая нода бросила рассыпь, эти точки как-то перемешались на кольце, и получилось, что каждой ноде теперь принадлежит не один сектор. В предыдущем варианте алгоритма каждой ноде доставался только один сектор этого круга. Теперь мы разбили круг на множество маленьких секторов, и каждой ноде принадлежит набор секторов — столько, сколько точек мы бросили. Такой разноцветный круг, состоящий из множества секторов разных цветов: можно представить, что каждая нода — это цвет, и вот такой пёстренький круг у нас получился. И задача поиска, какая нода у нас реплика, решается ровно так же: посчитали хэш, бросили точку на кольцо, пробежались, ближайшую точку, относящуюся к ноде, нашли, посмотрели, какой ноде она принадлежит — такая реплика. Тут ничего не поменялось. Но когда происходит добавление новой ноды, мы также вбрасываем множество точек. И если в прошлый раз мы откусывали от сегмента только одной ноды, то теперь будем откусывать 16 маленьких кусочков от каких-то других нод. И распределение становится более равномерным. А дальше начинается игра: чем больше точек мы накидаем, тем более равномерное распределение, но тем больше будут оверхеды на хранение всех этих метаданных, на подсчёт, на алгоритмический поиск и прочие оверхеды, о которых можно поговорить отдельно. В Cassandra изначально рекомендованные значения были где-то в районе 256 — в старых поколениях Cassandra.
[50:02] Александр: Точек на одну ноду, правильно? То есть это нормальная такая россыпь.
[50:07] Дмитрий: Да. И там есть всякие статьи, где оценивали: если равномерно накидать, то какое будет распределение, насколько равномерным оно будет между разными нодами. Но, как я сказал, это не бесплатное удовольствие, поэтому в какой-то момент мы сказали: окей, давайте мы точки будем кидать не случайно, а пытаться целиться в серединки этих отрезков. Будем кидать точки так, чтобы максимально равномерно разрезать сегмент на части, на более равномерные сегменты. За счёт этого мы можем получить всё это с гораздо лучшим распределением. И это позволит уменьшить количество точек. Сейчас дефолтные рекомендации — 8–16 точек на одну ноду, с учётом этого улучшенного алгоритма. Он улучшен за счёт того, что ты ноде подсказываешь, что будешь хранить данные с определённым replication factor. Чтобы ей сделать это хорошее, равномерное распределение, ей надо понимать, сколько ты ещё реплик будешь делать. Поэтому ты ей говоришь: знаю, у меня много keyspace’ов может быть и так далее, но давай под случай, когда 3 реплики, оптимизируемся для каждой записи. Replication factor 3 — вот под этот случай, пожалуйста, распределение по vnode оптимизируй, и у меня будет всё хорошо. За счёт этого дополнительного знания она может сделать более оптимально, и количество виртуальных нод будет меньше.
[51:50] Александр: Кстати, а количество реплик в keyspace я не могу менять после создания, да?
[51:55] Дмитрий: Можешь, но это требует некоторых дополнительных телодвижений.
[52:00] Александр: Ну вот, например, алгоритм consistent hashing уже, наверное, должен…
[52:05] Дмитрий: В consistent hashing он легко тебе ответит. Один из вариантов, который, скорее всего, там и делается: тебе нужно найти n реплик для твоей записи. Ты бежишь по кругу, как я сказал, находишь первую точку — это первая реплика. Дальше бежишь, пока не найдёшь точку, относящуюся к следующей ноде, — потому что, может быть, ты найдёшь несколько подряд виртуальных точек, относящихся к одной ноде. Бежишь по кругу, пока не нашёл точку другого цвета, — это вторая реплика. Ещё бежишь — нашёл третью ноду, третья реплика. Допустим, сказал replication factor теперь 4. Окей, пробегись по кругу и получи четвёртую точку — вот она, твоя четвёртая реплика. Весь вопрос, будут ли там данные, которые ты до этого дописал на предыдущие три реплики, — это уже задача ребалансинга, когда ты делаешь запросы вида «измени, пожалуйста, replication factor». Это отдельная задача.
[53:11] Александр: То есть мы ищем три — ну, в нашем случае три — реплики, и нам всё равно, какая из них primary, secondary. У нас нет такого, да? Мы просто три реплики.
[53:19] Дмитрий: В нашем случае нет понятия primary, да.
[53:25] Александр: Окей. То есть они все равноправные. Алгоритм consistent hashing, который работает в Cassandra, мне стал понятен. Мне кажется, мы очень подробно его объяснили, и теперь мы знаем, как выбираются те ноды, на которые будет происходить запись. Наверное, мы можем идти дальше. А что дальше у нас?
[53:47] Дмитрий: Да. Соответственно, мы знаем, что для того, чтобы записать данные в Cassandra, нам нужно записать на эти три реплики. И тут, возможно, одна из этих реплик вообще находится локально. Вспоминаем как раз про нашу оптимизацию на стороне клиента: клиент пытается, используя ровно тот же consistent hashing ring — это те метаданные, которые клиент тоже может подсмотреть, они в системных таблицах лежат, — понять, что ага, для того INSERT, который я делаю, вот та нода является репликой. Наверное, логично послать запрос в неё, тогда она меньше по сети будет ходить: один из запросов она может сама к себе сделать без похода по TCP-сокету, и вместо трёх запросов по сети сделает только два.
[54:32] Александр: Ну да, но базово мы думаем, что находимся в роли координатора сейчас. То есть мы на сервере, мы серверная нода-координатор, мы выбираем, куда писать.
[54:43] Дмитрий: Да, и нам надо теперь — мы знаем, что у нас есть три ноды, которые отвечают за нашу строчку, — послать в них запрос: запиши данные, которые написаны в INSERT. Это и есть базовая операция записи в Cassandra. Координатор зашлёт запросы в каждую из этих реплик; пока мы считаем, что это чёрный ящик, каждая из реплик ответит «всё, я записала», и мы отвечаем назад клиенту. Но мы хотели бы, наверное, чтобы наша система переживала выпадение нод. У нас распределённая система, и отказоустойчивость в Cassandra сделана как раз для того, чтобы переживать падение нод. А что, если одна из этих реплик ничего не ответит?
[55:29] Дмитрий: И тут появляется понятие, называемое в Cassandra consistency level.
[55:34] Александр: Это что-то новенькое, мы до этого не говорили про consistency level.
[55:38] Дмитрий: Да, ну теперь придётся. Это параметр, который клиент передаёт в запросе. Когда мы написали INSERT туда-то, туда-то, мы положили ещё один из параметров запроса, который называется consistency level. По сути дела, он говорит нам в случае записи, сколько нод должно ответить, чтобы мы считали операцию записи успешной. Это может быть ALL — то есть все ноды должны ответить, тогда мы считаем операцию успешной. Самый строгий, наверное, да? Если какая-то из них не ответит или ответит ошибкой, то ошибкой побежит и клиент: извини, не получилось. Это может быть consistency level ONE, в стиле «хотя бы одна нода пускай ответит». И есть разные уровни, все я перечислять не буду, но самое интересное…
[56:33] Александр: ANY, может быть?
[56:34] Дмитрий: Да, есть понятие ANY. Сейчас мы до него ещё, не знаю, прямо сейчас разберёмся; наверное, давай чуть отложим его. Но есть волшебный ANY, который вроде как означает «хотя бы куда-то запиши, а потом разберёмся». Сейчас мы чуть ближе к этому подойдём.
[56:51] Александр: Окей.
[57:00] Дмитрий: Но самый распространённый уровень, который чаще всего люди используют, — это некий компромисс. И этот компромисс называется quorum. А точнее, есть два вида quorum в Cassandra. Можно сделать quorum в локальном дата-центре: вспоминаем, что у нас кластер ещё и неоднородный, в нём может быть несколько дата-центров. И мы хотим — например, они далеко друг от друга находятся, между ними большая задержка — выполнять операции, дожидаясь ответа только от реплик в локальном дата-центре. Для этого нам подходит consistency level LOCAL_QUORUM, который говорит, что большинство реплик из локального дата-центра должно ответить успехом, чтобы мы ответили клиенту, что операция записи прошла.
[57:41] Александр: Хорошо. То есть мы так или иначе всё равно отправляем запрос на все реплики, но дожидаемся ответа от того количества, которое сконфигурировал клиент.
[57:50] Дмитрий: Да, ну не совсем прямо явно сконфигурировал — указал в виде consistency level. То есть он не пишет явно «2», когда хочет quorum; это прямо такое явное значение, называемое LOCAL_QUORUM.
[58:06] Александр: И enum’чик в Java.
[58:08] Дмитрий: Ну да, по факту да. То есть у тебя будет 5 реплик в локальном дата-центре — соответственно, будет 3 ноды достаточно, чтобы ответили.
[58:19] Дмитрий: Для одной ноды это 1, для двух нод — 2, для четырёх сколько?
[58:24] Александр: Сейчас, подожди. 1, 1, 2, 2… для четырёх — 3.
[58:29] Дмитрий: Ну да.
[58:30] Александр: И для пяти тоже 3.
[58:31] Дмитрий: Пополам плюс 1, условно говоря. Главное — округлить в нужную сторону, там всегда путаница возникает. Но это quorum. Слово quorum здесь понятно почему — это означает большинство.
[58:44] Александр: Majority — вот хорошее слово здесь было бы, потому что quorum нас в распределённые алгоритмы — Paxos и так далее — отсылает. Но это не то, это quorum не из Paxos.
[58:53] Дмитрий: Там есть разные quorum’ы. Есть, да, классический quorum в этих алгоритмах, примерно то же самое означает — вот это n пополам плюс 1. Там есть понятие ещё всяких супер-quorum’ов, например это даже три четвёртых n и так далее. Но это отдельная тема. Ну, так вот, вот так назвали. И второй вид quorum’а — это когда мы хотим собрать quorum с каждого дата-центра. То есть в каждом дата-центре, когда большинство сказало, что всё хорошо, тогда операция записи прошла успешно.
[59:23] Александр: То есть нам нужен quorum quorum’ов?
[59:26] Дмитрий: Не quorum quorum’ов, а в каждом регионе большинство должно проголосовать за, в каждом.
[59:34] Александр: Ага, в каждом. Не в большинстве регионов, а в каждом. Всё, понятно.
[59:40] Дмитрий: Так. Соответственно, у нас есть параметр, который говорит, от скольки нод ждать ответа. И тут можно даже математику делать. У нас есть параметр — количество реплик, replication factor, — извини, n это replication factor, и параметр consistency level, который говорит, скольких нод мы ждём ответа. Соответственно, мы можем сказать: чем выше у нас consistency level, тем менее мы толерантны к ошибкам. Если он ALL, то мы не выживаем, даже если одна из нод выпадет. Если quorum — например, LOCAL_QUORUM и replication factor 3, — то одна нода может выпасть, и мы будем по-прежнему возвращать ответы.
[1:00:14] Александр: Ну, вспоминаем. CAP-теория, мне кажется, самое время. Вот этот треугольничек. И мы говорим, что если consistency level выставляем максимальный, то у нас система CP: она консистентная, но не available в случае выхода любой ноды. А если говорим, что consistency level не такой строгий, то мы этот слайдер двигаем, и система становится более доступной, то есть толерантной к фейлорам. Но консистентность может быть не такой строгой, потому что… Вот, кстати, страдает ли в этом случае консистентность? В случае если majority не должна…
[1:01:14] Дмитрий: Для того чтобы говорить о консистентности, мы должны говорить ещё и об операциях записи. К этому мы сейчас чуть вернёмся. Мы пока говорим про операции чтения… сейчас — про операции записи. Вот спойлер: в операциях чтения тоже есть consistency level, и там будет как раз пересечение, и это, может быть, даже немного напомнит тебе Paxos. Но пока мы чуть забежим вперёд, а пока говорим просто про записи. Соответственно, чем меньше в данном случае consistency level, тем больше количество фейлов мы можем переживать. Итак, мы рассылаем запросы во все реплики. От того количества реплик, которое соответствует нашему consistency level, мы ждём успешный ответ. Там прямо что-то типа CountDownLatch в Cassandra сделано: мы поставили нужное нам количество, прибегают ответы от реплик, они этот count уменьшают, и когда мы дошли до нуля — говорим, что клиенту можем отправить ответ «успешно».
[1:02:19] Дмитрий: Но какие-то из реплик могли не ответить, потому что они мёртвы. И тут есть две ситуации. Мы можем знать, что они точно уже мёртвы, ещё на этапе, когда не начали выполнять отсылку запросов в реплики, — потому что мы каким-то способом мониторим другие реплики. Этот способ в Cassandra называется gossip.
[1:02:45] Александр: Так.
[1:02:45] Дмитрий: То есть когда нода Cassandra попадает в кластер, она начинает взаимодействовать с другими нодами путём неких heartbeat’ов. Она периодически посылает heartbeat’ы в какие-то из других нод Cassandra. Собственно, чтобы в самом начале начать её считать… Для этого ей нужно узнать про другие ноды. Ей нужен адрес некой входной точки в этот существующий кластер Cassandra. Допустим, у нас уже есть кластер, у него есть какие-то ноды. Нам надо понять, с кем вообще новой пришедшей ноде общаться. И для этого в конфигурации Cassandra есть понятие seed-node. Когда ты конфигурируешь Cassandra, на стороне каждого сервера ты должен указать список нод, которые являются seed’ами. Это точки входа, с которыми она начнёт взаимодействовать, когда присоединяется к кластеру. Можно провести аналогию с contact points в клиентах — то же самое, что мы говорили в клиентах: там не должен быть полный перечень нод, там должен быть некто, с кем мы можем пообщаться. Из этих перечисленных нод мы можем вытянуть всю информацию о том, какие вообще ещё ноды есть в кластере, и узнать полную топологию. И дальше мы начинаем периодически рассылать heartbeat’ы в другие ноды, чтобы они знали, что мы живые.
[1:04:09] Дмитрий: То есть мы подсоединяемся к seed, узнаём, какие вообще есть другие ноды в кластере, и начинаем рассылать heartbeat’ы и в эти seed, и в другие ноды. То есть seed — это точка входа в кластер, которая нужна, на самом деле, только в самом начале. Потом они могут сдохнуть и не работать, всё равно кластер останется доступным. То есть это не какой-то мастер, координатор и так далее. Это такая точка, куда ты можешь первый раз прийти и узнать вообще, кто в кластере есть. Потому что если ты не знаешь ни одного адреса, то как ты придёшь и всем скажешь «привет, я тут новенький», если вообще никого не знаешь.
[1:04:50] Александр: Собственно, да. По какому айпишнику?
[1:04:51] Дмитрий: Ну, кстати, есть разные способы динамической топологии и всё такое, когда ты просто поднимаешь ноду. Во втором Ignite такое реализовано: когда ты просто поднимаешь локальные ноды, и они как-то узнают друг о друге.
[1:05:06] Александр: Ну, у тебя тогда есть, скорее всего, какой-то паттерн типа whiteboard, когда ты на какой-то доске объявлений пишешь, например, что организуешь кружок по интересам под названием Cassandra Cluster, и он будет проходить по такому-то адресу в такое-то время, — и там тогда все собираются. В нашем случае такой доской объявлений может быть какая-то система типа ZooKeeper, etcd или ещё что-то, где все публикуют информацию, и каждый может прийти в эту точку и узнать про остальных. У этого есть свои плюсы и минусы.
[1:05:44] Дмитрий: Собственно, Cassandra немного по-другому сделала. Там были, по-моему, даже идеи по поводу ZooKeeper в качестве варианта, но по нему не пошли, потому что это ведёт нас, опять-таки, к некоторой централизации, которую Cassandra хочет избежать. На самом деле этот gossip у нас, получается, решает две задачи. Первая. Вот когда мы говорили про это кольцо с нодами и так далее — нам нужно, чтобы все ноды знали про все ноды. Эта задача в computer science называется membership: есть некая группа, и нам надо знать, кто в эту группу вообще попал. А вторая задача — это узнать, кто в этой группе живой в настоящий момент. Потому что мы сейчас хотим запросы посылать, а что посылать, если с той стороны мы знаем, что товарищ дохлый, — нам надо что-то другое. То есть мы поддерживаем группу в актуальном состоянии. Вот эта первая задача, membership, на самом деле в Cassandra сейчас уже решается немного по-другому, но это пока спойлер. Пока думаем, что gossip у нас есть везде.
[1:06:58] Дмитрий: Gossip неплохо работает, но надо понимать, что если мы будем heartbeat посылать от каждой ноды к каждой, то когда у нас в кластере тысяча нод, становится как-то не очень хорошо. Все всем посылают — это уже просто перемножение.
[1:07:10] Александр: Да, у нас довольно много спама почти залетает.
[1:07:16] Дмитрий: Лишнего спама, он не нужен, столько много. Поэтому в gossip что мы можем делать? Мы можем общаться, не обязаны посылать heartbeat каждому. Это называется эпидемические алгоритмы — такой термин используется. То есть чтобы каждый заболел гриппом, не обязательно каждому на каждого чихать: я могу на тебя начихать, а ты потом на кого-нибудь другого начихаешь. Вот gossip — это то же самое: можно чихать в случайных людей. Если эти люди между собой достаточно хорошо пересекаются по чихам, то рано или поздно все начнут чихать. Там есть математика за этим стоящая, можно посмотреть скорость распространения этой инфекции.
[1:08:02] Александр: Да, инфекционный алгоритм, наверное, ещё иногда это называют.
[1:08:05] Дмитрий: И вот gossip в Cassandra — это некий такой алгоритм, который говорит: надо почихать в сторону seed-ноды, потому что это такая точка, с которой много кто наверняка общается, — такой суперраспространитель в терминах инфекции. Ну, в общем, COVID был, там прямо есть такой термин: тот, кто много общается с остальными, — суперраспространитель. Вот на него выгодно почихать, потому что он наверняка с другими будет обчихиваться. И каких-то ещё нод выбрать случайно. Но это количество мы выбираем как какое-то небольшое константное. По-моему, говорят, три. То есть даже если в кластере тысяча нод, то я пообщался с четырьмя нодами. Я не помню конкретное число, три или четыре, но в общем константным числом нод. И за счёт этого распространения инфекции по кластеру достаточно быстро все узнают про меня — что я там живой или что я перестал быть живым.
[1:08:55] Дмитрий: Это общение, причём двустороннее в gossip. То есть каждая нода не только… там, на самом деле, общение, по-моему, из трёх раундов состоит: я тебе сообщаю о своём состоянии и состоянии нод, которых я знаю, ты мне говоришь это обратно, и я, по-моему, ещё раз тебе подтверждаю, что получил этот ответ, что-то такое.
[1:09:15] Александр: Там получается, что раз в секунду выбирается жертва из какого-то подмножества, ей отправляется сообщение, которое называется GossipDigestSyn. И оно получает ответ… нет, GossipDigestAck. И на этот ack отсылается…
[1:09:39] Дмитрий: Syn-ack.
[1:09:40] Александр: Да, типа ещё один GossipDigestAck2. Это как в TCP есть трёхсторонний хендшейк — вот примерно та же история. И в итоге, чтобы обе стороны обменялись информацией. То есть это не одностороннее распространение информации, а ещё и двустороннее: если я пообщался с seed-нодой, то не только seed-нода про меня узнала, но и я кое-что из неё узнал про других.
[1:10:03] Дмитрий: Да.
[1:10:05] Александр: Инфекцию-то мы не одну распространяем, а их там много гуляет, поэтому…
[1:10:08] Дмитрий: Да, как-то так. Весь геном, накопленный к этому моменту.
[1:10:14] Александр: У нас, получается, инфекция — это хорошая аналогия в плане, если ты смотришь на такие инфекционные алгоритмы, а у нас это просто информация. Какая-то нода узнала новую информацию, что-то у неё получилось, и они обмениваются информацией. По сути, происходит обмен информацией вот таким способом.
[1:10:25] Дмитрий: Да, вирусы гуляют по сети.
[1:10:33] Александр: И это, на самом деле, очень простой для понимания, хорошо визуализируемый, прикольный такой — ну, не протокол, наверное, а способ общения. Хотя, возможно, и протоколом можно назвать.
[1:10:45] Дмитрий: Это и есть протокол.
[1:10:47] Александр: Я рад, что мы про него поговорили, gossip.
[1:10:49] Дмитрий: Да, gossip, собственно, на русский переводится как «слухи». Кто-то кому-то что-то сказал, и рано или поздно весь коллектив узнал про это.
[1:11:01] Александр: Одна бабка сказала. Оно хорошо работает в IT.
[1:11:05] Дмитрий: Вот. Итак, мы знаем, что какие-то ноды явно, скорее всего, не очень живые, а какие-то живые. Мы тут упростили — там есть ряд моментов, которые я поднимать пока не буду, чтобы не удалиться на полчаса в обсуждении этой штуки. Но мы знаем, что, возможно, какие-то ноды из наших реплик, с которыми мы общаемся, мертвы. И мы можем сделать некие предварительные подсчёты и сказать: так, ты меня просишь как клиент сделать LOCAL_QUORUM. Это значит, что мне две ноды из трёх в локальном дата-центре должны ответить. А я смотрю — по моим данным сейчас в дата-центре только одна реплика доступна, может, я сама, а соседей-то не видать, они мёртвые. Я, конечно, может, и послал бы в них запрос, но они, скорее всего, с большой вероятностью вернут мне ошибку. Что я сейчас буду заниматься этой ерундой? Я сразу тебе отвечу и скажу: извини, запись невозможна, потому что живых нод недостаточно. Это как раз одна из ошибок, которую вы как клиент получите, — что-то типа write exception, not enough — недостаточно нод для обеспечения того consistency level, что вы попросили. И это как раз та операция… вспоминаем прошлый рассказ про ретраи: это тот случай, когда ретраить можно безопасно, потому что ещё никакая запись даже не пыталась выполняться. Так что точно никто ничего не записал, и можно ретраиться спокойно.
[1:12:35] Александр: И идемпотентность как бы сохраняется.
[1:12:37] Дмитрий: Да, и это явно отдельная ошибка, которую сервер сообщит вам отдельным кодом ошибки. Вот. Но допустим, всё-таки какое-то количество реплик у нас живо, и мы можем попытаться послать в них запросы. Мы послали, и, допустим, они не ответили. Или не все из них ответили, и не собрался quorum. Вроде как все три живы, послали, а за установленное время никто не ответил. И тогда мы клиенту будем говорить: извини, не удалось собрать достаточно реплик, чтобы эту запись выполнить, — но запись мы уже сделали. Это будет другая ошибка. То есть, возможно, какие-то из этих реплик данные записали. Если мы говорим про базовые операции записи, не касаемся сейчас никаких lightweight-транзакций и прочего добра, то Cassandra в этом случае ничего откатывать не будет. Если у вас произошла такая попытка записи в какие-то из реплик, какие-то из реплик действительно могут данные записать. То есть было три реплики, одна из них ответила успешно, точно записала. Другая не ответила ничего — и вы не знаете: может быть, она записала и потеряла ответ уже на пути назад, а может быть, ответ даже не долетел, она вообще выключена и ничего не записывала. И есть третья реплика, для которой мы даже не пытались записать, потому что знаем точно, что она мёртвая.
[1:14:05] Дмитрий: У нас нет никакого механизма в Cassandra в базовых записях отката в стиле «сейчас мы всем пошлём нотификацию вторым раундом, что они все записаны, и тогда можно точно записать», — вот этот двухфазный коммит. Между этими репликами нет ничего такого: что записалось, то записалось. Поэтому, если вы получили от сервера ответ write timeout exception — мы не дождались ответа, — это не означает, что данных там нет. Это означает, что, возможно, данные там есть, а возможно — нет. Последующими чтениями вы узнаете, есть такое или нет.
[1:14:43] Дмитрий: Казалось бы, как всё плохо и так далее. Но на самом деле в обычной, некластерной системе вы можете оказаться в такой же ситуации. Вот вы общаетесь с вашим любимым, полностью консистентным Oracle: одна нода, никакого кластера. Вы сделали операцию записи, сделали коммит вашей транзакции с точки зрения клиента, и на коммит ответ вам не пришёл. Закоммитилась эта транзакция или нет — вы же не знаете. Может быть, ваш commit-request потерялся на пути к серверу, а может быть, ответ сервера о том, что коммит произошёл успешно, потерялся на пути к вам. И пока вы не сделаете следующий SELECT из базы, вы не узнаете, в каком состоянии оказалась база данных. Даже в таких простых случаях у нас возникают эти распределённые проблемы — то, что называется задачей о генералах, — и в общем случае вы никак не узнаете это, пока не спросите про состояние. Вот тут то же самое. Поэтому надо понимать, что записи в Cassandra будут происходить без откатов.
[1:15:53] Дмитрий: То есть вот эта обработка ошибки на стороне клиента от сервера во время записи играет очень важную роль. И в зависимости от того, какая ошибка там действительно есть, — мы сейчас две разобрали: когда, например, просто нет реплик, на которые координатор может отправить запросы, он на них не отправляет и говорит «sorry, not enough nodes», и мы на клиенте понимаем, что записи не было в этот момент, потому что видим семантику ошибки и можем её интерпретировать. А если у нас там «не дождался ответа от quorum», то это, вообще говоря, другая ошибка, и у неё должна быть абсолютно другая обработка, если мы хотим там записывать, например, только один раз. Нам надо будет пойти, подождать, почитать, увидеть, что реально записалось, или увидеть, что не записалось. В общем, это понимание для нас как для инженеров, которые работают с Cassandra, очень важно, потому что от этого может зависеть, вообще говоря, логика, чьи-то деньги. Чтобы потом не было такого удивительного сюрприза: вроде бы мне вернулась ошибка, а потом я пошёл и на каких-то нодах нашёл свои данные. Да, такое возможно. Какая ошибка — да, ты посмотри. Это важно.
[1:17:02] Александр: Хорошо, что мы про это поговорили.
[1:17:04] Дмитрий: И вот в нашем случае мы разговаривали — то есть мы собрали достаточно нод для quorum, отослали данные, например, в две из них, а для третьей ноды мы знаем, что она выключена. Но операцию записи мы хотим выполнить: в другие реплики послали, и мы хотим, чтобы в этой третьей реплике записи оказались. Но она сейчас выключена. Что мы будем делать? Мы хотим сделать себе пометку о том, что в эту ноду нужно отослать данные, когда она вернётся в строй, когда снова придёт к нам и скажет «я жива». Вот в тот момент мы хотим проиграть в неё эти отложенные изменения. Это механизм Cassandra — один из механизмов (в Cassandra много механизмов, которые пытаются обеспечить консистентность хотя бы в конечном счёте). Это, по-моему, называется hinted handoff в Cassandra.
[1:17:56] Александр: Да, это именно он. И это круто. То есть, на самом деле, мы себе пишем некий файлик. Каждая Cassandra — на координаторе.
[1:18:06] Дмитрий: Ну да, поскольку каждая нода — координатор, то каждая нода в принципе таким занимается.
[1:18:10] Александр: У ноды в роли координатора в конкретном запросе появляется файлик, который хранит данные для какой-то другой ноды. Этот координатор по логике нашего consistent hashing не должен эти данные хранить, но он хранит для того, чтобы потом их отослать и как бы дозаписать данные условно. Блин, это очень сильная идея, на самом деле. Да, то есть это получается мини-коммит-лог, который я в отложенном режиме проигрываю в другую ноду.
[1:18:39] Дмитрий: Этот файлик там лежит, естественно, per-node — поскольку каждой ноде мы можем в свои записи посылать, всё шардировано, поэтому файлик per-node. Дальше он начинает проигрываться — понятно, там всякие rate limiter’ы есть и так далее, чтобы не убить другую ноду, — но потом, eventually, мы туда проиграем. А теперь возвращаемся к нашему разговору про consistency level. Помнишь, я упоминал волшебный уровень ANY, который не ONE, а ANY, он более слабый, чем ONE? ANY означает, что если ни одной реплики нет, то запиши для них хинты и ответь, что всё прошло успешно. То есть мы реально ни в кого ничего не записали — мы не записали ни в одну из реплик, но заперсистили намерение записать на координаторе. Этот ANY — это я сам как координатор, но при этом чтением ты это не увидишь: пока оно не проиграется, ты не увидишь это изменение. Вот этот волшебный consistency level ANY, самый слабый, как раз означает: пиши в хинты, проиграется потом, а клиенту ответь, что куда-то мы сохранили. Операция write успешна для вот этого слабенького consistency level.
[1:19:48] Дмитрий: Итак, мы поговорили про логику координатора. У нас реплики пока были чёрными ящиками. И тут, наверное, основные моменты мы все, пожалуй, разобрали. Понятно, что есть ещё всякие вещи, но, думаю, можем спускаться на уровень реплик.
[1:20:07] Александр: Идём к репликам, погнали.
[1:20:09] Дмитрий: Да. Соответственно, мы посылаем запрос в реплику. Что там дальше происходит? А дальше нам нужно повторить всё то же самое, что мы говорили для клиента и сервера. В нашем случае опять у нас есть нода Cassandra, общающаяся с другой нодой Cassandra, — у нас есть клиент-серверное взаимодействие. Реквест посылается с этой стороны, получается на той стороне, опять Netty, опять некий протокол, но протокол уже другой. Это протокол межсерверного взаимодействия, и там формат общения идёт уже не в CQL, а в некотором внутреннем представлении. Там тоже некий бинарный формат, где этот запрос закодирован; в случае записи эта сущность называется mutation. То есть Cassandra называет вашу запись в своей терминологии мутацией. Это INSERT, DELETE, UPDATE — что угодно из этого мутация. Это сериализованное представление уже разобранного, разложенного запроса с некими внутренними метаданными. Оно сериализуется в какой-то бинарный формат — там есть версионирование этого формата, чтобы, как раз то, что мы обсуждали, разные версии сервера могли между собой общаться, — и пишется сквозь отдельный Netty-стек в другой сервер. И это будет другой порт. На самом деле сервер Cassandra открывает два основных TCP-порта: 9042 — это порт, которым он общается с клиентами, и 7000 по умолчанию — порт для межсерверного взаимодействия, где нода Cassandra с нодой Cassandra общаются по этому другому протоколу. Там же gossip подлетает — всё это там происходит.
[1:21:44] Дмитрий: Мы прошли сквозь этот стек Netty, попали на клиентскую сторону. Тут всё то же самое применимо, на самом деле, про rate limiting, что мы с тобой обсуждали. Потому что мы не хотим, чтобы один сервер перегрузил другой сервер. Поэтому механизм с ограничением количества запросов также сделан и здесь — эти алгоритмы переиспользуются. Прибежали на клиент… прибежали на реплику, теперь начинается операция локальной записи. И тут мы уже должны поговорить про структуру данных, которая при этом используется. Первый шаг чем-то похож на реляционную базу данных: для начала мы хотим всё это сохранить в относительно персистентном виде, и мы запишем данные в commit log. Мы получили мутацию, и первое, что с ней делаем, — записываем в commit log.
[1:22:16] Александр: Да.
[1:22:39] Дмитрий: Это некий локальный ротируемый набор файликов. В этом плане Cassandra немного отличается, наверное, от Postgres и так далее: тут commit log — набор фиксированный, ограниченный по размеру. У нас есть суммарный объём диска, который мы под этот commit log хотим отвести. Когда он заканчивается, мы те другие структуры данных, куда мы пишем основные данные (о которых сейчас будем разговаривать), сбрасываем на диск, убеждаемся, что они сброшены, и после этого можем commit log заротировать. За счёт этого достигается ограничение по размеру на суммарный commit log. То есть мы форсируем сбросы. Интересный момент: как и в случае с in-flight реквестами, где мы ограничивались не по количеству, а по общему суммарному размеру сообщений в памяти, — так и здесь, в случае с commit log, мы тоже указываем размер.
[1:23:50] Александр: Мне кажется, это такая сквозящая сквозь идею Cassandra штука — что мы хотим ограничивать именно размер. Прав ли я?
[1:23:59] Дмитрий: Ну, я не могу сказать, что она прямо везде сквозная, но это безопасная идея в плане операционном. Если кто-то не смотрит за сервером, не получится так, что ты очень легко выел там всё место или всю память. То есть такой естественный защитный механизм, который более устойчив к проблемам и к неожиданному росту нагрузки, — проще для этого сайзинг делать. Вот в этом плане да.
[1:24:27] Александр: Окей. Ротируемые commit log файлы. Вот мы их ротируем.
[1:24:33] Дмитрий: Мы туда записываем некое тоже бинарное представление. На самом деле, если посмотреть кишочки, то это то же самое сериализованное представление, что мы отправляли между серверами. То есть эта мутация сериализуется в байтики и кладётся в append-only файл. И тут мы можем сделать оптимизацию: не два раза сериализовать — чтобы в сокет отправить, а чтобы записать в commit log, — а один раз.
[1:25:03] Александр: Сразу, напрямую.
[1:25:05] Дмитрий: Да, записать в byte buffer, а потом заслать его и туда, и туда. Есть zero abstraction идея в C++, в Rust — вот мы там zero serialization делаем в Cassandra.
[1:25:13] Александр: Zero serialization, точнее? Сериализация.
[1:25:17] Дмитрий: Сериализация.
[1:25:19] Александр: А, ну подожди, мы же представление из памяти хотим представить в виде набора байтиков, чтобы записать в сокет, чтобы записать на диск. И мы прямую сериализацию…
[1:25:31] Дмитрий: Тут скорее, знаешь, наверное, что ещё важно. Важно отметить, что это один и тот же способ интерпретации этих байтиков, что и в сети, то и на диске. Поэтому нам в целом-то и не надо их туда-сюда гонять, для commit log. Дальше будут другие представления на диске. Вот, мы пишем commit log.
[1:25:52] Дмитрий: Тут важный момент. Cassandra спроектирована изначально как распределённая база данных, и она всё-таки рассчитывает, что у вас реплик много, и выбирает определённый trade-off: надёжность против производительности. В частности, в точке с commit log. Потому что классический commit log в какой-нибудь реляционной базе данных по умолчанию пишется синхронно на диск — у вас есть fsync. Если вы записали что-то в commit log, и база получила от commit log acknowledgement, что запись произошла, мы уверены, что данные попали на диск, что они были сброшены на диск, а не застряли где-то в буфере операционной системы в памяти. Мы для этих данных делали операцию fsync. Это не значит, что мы на каждый чих должны делать fsync — мы можем группировать разные операции вместе, — но когда мы сказали, что операция записи в commit log прошла успешно, то для нашей записи в том числе fsync был сделан.
[1:26:53] Александр: Да, это очень важно. Это очень важно для понимания инженеров: в какой именно момент времени мы говорим, что данные — по крайней мере на уровне реплики, в которую мы сейчас поместили нас и наших слушателей (мы сейчас сидим в реплике, в которой координатор сказал «записывай», и мы этот запрос на запись, мутацию в терминах, увидели), — этот commit log, на всякий случай повторю, находится на стороне реплик. То есть мы заслали запрос в три реплики, и каждая из реплик сама независимо пишет сначала в commit log. Он не какой-то общий, он не на координаторе.
[1:28:04] Дмитрий: Это, вообще говоря, можно рассматривать ментально — вот мне так проще — как один файл на диске у каждой реплики, грубо говоря, commit log: просто они файлик пишут. И вот до тех пор — даже когда у тебя там открыт стрим, ты туда напихал байтов любым API, работая с файлами, ты типа записал, — но до тех пор, пока ты не вызвал fsync (причём fsync тоже разный бывает) и не получил от операционной системы «ок», — вот только этот «ок» ты можешь переслать обратно координатору и сказать: я записал. Вот только в этот момент. До этого считается, что ты ещё не записал.
[1:28:04] Александр: То есть мы не можем себе позволить, просто получив реквест на запись, сразу отослать респонс и асинхронно писать на диск.
[1:28:16] Дмитрий: Да, можем, но это уже не те гарантии. И в классических базах данных эти гарантии стараются соблюдать, потому что нода-то одна: если она потеряла — то больше ничего нигде нет. В Cassandra мы говорим, что вы деплоитесь всё-таки в определённом окружении, у вас реплик много, поэтому мы скорее будем рассчитывать на то, что все реплики не выпадут сразу, и не хотим по умолчанию делать fsync, потому что это дорогая вещь, она сильно нас тормозит, ограничивает нашу пропускную способность как системы. Поэтому по умолчанию мы его не делаем синхронным. Вы можете настроить его в Cassandra синхронным, но это не дефолтная настройка. По дефолту мы говорим, что делаем его регулярно, по умолчанию каждые 10 секунд. Каждые 10 секунд Cassandra будет fsync‘ать всё, что там накопилось, но не обещает вам синхронизацию каждого входящего запроса. Это trade-off. Понятно, что у вас есть риск потери, но мы говорим, что поскольку реплик много, мы рассчитываем, что этот риск минимизируется за счёт репликации.
[1:29:00] Александр: Да, и это trade-off, который в computer science, в нашей сфере, постоянно везде — мы всегда, по сути, балансируем между trade-off’ами. У нас есть жёсткие ограничения внешней среды: мы не можем обрабатывать миллионы записей в секунду как система, делая на каждую запись fsync. Это просто физически невозможно, потому что fsync — долгая операция. То есть нам либо надо как-то батчевать, и тогда latency нужно подкручивать. Принимая каждое решение, мы просто балансируем между набором ограничений, которые у нас есть. Мы, по сути, не можем обрести вселенную и сказать, что мы вот такие быстрые и надёжные. Ну, не бывает. Выбери, как говорится, два из трёх. И Cassandra по дефолту выбирает, что мы будем быстрые-быстрые.
[1:30:08] Дмитрий: Cassandra выбирает скорость, очевидно, на запись в своей архитектуре. Но мы нивелируем негативные стороны выбранного trade-off тем, что нас много, и статистически, скорее всего, плохой ситуации отказа глобально не наступит. И данные всё-таки запишутся.
[1:30:29] Александр: Здесь как бы в эту систему ограничений вносится ещё и статистика со стороны Cassandra. Это тоже интересно, такое наблюдение. Я вот для себя сейчас его сформулировал.
[1:30:41] Дмитрий: Ну, если вы прямо страховщик, вы можете, конечно, включить этот синхронный режим fsync через конфигурацию, но не удивляйтесь, что производительность у вас упадёт, скорее всего, на порядок.
[1:30:52] Александр: Да.
[1:30:54] Дмитрий: Мы записали в наш commit log, и дальше мы пишем в основную структуру данных, из которой мы, собственно, потом можем читать. И тут Cassandra идёт от концепции: мы хотим избежать случайных записей на дисках. Как мы можем это сделать? Допустим, мы хотим, чтобы на диск мы писали…
[1:31:18] Александр: Случайная запись на диск — ты говоришь про random seek, правильно?
[1:31:22] Дмитрий: Random write.
[1:31:24] Александр: Ну, seek — это уже термин вращающихся дисков. То есть во вращающихся дисках вот это есть. Или, да, в файлах всегда. Я понял, в чём ты.
[1:31:33] Дмитрий: То есть, особенно для вращающихся дисков… Cassandra изобреталась, когда HDD были ещё популярны — они и по-прежнему есть. Да и SSD тоже не то чтобы прямо очень счастливо обслуживают ваши случайные записи. В общем, мы бы хотели писать большими пачками, последовательно, избегать случайных записей. Тут мы можем пойти от такой задачи и попытаться понять, какую структуру данных для этого нам стоит сделать. Если мы хотим что-то писать большими пачками, наверное, мы должны это всё копить где-то в памяти, в каком-то буфере, а потом периодически сбрасывать на диск.
[1:32:15] Александр: Ну, если у тебя есть пачка, она должна как-то собраться. Она не приходит откуда-то сразу. Нам клиенты пишут не пачками, у нас запросы.
[1:32:26] Дмитрий: Да, они пишут по одному запросу, мы где-то в памяти это всё агрегируем, батчуем, а потом рано или поздно память кончается, да и мы хотим переживать выключение и падение, поэтому нам нужно сбрасывать всё это на диск. И когда мы сбрасываем эту пачку на диск, мы её отпишем целиком, последовательно, — вот она, и нет никакой случайной записи. Сейчас мы, не углубляясь в детали, уже погрузились в область рядом с LSM-деревьями, но на самом деле есть и другие алгоритмы хранения данных, которые тоже оперируют такими же идеями.
[1:32:41] Александр: Первое упоминание LSM-дерева в выпуске про Cassandra случилось на втором часу второй части.
[1:32:52] Дмитрий: Да. То есть мы сформулировали, что хотим и почему. Но дальше мы хотим ещё из этого делать чтение, поэтому просто в буфере копить это…
[1:33:25] Александр: Буфер, который копит только для того, чтобы писать, — это только что мы сделали. Commit log был таким буфером. Что нам тут ещё что-то делать? Вот же мы всё записали, только батчевания не было, но давайте в памяти сделаем — и всё.
[1:33:39] Дмитрий: Проблема в том, что мы потом ещё данные хотим читать, поэтому хотелось бы не читать их в формате «перебери все записи в памяти в буфере», а как-то побыстрее. А когда мы хотим данные читать, нам логично сделать какую-нибудь индексируемую структуру данных в памяти. Самая популярная структура данных — это некие мапы, хэш-мапы и так далее. И тут мы приходим… Я, наверное, попытаюсь прийти сразу к результату, чтобы слишком не затягивать рассказ о предпосылках. Мы приходим к следующей иерархической структуре данных. Cassandra — это мапа, мапа, мапа.
[1:34:21] Александр: Вау!
[1:34:23] Дмитрий: Да. На верхнем уровне у нас первая мапа — это по partition-ключу. В ноду Cassandra по-прежнему прилетает много partition-ключей, поэтому нам всё равно надо хранить для каждого partition-ключа что-то. Вот это первый уровень мапы. При этом эта мапа должна быть конкурентной. У Cassandra не thread-per-core архитектура, не Scylla: несколько потоков могут одновременно делать операции записи с одной структурой, вообще говоря. Поэтому мы хотим сделать некую concurrent-мапу. Есть две реализации здесь. Какие ты знаешь concurrent-мапы в Java и Scala?
[1:35:11] Александр: Так, ConcurrentHashMap. А ещё какой-нибудь? Блин, ты мне сейчас задал… Подожди, идею.
[1:35:16] Дмитрий: Смотри, давай подсказку сделаю. Мне для дальнейших нужд — я потом объясню зачем — хорошо бы по partition-ключам ещё уметь сортировать. То есть у меня partition-ключи comparable. Потом, когда я буду писать данные на диск или работать с большими объёмами данных, удобно, когда они отсортированы.
[1:35:37] Александр: Ну тогда TreeMap должна быть. Но concurrent.
[1:35:40] Дмитрий: Concurrent, yes.
[1:35:42] Александр: В общем…
[1:35:44] Дмитрий: Стандартная JDK-реализация для concurrent-мап — это ConcurrentSkipListMap.
[1:35:51] Александр: Skip list, у меня правильно занят, да.
[1:35:53] Дмитрий: Собственно, вот один из вариантов, как это сделать, который в Cassandra был до недавних относительно пор единственным. Исторически были другие варианты, но они уже отпали. Это ConcurrentSkipListMap на первом уровне. Потом появился второй вариант, где была сделана следующая идея: вот эти partition-ключи очень часто очень похожи, у них могут быть одинаковые начальные кусочки, и мы бы не хотели тратить ресурсы на то, чтобы хранить их по отдельности. У них одинаковые префиксы, и поэтому мы можем использовать другой тип мапы, который называется — уже стандартный, к сожалению, в Java нет такого — префиксная мапа или префиксное дерево. Или в английской литературе чаще всего называют trie.
[1:36:44] Александр: Да, trie.
[1:36:45] Дмитрий: T-r-i-e. Или по-русски префиксное дерево. То есть мы делаем некое дерево: вот, например, у меня есть два слова, которые начинаются одинаково, а заканчиваются по-разному. Одинаковую часть я выношу в общую ноду этого дерева, а суффиксы, которые отличаются, — это его дети. И так далее. Я строю эту структуру данных итерационно, у меня получается всё больше не уровней, но общие части слов я храню один раз в этих общих нодах. Этот алгоритм, понятно, в более сложном виде был реализован в Cassandra. Собственно, сейчас в версии 5.0, в релизной версии, он доступен, можно им пользоваться. И поскольку такое дерево в полноценном concurrent-варианте строить очень тяжело, она сделана шардированной в памяти. То есть она шардирована — это не те же самые шарды, что в Cassandra, это просто нарезано на кусочки, и каждый кусочек мы закрываем блокировкой: только один поток в один момент времени может писать этот шард. Но поскольку их много, мы тем самым потоки наши параллелим. Если мы считаем, что ключи размазаны относительно равномерно, то мы можем эти шарды обновлять независимо. И так у нас на верхнем уровне либо ConcurrentSkipListMap, либо префиксное дерево. А в качестве значений там лежат партиции.
[1:38:14] Александр: Так, это первая мапа.
[1:38:16] Дмитрий: Да. И вот из этой трёхуровневой мапы мы раскрутили первый слой. Дальше второй уровень — это вторая часть, clustering-ключ. У нас был partition-ключ, clustering-ключ. Верхний уровень мапы приводит нас в партицию на конкретную ноду. Мы сейчас внутри конкретной ноды сидим, именно ноды Cassandra. Внутри неё хранится какой-то набор партиций, ключом которых является partition-ключ. А внутри партиции у нас находится набор строчек, индексированных clustering-ключом. Собственно, вот второй уровень мапы: ключом является clustering-ключ, а значением — значение строчки, то есть колонки и их значения. И эта мапа, как мы обсуждали при создании таблицы, — clustering-ключ это тоже отсортированный ключ. То есть мы снова хотим свойство сортировки, только здесь уже по пользовательскому типу, не по нашему внутреннему partition-ключу, который для наших нужд сделан, а здесь мы сортируем упорядоченно с точки зрения пользователя. Это чтобы ему потом было удобно делать запросы, потому что мы хотим, чтобы он мог вытаскивать диапазоны данных внутри партиции, — поэтому нам удобно хранить их в отсортированном виде.
[1:38:56] Дмитрий: Вот здесь второй уровень мапы. И здесь я тебя напугаю, потому что здесь тоже нестандартная коллекция используется. Здесь не TreeMap и даже не skip list map. Здесь нам нужно обеспечить некое свойство консистентности. Cassandra, конечно, база данных нового поколения, и она по консистентности относительно расслаблена по сравнению с реляционными базами данных, но всё-таки какую-то консистентность она хочет соблюдать. В частности, когда ты делаешь вставку в Cassandra, ты можешь в своём батче вставить несколько строчек в одну партицию. У тебя partition-ключ одинаковый, clustering-ключи разные, и это всё склеено в один batch-запрос.
[1:40:35] Дмитрий: И хотелось бы по-хорошему, чтобы эти штуки с точки зрения читателя либо появились все вместе, либо мы их не видели вообще, — то есть не видели промежуточной записи. А уж тем более такого бы мы не хотели для отдельных колонок: у тебя одна строчка, и колонки внутри этой строчки по отдельности появляются с точки зрения читателя — что-то там обновилось, а что-то нет. Ты хочешь некий такой либо старое целиком представление, либо новое, но не микс.
[1:41:03] Александр: У меня есть пример консистентности. Например, я записываю в одном батче, отправляю типа INSERT‘ы в какие-то данные — не знаю, в корзине. И я это делаю одним запросом. И я хочу сумму по корзине иметь либо нулевую, либо сразу 10. Я не хочу, чтобы у меня с каждым запросом там 1, 2, 3, 3, 3, 5, 10 произошло. То есть это консистентность, про которую мы говорим.
[1:41:33] Дмитрий: Да. У меня, кстати, как раз похожий случай из рабочих примеров есть, где у меня некие дельты обновления некого счётчика или ещё чего-то, которые я храню по отдельности, есть в той же партиции summary — некий суммарный баланс. Я хочу эти дельты переносить атомарно из одной части в другую: убрать дельту и добавить на финальную сумму. Чем-то похоже на твой пример с корзиной. Я хочу, чтобы я убрал здесь, положил там, и это была атомарная операция, а не так, что я здесь убрал, а ты видишь, что я убрал, а положить то, что я положил, ты не видишь ещё из-за того, что не успел прочитать какую-то новую версию объекта. И нам должно получаться здесь некое такое упрощённое версионирование. В реляционных базах данных это через всякие локи, MVCC — multiversion concurrency control — решается. Здесь мы хотим упрощённый вариант чего-то подобного. И нам нужна некая мапа, которая как бы давала возможность прочитать старый вариант, а потом переключиться целиком на новый, и при этом общие части старого и нового варианта мы бы не хотели полностью переписывать. То есть я не хочу весь этот clustering-блок, мою партицию целиком каждый раз переписывать, создавать с нуля, потому что там уже тысячи строк будешь хранить. И такая перезапись — дорогое удовольствие. Поэтому мне нужна структура, которая позволяет такое версионирование сделать. И в Cassandra сейчас реализовано там B-дерево.
[1:42:15] Александр: Слушай, если бы ты мне предложил выбрать…
[1:43:24] Дмитрий: В памяти.
[1:43:24] Александр: Господи, я бы не угадал. Я не угадаю ещё, но это B-tree.
[1:43:33] Дмитрий: Так, ну и интересно.
[1:43:35] Александр: Оно сортированное.
[1:43:37] Дмитрий: Так, да, это факт. Оно хорошо скейлится под большие объёмы данных в том плане, что… Как это? B-дерево изначально появилось в дисковых структурах данных, потому что мы хотели минимизировать количество операций с дисков. А теперь, как говорится, память — это новый диск.
[1:43:58] Александр: Да.
[1:43:59] Дмитрий: Представь в качестве цитаты. По-моему, много кто такое говорил, поэтому авторство точно не моё. С точки зрения… Разрыв между скоростью ЦПУ и скоростью памяти большой, и он скорее только растёт, чем уменьшается. Соответственно, когда мы работаем со структурой данных даже в памяти, нам лучше бы поменьше читать из памяти. То есть у нас раньше было память и диск, а теперь это память и кэши процессора — та же самая структура. Нам бы побольше вытягивать в кэши на уровне ЦПУ и поменьше работать с основной памятью.
[1:44:43] Александр: И самое интересное, что решение этой проблемы одно и то же. Даже не говоря про B-tree, а про идею страниц, которая использовалась сначала для дисков, и теперь также мы со страницами работаем просто в оперативной памяти, и эти же страницы к нам попадают в кэши потом.
[1:44:58] Дмитрий: Да, там concurrent-версии единой страницы пока нет, по крайней мере в текущей версии. Но некая группировка записей, такие вот блоки, есть.
[1:45:07] Александр: Блочная структура такая.
[1:45:09] Дмитрий: Да, такая блочная структура. То есть дерево не бинарное, а такое сильно ветвистое. Соответственно, каждый блок, каждый кусочек даёт больше информации при работе с ним. Это раз. Точнее, это два. Она сортированная — нашу задачу решает. Она вот такая CPU-cache-friendly.
[1:45:33] Александр: Cache-friendly, да.
[1:45:34] Дмитрий: Она достаточно масштабируемая в плане размеров. Мы могли бы взять сортированный массив, но сортированный массив хорошо работает, когда структура маленькая; как только у тебя там тысячи элементов, то каждая вставка в него — это перекопирование этого массива целиком, чтобы воткнуть что-то, вставить свободное место. По количеству строчек он не очень хорошо масштабируется, поэтому в этом плане B-дерево умеет раздуваться, если надо. И в этой древовидной структуре как раз можно обеспечить механизм версионирования, по крайней мере в той версии, которую реализовали в Cassandra: те блоки, которые не меняются, шарятся, а те блоки, которые меняются, как раз приклеиваются…
[1:46:26] Александр: У тебя, допустим… Пинится? Или нет?
[1:46:29] Дмитрий: У тебя было старое дерево, ты в него что-то записал, те блоки, которые поменялись, ты скопировал, а те, что не поменялись, оставил в старом дереве. То есть у тебя поверх старого дерева наслоилось такое новое дерево.
[1:46:44] Александр: Ага.
[1:46:44] Дмитрий: Причём в какой-то момент времени, за счёт того что это, по крайней мере сейчас, всё живёт в хипе, вот эти старые блоки, которые больше никто уже не читает, не использует, они самоотмирают, потому что на них никто в ссылках не хранит. Их garbage collector подбирает. То есть тут мы ленимся и не собираем их самостоятельно. Но если ссылка на него есть — например, поток, который читает, — вот он будет читать свою старую версию целиком этого дерева, поскольку обновление идёт с головы этого дерева. Мы сейчас совсем в детали этого алгоритма уходить не будем, но, в общем, дерево форкается от головы, и в нём меняются только те блоки, которые по пути от головы до листьев, которые поменялись, копируются. Всё остальное остаётся тем же самым. За счёт этого мы относительно дёшево можем сделать вот это multiversion-дерево с точки зрения взаимодействия потоков чтения и записи. Вот это второй уровень. Это чудо-B-дерево.
[1:47:46] Александр: Так, то есть смотри, получается, второй уровень мы разобрали, чудо-B-дерево. И получается, что Cassandra это как бы не map of map, а ConcurrentSkipListMap of map. Либо trie of B-tree that holds… что?
[1:48:05] Дмитрий: Вот, а дальше у нас последний, собственно, уровень мапы. Мы уже по классическому ключу определили конкретную строчку. Но строчка-то у нас — это не просто байтики, это набор колонок со значениями. Там же много колонок, может быть, дофига колонок, особенно если какая-нибудь сложная структура данных типа коллекции — с точки зрения коллекции каждый элемент это такая мини-колоночка. Для каждой колонки нам надо хранить некие метаданные. И ты не обязан вставлять все значения сразу: можешь сначала половину колонок вставить, потом другую половину, а потом обновляешь что-нибудь — вообще только одну колонку обновил из 50. Поэтому это ещё сущность, которая меняется. И здесь колонки нам снова удобно с точки зрения обновления хранить в отсортированном виде, чтобы их мержить было можно. Старую версию колонок и новую версию колонок мы хотим уметь мержить, а мержить удобно, когда обе структуры данных отсортированы, чтобы ты мог сделать один проход сортировкой слиянием — то есть смержить блоки, которые заранее отсортированы, просто двигаясь по ним параллельно. Поэтому тебе удобно, чтобы эта структура была отсортирована. Она должна быть индексированной, потому что ты хочешь по конкретным колонкам уметь что-то доставать: ты говоришь «хочу достать вот эти колонки» и хочешь из этого множества колонок эффективно заселектировать что-то. То есть по имени колонки я хочу взять её значение в строке. Поэтому это, собственно, мапа. Снова.
[1:48:57] Александр: Вот это третий уровень, потому что ключ — это название… Слово «сортинг» мне подсказывает, что это TreeMap всё-таки.
[1:48:57] Дмитрий: Вот она сортит, она мапа, но ты не угадал. Тут идея такая: мы уже реализовали сортированную мапу. Зачем ещё одну городить? Поэтому ещё одно B-дерево и всё.
[1:49:59] Александр: О как! Блин.
[1:50:02] Дмитрий: Чтобы тебя так сильно не пугать, на самом деле, если распарсить B-дерево в упрощённом случае, когда колонок не очень много — там, я не помню, элементов в этом B-дереве когда мало, до 16, я не помню конкретную константу, — то это B-дерево, состоящее всего из одного уровня, а что внутри лежит, по сути дела сортированный массив. То есть по сути там внутри просто сортированный массив. Не надо совсем уж бояться: в простых случаях, когда ты не награбил кучу колонок, ты по сути дела работаешь с сортированным массивом. Сортированный массив в случае этой реализации Cassandra — это, можно сказать, частный случай B-дерева, вырожденного, состоящего из одного уровня. Поэтому не всё так страшно. Вот она, собственно, мапа, мапа, мапа.
[1:50:54] Александр: Получается, смотри, мы разобрали прямо по структурам данных. И особо талантливые…
[1:50:59] Дмитрий: В памяти.
[1:51:01] Александр: Особо талантливые и способные слушатели могут даже, наверное, это себе представить аккуратненько, с хорошим воображением. Но в целом понятно.
[1:51:11] Дмитрий: И это то, что у нас по сути называется memtable. Или нет?
[1:51:17] Дмитрий: Да, да, собственно, то, что мы только что проговорили, — эта структура данных в совокупности в памяти называется memtable, и у неё могут быть разные реализации.
[1:51:26] Александр: Да, и, ребят, кто немножечко, может быть, чувствует, что плывёт или что недостаточно знаний, базы, — у меня целый выпуск про LSM-деревья есть, «Тысяча фичей», LSM-3, гуглите, — и там как раз эта часть просто бинарным деревом. Я говорю: бинарное какое-то дерево поиска, которое позволяет вот так находить. А по факту, видите, в реальных системах получается, что там трёхуровневая структура данных, причём не самая тривиальная — на каждом уровне используется структура. Супер интересно. Но мы сейчас рассмотрели то, что в памяти.
[1:51:57] Дмитрий: В памяти, в оперативной, CPU-cache-friendly такая штука.
[1:51:59] Александр: Что с ней происходит? Что у нас дальше?
[1:52:07] Дмитрий: Да, собственно, в принципе, вот мы в неё вставляем. Вспоминаем, что мы говорим про операцию вставки: естественно, мы идём в первый уровень, если там партиции ещё нет по этому ключу — создаём партицию, если есть — вставляем такой putIfAbsent. Ну, putIfAbsent или там… что получилось, если она уже была. Дальше берём эту вторую структуру данных, B-дерево, и в неё вмерживаем нашу строчку. Мы берём существующее B-дерево, берём нашу мутацию, которая, на самом деле я не сказал, тоже B-дерево.
[1:52:42] Александр: Мутация?
[1:52:44] Дмитрий: Мы незаметно, на самом деле, поработали уже с B-деревом, когда создали объект mutation, потому что там надо тоже эти разобранные данные разложить в сортированном формате, удобном для слияния. А удобно сливать два дерева. Опять-таки нам лень писать ещё одну сортированную коллекцию — возьмём просто существующее B-дерево, благо написано неплохо. И мы должны слить два B-дерева, одно на другое наложить. И получившийся результат записать как новое значение этой партиции.
[1:53:17] Александр: В памяти всё ещё пока что.
[1:53:19] Дмитрий: Да, всё ещё в памяти. И вот здесь происходит атомарный свитч, на самом деле.
[1:53:22] Александр: Atomic swap.
[1:53:25] Дмитрий: Да. В ConcurrentSkipListMap там по сути делается atomic reference, перещёлкивается. В trie-дереве там лок на каждый шард, собственно, под этим локом всё это выполняется. И мы замещаем в этом дереве нашу структуру данных с новым значением. И тут второй важный, другой важный момент, отличающий Cassandra от других баз данных. Там надо помнить, что иногда можно попасться на всём этом. Как мы мержим эти данные? Вот у тебя есть старая строчка, что там уже была с какими-то значениями, — какие-то колонки с какими-то значениями. И новая строчка.
[1:54:11] Александр: Так, старая, новая.
[1:54:13] Дмитрий: Можно замержить по принципу «просто всё, что новое прилетело, то и хорошо». Если старого значения не было, а новое прилетело — понятно, что делать, берём новое. Но вот что делать, если старое значение было и новое значение было? Я сказал уже «старое» и «новое», но на самом деле в распределённой системе это особенно нетривиальные понятия, потому что мы затронули такую тонкую материю, как время.
[1:54:33] Александр: О, боже.
[1:54:38] Дмитрий: Потому что, может быть, ты что-то делал, и я в этот момент что-то делал. И кто из нас раньше, а кто позже — это ещё вопрос. А когда мы ещё говорим про запись в разные реплики (мы же пишем в три разные реплики), нам хотелось бы, чтобы на каждой реплике получились одинаковые результаты, а не так, что на первой реплике твоя запись прилетела первой. Вот мы с тобой два человека, которые пишем в одну и ту же строчку примерно в одно и то же время. Наши запросы прилетели в два разных координатора, а может, даже в один. А дальше координаторы распылили на эти реплики. На каждую из трёх реплик полетели твоя запись и моя запись. А дальше они же никак не упорядочены между собой, нет какого-то глобального упорядочителя. То есть у тебя на первой реплике могла прилететь первой твоя операция, а на второй — моя. То есть нет какого-то глобального счётчика, который между нами синхронизирован и на который мы можем ориентироваться.
[1:55:37] Александр: Да, то есть вот этот глобальный счётчик — это задача уровня Paxos, консенсуса и так далее в распределённой системе. Это как глобальные транзакции реализуются через такой, называется, один из вариантов — секвенсор, когда есть некий координатор, который упорядочивает твои транзакции глобально в некую линейную систему. Это определённый порядок, и за счёт этого обеспечиваются все распределённые транзакции.
[1:56:01] Дмитрий: Вот здесь мы этого оверхеда не хотим. Мы про базовые операции записи говорим. Но мы хотим, чтобы каждая реплика в итоге получила один и тот же результат. А здесь мы только за счёт того, что наши запросы в плане конкуренции прилетели в разном порядке, получаем, что реплики уже разъехались. Поэтому просто по принципу «кто первый прилетел, того и положим» не работает. Нам нужен некий порядок. В Cassandra выбран один из самых простых вариантов. Возвращаясь в самое-самое начало: когда ты как клиент делал INSERT в таблицу, ты на самом деле явно или неявно указываешь таймстемп этой операции. То есть у тебя на этом INSERT написано время записи.
[1:56:52] Александр: Локальное для клиента.
[1:56:54] Дмитрий: Да. Ты можешь указать его явно. Ты можешь указать алгоритм его вычисления — это называется таймстемп-генератор или как-то так в клиентском коде Cassandra. То есть ты можешь положить этот таймстемп, и он полетит на сервер. Ты можешь сказать «используй серверный таймстемп», но это не очень хороший подход по ряду причин. Отсюда следствие. Когда мы работаем с Cassandra, вы часто видите сразу, что часы должны быть синхронизированы. В определённых системах часы желательно бы синхронизировать, а в Cassandra это особенно важно. Часы надо синхронизировать. Но более того, речь идёт не только про серверы Cassandra, но и про клиентов тоже. Ваши приложения тоже должны иметь синхронизированные часы, потому что таймстемпы летят от ваших приложений.
[1:57:31] Александр: Вот это, конечно, задачка. Да, такая уже…
[1:57:52] Дмитрий: Задачка скорее в текущих условиях айтишная. То есть у вас должны стоять соответствующие настроенные пакеты с демонами, которые синхронизируют время.
[1:58:03] Александр: У тебя должен быть контроль над клиентами.
[1:58:07] Дмитрий: Да, если клиенты хотят получить что-то вменяемое.
[1:58:12] Александр: Ну да. Интересно. Часы синхронизированы. Никаких там гибридных таймстемпов не используется?
[1:58:18] Дмитрий: Прямо используется таймстемп Unix. В классических записях — нет. Мы рассчитываем на синхронизацию, на упорядоченность по времени. Во всяких транзакционных вариантах логики там своя история — появляются эти логические часы Лампорта в том или ином виде. Но мы говорим про обычные записи. Соответственно, когда ты пишешь, у тебя по факту есть таймстемп того, что ты записываешь. Более того, с точки зрения Cassandra таймстемп для каждой ячейки, вообще говоря, свой. Потому что ты мог сделать INSERT для набора колонок, а потом обновил одну из этих колонок — у этой обновлённой колонки таймстемп должен сохраниться другой. Поэтому ей надо хранить таймстемпы для каждой ячейки данных. Вот в этом memtable у каждой ячейки данных свой таймстемп. И когда ты делаешь мерж, ты не просто старое заменяешь новым — с точки зрения «то, что было, это старое, а то, что пришло, это новое». Ты мержишь по таймстемпу. У тебя сравниваются таймстемпы: то, что было уже в memtable, и то, что ты принёс. Если то, что ты принёс, старее, чем то, что записано в memtable, победит то, что записано в memtable. Это как бы один из базовых краеугольных камней Cassandra — идея last write wins по таймстемпу. Она там через разные уровни проходит, но здесь то же самое. За счёт этого, даже если запросы прилетят в разном порядке на эти реплики, результат мержа будет консистентным. В каждой реплике в итоге получится одно и то же, несмотря на то что переставились вот эти операции, потому что есть линейный порядок по времени.
[2:00:03] Александр: Понятно, понятно. Это важный момент.
[2:00:07] Дмитрий: Логика мержа дерева, точнее, мапы дерева, понятна плюс-минус, про которую происходит в оперативной памяти. Вот мы их подмержили. Мы пишем, пишем, пишем, рано или поздно память-то… Память у нас закончится.
[2:00:17] Александр: Мы не можем бесконечно писать.
[2:00:23] Дмитрий: Поэтому мы должны всё это сбросить на диск, чтобы потом было ещё удобно читать, но не только. У нас есть эта мапа-мапа-мапа в памяти. И вспоминаем, что у нас по первому ключу тоже была сортировка. По partition-ключу — я говорил в самом начале, что это ConcurrentSkipListMap или trie — по нему можно бежать в каком-то внутреннем порядке этих ключей. Вот мы оббегаем эту мапу по каждому partition-ключу и каждую партицию записываем на диск в некотором представлении. В некотором сериализованном представлении, отдельном от того, что мы делали для commit log, но самое главное — оно упорядоченное. Относительно partition-ключей, а внутри партиции — относительно clustering-ключа.
[2:01:21] Александр: Ага, то есть это сортировка двухуровневая.
[2:01:24] Дмитрий: Ну да. Я бегу по внешней мапе, по partition-ключам, которые отсортированы, записываю partition-ключ, потом у меня для этого partition-ключа набор строчек. Эти строчки я отсортировал по clustering-ключу и, соответственно, пробежался по каждому из этих clustering-ключей, записал строчку в том порядке, в каком она отсортирована. Дальше перехожу к следующей партиции, пробегаю по порядку строчки, ассоциированные с clustering-ключами, и так далее. И вот это я всё пишу в том некоем сортированном порядке, что мне важно. Собственно, поэтому эта структура называется SSTable — sorted string table. Sorted — это как раз идея о том, что ключи на диске не рандомно раскиданы, а отсортированы.
[2:02:09] Александр: Ага, окей.
[2:02:12] Дмитрий: Дальше там, я не знаю, будем ли мы углубляться в то, как оно устроено или нет. Там всякие весёлые оптимизации есть: как бы писать поменьше, с одной стороны; как бы нам это сжать, с другой; а с третьей — как нам ещё всё это потом эффективно начитать, поэтому нам нужны индексы будут. Но пока, наверное, это опустим и поговорим про последнюю часть LSM-деревьев.
[2:02:37] Дмитрий: Мы сбрасываем данные на диск, записываем что-то в файлик. Потом мы снова что-то пишем, снова всё накапливается в памяти, снова сбрасываем на диск, снова пишем в файлик. То есть система непрерывно порождает вот эти новые файлики. Мы на этом этапе не делаем случайные записи, мы никакие существующие файлики не модифицируем. Мы каждый раз порождаем immutable новый файлик. То есть мы не берём существующий файлик и что-то туда перезаписываем — мы создаём новый каждый раз. У нас начинает плодиться этот набор файликов. Это, собственно, куски этого LSM-дерева уже на диске, его фрагменты. То есть у нас есть фрагмент в памяти, называемый memtable, и множество SSTable-файликов на диске.
[2:03:10] Дмитрий: В принципе, как бы могло бы и работать. Проблема в том, что этих файликов становится всё больше и больше. Читать из них уже эффективно не получается, с одной стороны. С другой стороны, если вы всё время обновляли одну и ту же строчку, условно говоря, или там небольшой набор строчек, у вас эта строчка будет на диске во множестве экземпляров и будет сжирать место. То есть у вас будет дублирование этой строчки и съедать место на диске, и вы бы хотели от него избавиться. У нас будет одна актуальная, самая последняя записанная в самом последнем файле строчка, а все предыдущие файлы неактуальны, но мы этого не знаем, и мы вынуждены их каждый раз читать, логику запускать и находить последнее. И место на диске оно занимает.
[2:03:57] Александр: И место, да. И CPU, и space.
[2:04:11] Дмитрий: Рано или поздно вы скушаете всё место на диске, и ничего не останется. Поэтому нам нужно что-то с этим делать. И вот последний краеугольный кусок LSM-деревьев — это что же с этим делать? Мы хотим как-то эти файлики на диске схлопывать. Собственно, процесс утрамбовывания нескольких файликов вместе называется compaction для LSM-деревьев. Есть разные способы это делать, но идея, объединяющая их все, следующая: мы выбираем какие-то кандидаты. Файликов много, можно по-разному выбирать кандидатов, и это нетривиальная задача — как нам выбрать удачное сочетание файликов, чтобы эту операцию оптимизировать.
[2:04:56] Александр: Ну да. Представляем, что эти файлы плодятся каждые, не знаю, 10 секунд просто. Тык-тык-тык — новый файл, новый файл. От интенсивности записи. Но что с этими файлами делать? Мы compaction запускаем.
[2:05:09] Дмитрий: Да, мы запускаем для них compaction. Compaction берёт на вход несколько файликов, никаких не меняет и делает по сути дела сортировку слиянием — склеивает несколько файликов в результирующий файлик. И здесь мы пользуемся тем, что они были отсортированы — почему я так сильно подчёркивал, что они должны быть отсортированы. Мы можем читать старые файлики, бежать по ним, перебирая строчку за строчкой, строчки между собой мержить. То есть у меня получается сортировка слиянием, а именно её часть, связанная со слиянием. Мне надо слить две одинаковые строчки, у которых ключи совпали. Если для одной строчки нет других кандидатов, я просто её записываю на выход как результат этого алгоритма. Если же две строчки с одинаковыми ключами мне попались, я должен их слить. И тут снова работает та логика, которую мы только что для memtable обсуждали: мы объединяем эти ячейки, побеждает снова та ячейка, у которой таймстемп выше. Ровно та же идея. Мы слили, получилась строчка-результат, и результат мы записали на диск. В этом алгоритме мы делаем последовательное чтение с диска, делаем последовательную запись на диск, а в памяти держим только текущие строчки. Нам не надо вычитывать всю таблицу в память, которая может быть — этот файлик может быть очень большой. Мы читаем строчку за строчкой.
[2:06:34] Александр: Да, в этом как раз и основное преимущество merge sort — что она позволяет сливать большие файлы, которые каждый по отдельности или вместе в памяти единовременно могут и не уместиться. Но merge sort всё равно их переварит.
[2:06:48] Дмитрий: Да, собственно, это называется алгоритмами внешней сортировки. И мы пишем результат в новый файлик, и как только новый файлик мы сформировали — он содержит всю нужную нам информацию, — старые файлики нам больше не нужны, мы их можем удалить. В результате мы n файликов склеили в моём примере в один файлик. В общем случае, там в зависимости от алгоритма compaction может получиться несколько файликов: возможно, мы хотим их не слишком большими делать и нарезать на какие-то кусочки. И дальше там есть целое дерево этих алгоритмов — Leveled Compaction, Size-Tiered Compaction, Time Window Compaction, — которые оптимизированы под разные сценарии. Потому что это, говорю, правильная задача, но нетривиальная. И как в CAP-теореме, там можно выиграть в чём-то одном, но проиграть в другом. Есть в базах данных тоже своя, так сказать, CAP-теорема под названием RUM — Read, Update, Memory. То есть как быстро мы можем читать, как быстро можем писать, и сколько нам для этого памяти нужно или места на диске. Вот с compaction то же самое: есть алгоритмы, эффективные под запись, есть эффективные под чтение данных. И ещё вопрос, сколько мы готовы места на диске потратить. Это иногда называется write amplification factor — насколько часто мы вообще с диском работаем.
[2:08:24] Дмитрий: Про эти алгоритмы можно разговаривать долго. Подчеркну только, что, наверное, интересно: в пятой версии Cassandra появился новый алгоритм — Unified Compaction Strategy, — который, судя по всему, превосходит существующие, может быть, за исключением Time Window Compaction, который специально под конкретный кейс был сделан, и работает существенно быстрее. Но логика за ним внутри очень нетривиальная. Если в двух словах сказать, просто ключевые слова без понимания, — он использует понятие плотности таблиц. То есть он пытается объединить таблицы с похожей плотностью. Другие алгоритмы пытаются… Простой алгоритм: возьмём файлики одинакового размера, или почти одинакового размера, и сольём их вместе. Большой файлик с маленьким сливать невыгодно, потому что мы большой файлик будем почти полностью перечитывать и перезаписывать как есть. Лучше сливать файлики одинакового размера. Вот там эта идея расширена и введено понятие плотности файликов. То есть количество ключей на диапазон вот этого кольца, о котором мы говорили в самом начале.
[2:09:32] Александр: Очень продвинутая тема.
[2:09:34] Дмитрий: На эту тему есть доклады у автора этого алгоритма.
[2:09:38] Александр: Ключевые слова какие, получается, гуглить, если я захочу?
[2:09:41] Дмитрий: Unified Compaction Strategy.
[2:09:43] Александр: Отлично. Слушатели, запомнили.
[2:09:46] Дмитрий: Автор, по-моему, его из Болгарии. Собственно, вот за счёт чего мы не плодим файлики до бесконечности.
[2:09:50] Александр: И вот наш пазл с точки зрения записи сложился. Мы пишем данные в memtable, копим, копим, копим, сбрасываем memtable на диск — получается SSTable. И для SSTable периодически в бэкграунде запускается процесс compaction, который его схлопывает. Чем-то он даже похож на механизм автовакуума в Postgres. Наверное, можно попытаться провести аналогию: тоже убирает некие старые версии строчек, назовём это так.
[2:10:28] Дмитрий: И тут, наверное, про операцию записи в основном всё. Последний момент, который стоит осветить, — это операция удаления.
[2:10:38] Александр: Так, да. Это же тоже мутация, но она особенная.
[2:10:41] Дмитрий: Да. Казалось бы, удалить — как это, ломать не строить, удалять проще, чем добавлять.
[2:10:48] Александр: На самом деле нет.
[2:10:49] Дмитрий: На самом деле в распределённых системах — и не только — это сложнее сделать. Давайте представим, что в том, что я описал, мы бы просто удаляли данные. Вот у нас прилетает запрос на удаление, разлетается по репликам. Мы записали в наш commit log, что надо что-то удалить, — там просто в бинарном виде записываем операцию, а не результаты: удалить то-то. Ну, окей, записали. А в memtable мы пришли, и, предположим, там что-то было — мы удалили. Ну, хэш-мапа, в ней же есть операция удаления. Удалили там — и всё, реализовали то, что хотели. Нет, потому что мы наткнулись на первую проблему LSM-деревьев: что этот memtable — неполное представление того, что у меня есть. То, что записи в memtable нет, не означает, что у меня нет её суммарно в базе данных. Потому что она может лежать в одной из тех таблиц, что я сбросил на диск, — вот в тех файликах, которые только что мы сбрасывали. В какой-то из них лежит старая запись. Потому что данные, которые в этом дереве in-memory лежат в memtable, мы физически не можем все загрузить, потому что мы храним много данных. И поэтому мы храним какое-то подмножество данных. Подмножество означает, что оно не представляет весь набор данных. Соответственно, отсутствие строчки в этом подмножестве не означает отсутствие этой строчки во всём множестве.
[2:12:19] Александр: Да, всё так.
[2:12:20] Дмитрий: И мы не хотим идти и эту строчку искать и удалять в этих файликах рандомной записью. Это надо мне… значит, эта структура данных не предназначена под обновление, и я не смогу там ничего обновить. Что делать?
[2:12:33] Александр: Ну да, я предположу за слушателя. Слушатель, возможно, нам нужно не просто удалять из мапы, делать remove, потому что это недостаточно. Нам нужно где-то — либо в какой-то дополнительной структуре данных — сохранять информацию о том, что вот это было удалено.
[2:12:52] Дмитрий: Да.
[2:12:55] Александр: Пожалуйста, сделай что-то с этим. Либо мы хотим внутри структуры данных нашей мапы вместо удаления сделать put — то есть сделать апдейт существующей записи, но со специальным значением, какую-нибудь туда вставить пометочку, типа «это удалено».
[2:13:11] Дмитрий: Да. Всё так, и это на самом деле неявно происходит много где. Вы приходите в какую-нибудь социальную сеть, говорите «удалить мой аккаунт». Думаете ли вы, что там всё сразу удалится?
[2:13:23] Александр: Нет, не думаю.
[2:13:25] Дмитрий: Там поставится просто логическая пометка о том, что аккаунт удалён. То есть, грубо говоря, на существующих реальных записях в дополнительной служебной колонке вам пометят, что статус — удалён. И, может быть, через какое-то время очистят реально целиком. Или то же самое, когда вы удаляете данные с диска: вы удаляете файлик файловой системы. Думаете, что там всё реально удалилось? Нет. В метаданных файловой системы пометка о том, что этот блок свободен, была поставлена. Файла больше нет, но байтики с вашими данными по-прежнему там лежат, пока они не перезатрутся чем-то другим. И этим, кстати, можно пользоваться. Вы удалили данные с вашего ноутбука. Думаете, никто не найдёт? Берём диск, расковыриваем его специальными утилитами и вытаскиваем содержимое всех файликов, что вы удалили. Недавно, по крайней мере, пока они не перезаписались. Увлечённый инженер-хакер с правильным набором тулов сможет добраться до данных. Но, правда, у этого действия тоже есть срок годности: скорее всего, рано или поздно перезапишется. Поэтому чем быстрее, тем лучше — скорое вмешательство повышает вероятность нахождения данных.
[2:14:39] Дмитрий: Поэтому, да, Cassandra и LSM-деревья в принципе так и делают. Парадокс: удаление — это вставка в этих структурах данных.
[2:14:53] Александр: Да.
[2:14:54] Дмитрий: То есть мы вставляем… Вы говорите «удалить вот эту строчку». Мы говорим: окей, это вставка. Удалить такой-то ключ — partition-ключ, clustering-ключ. А в конце дальше у нас пометка о том, что это удаление. И там снова таймстемп. Вспоминаем, что мы всё мержим по таймстемпу. Операции записи с операциями удаления могут пересекаться, в том плане что накладываться одна на другую, — побеждает та операция, у которой таймстемп выше. У операции записи есть таймстемп, и у операции удаления есть таймстемп. Поэтому мы их можем мержить, поэтому все эти алгоритмы мержа снова здесь работают. И за счёт этого мы эту tombstone… Я слово сказал, как это называется в базе данных. Это называется tombstone, могильная плита. Мы пишем, что всю строчку похоронили, выбили там дату смерти её. И потом благополучно это всё накопили в нашем memtable. То есть с точки зрения LSM-дерева это почти такая же строчка, как и обычная: мы накопили её в памяти и потом сбросили на диск. Дальше, когда мы делаем compaction, те же правила мержа работают и там. У нас есть живая строчка, она встретилась с мёртвой строчкой — та, у которой таймстемп выше, победила. Может, у вас было удаление, вставка, потом снова удаление, потом снова вставка. Такие последовательности никто же не отменяет для вашей строчки, для вашего первичного ключа. Поэтому всё это как-то смержилось, получился какой-то результат. Финальным результатом может оказаться вот эта пометка о том, что что-то удалилось. И она будет записана физически на диск.
[2:16:33] Дмитрий: Мы не можем её в этот момент просто удалить и забыть про неё. Мы её будем писать на диск, это будет занимать какое-то место. Нам рано или поздно нужно сделать сборку мусора на файловой системе — нам нужно вот эти tombstone удалить. Они живут какое-то определённое количество времени, после которого мы считаем, что они доступны для удаления. И два способа их удаления, собственно, Cassandra делает. Либо она, когда очередной раз делает compaction, всё равно пробегается по этим данным. Она видит таймстемпы этих удалённых записей. Видит, что запись у неё старая, tombstone можно удалить. Можно ли?
[2:17:18] Александр: Так, вопрос подвозный. Со звёздочкой, я сказал такой, да.
[2:17:24] Дмитрий: Ты имеешь в виду во время compaction, увидев запись tombstone… Допустим, ты увидел tombstone, он был сделан, например, 20 дней назад. Ты решил, что tombstone вроде как только 10 дней можно хранить, а дальше мы без них обойдёмся. Можно ли его так просто взять и удалить?
[2:17:43] Александр: Ну смотри, я при создании нового файла, при слиянии, сортировке, увидев tombstone, я просто его скипну.
[2:17:54] Дмитрий: Да-да-да, я про это. Могу ли я его скипать? Или всё-таки я должен его записать, если что-то где-то как-то…
[2:17:59] Александр: Что-то где-то не так. Но ты сказал, что 10 дней мы храним эти tombstone.
[2:18:00] Дмитрий: Точно храним. Раньше 10 не удалим.
[2:18:10] Александр: А после? Этот период тоже прошёл.
[2:18:12] Дмитрий: Мы можем себе позволить с точки зрения клиента реально удалить эту запись. Мы перед клиентами… Зачем мы храним 10 дней? Мы сейчас поговорим чуть позже. Для этой цели, пока мы говорим про локальную ноду, нам не важно, сколько дней мы храним. Тут проблема в другом. Если мы удаляем эту запись, то мы как бы про неё забываем, о её существовании фактически. Мы говорим, что мы не знаем, что она была — была не удалена, её не было. Как любой другой никогда не вставленной записи в нашу базу данных, она приравнивается к вот этим. И тогда, сделав такое и забыв реально навсегда, я могу, возможно, не отличить… То есть ноды, мои пиры, которые, возможно, реплики…
[2:19:01] Александр: Забудь пока про другие ноды.
[2:19:01] Дмитрий: Мы говорим просто про локальный алгоритм на одной ноде.
[2:19:06] Александр: Я подвоха не могу найти, честное слово. Но твоя улыбка говорит, что что-то есть.
[2:19:11] Дмитрий: А он есть. Конечно же, есть.
[2:19:13] Александр: Ну давай.
[2:19:14] Дмитрий: Смотри, когда мы делали алгоритм compaction, у нас на диске множество файликов изначально. Мы какие-то из них выбираем, чтобы делать compaction. Там могут гигабайты файликов валяться — ты не компактишь сразу все, ты выбираешь некое подмножество, которое тебе удобно; кто-то там пересекается по ключам ещё. Начинаешь компактить. Но у тебя рядом может лежать ещё один SSTable, который в это compaction-множество не попал. А что, если там лежит та же строчка, но ещё старее? У тебя была старая строчка, поверх неё потом в другом файле попал tombstone для этой же строчки, но более свежий. Когда ты читаешь, ты видишь и то, и то, — как, видимо, мы поговорим в следующий раз, — склеиваешь и понимаешь, что данных нет. Tombstone побеждает живую запись, результата записи нет. Tombstone означает: клиент не увидит ничего, но база понимает это. А теперь ты tombstone свой во время compaction решил убирать. И та запись, которая была в терминах Cassandra shadow — затенена tombstone’ом… Ты tombstone, который тень на неё отбрасывал, выкинул. Эта запись волшебным образом снова в действии, она снова живая — как это, зомби вылезли из могил. Даже не в распределённой системе данных.
[2:20:47] Александр: Я понял проблему. Получается, что мы не можем даже в рамках локального процесса compaction позволить себе действительно скипнуть запись с tombstone и не писать её в следующий результирующий файл, потому что мы не можем гарантировать, что нет где-то рядом записи апдейта для того же самого ключа, который мы сейчас хотим удалять. И этот апдейт возродится из мёртвых, потому что мы запись о его удалении сейчас удалим, и получится, что мы его воскресили. Нам для полной картины нужно иметь все апдейты, все мутации, которые происходили для конкретной строчки, — и только в этом случае мы можем принять решение о том, что мы её хотим удалять, потому что знаем всю её историю. А пока мы всю историю не знаем, мы боимся shadowing. Правильно?
[2:21:11] Дмитрий: Да и нет. Ты прав в том плане, что мы не можем просто так её удалить, потому что возникнет эта проблема. Но это не совсем так, что единственный способ решить эту проблему — чтобы всё поучаствовало в compaction. Это называется full compaction, или major compaction, когда мы все таблицы сгребли и начали сквозь эту мельницу слияние проводить, — только тогда мы можем удалить tombstone. На самом деле мы тут можем довольно часто выкинуть его, сделав дополнительные, достаточно дешёвые проверки.
[2:22:12] Александр: Так.
[2:22:13] Дмитрий: И там их несколько, эта логика есть в Cassandra, и там есть несколько вариантов. Вариант первый: ты можешь посмотреть на таймстемпы разных файликов. У каждого SSTable на диске есть набор метаданных, ассоциированный с данными, которые лежат в этой таблице. И там, например, есть для каждого файлика минимальный таймстемп и максимальный таймстемп тех данных, что записаны в этом файлике. Это первая информация. Ты можешь сказать: окей, у меня есть tombstone такой-то даты, а все файлики, которые лежат и не участвуют в compaction, они старее… Ой, извини, они моложе моего файлика.
[2:22:54] Александр: Ага, то есть этот tombstone он самый старый у меня.
[2:22:57] Дмитрий: Да, он не может отбрасывать тень, потому что явно старее, чем максимальная запись в другом файлике. Поэтому удалить я его могу беспрепятственно, поскольку он очень старенький. Внизу этого слоя тарелочек он лежит самым нижним, он никого не затянет. Это первый вариант. Второй вариант: у тебя, опять-таки, данные отсортированные в SSTable, и в этих метаданных есть информация о том, какой самый маленький ключ и самый большой ключ в каждом файлике лежит. Min/max values. Ты можешь посмотреть на ключ tombstone, сравнить его с min/max других файликов, не участвующих в твоём compaction. Если он не попадает ни в один из этих диапазонов… У тебя может быть серия данных базы данных, но у этих данных вообще пересечений нет, всё прекрасно разложилось на разные файлики и не пересекается. Окей, он физически не может перекрыть никого другого, потому что отрезки ключей разные. У тебя ключ этого tombstone не пересекается с отрезком min/max ни одного из файликов на диске. Эвристика такая, да. Другая эвристика, которая тоже позволяет сказать, что этот tombstone безопасно выкинуть, потому что нет пересечений.
[2:24:16] Дмитрий: Ну и наконец, третий способ, который Cassandra использует, если эти два не помогли: у меня есть способ для каждого файлика быстро и эффективно проверить, есть ли partition-ключ в этом файлике. Он нам пригодится потом для записи… этот способ нам пригодится потом для чтения, но он полезен и здесь. Мы можем быстро сходить для другой таблицы и спросить: а в этой таблице есть что-нибудь с partition-ключом таким-то? И нам скажут «да» или «нет». Причём эта структура данных может быть вероятностной: она может нам сказать «точно нету» или сказать «может быть, есть». Но мы читать при этом не будем — то есть мы с диска не читаем, всё это в памяти находится. Если эта структура данных нам сказала, что там точно ничего нет, значит, мы можем, опять-таки, tombstone удалить. Мы сейчас говорим про структуру данных под названием Bloom filter. Она в основном используется для операции чтения, для того чтобы все файлики не перебирать. Но вот она чудесным образом помогает нам и здесь.
[2:24:44] Александр: Ага, в процессе удаления.
[2:24:44] Дмитрий: Да, в процессе удаления tombstone мы можем заглянуть, проверить другие SSTable на наличие пересекающихся записей, используя Bloom filter этих других таблиц. И если одно из этих условий нам позволяет сказать, что точно ничего не пересекается, мы можем откинуть эту запись. Это на практике весьма и весьма часто случается. Ну а в противном случае мы ждём, когда либо таблицы сойдутся в общем compaction — старой, новой — и пересекутся, либо в крайнем случае ты руками можешь запустить полный compaction и форсировать этот процесс, чтобы удаление всё-таки произошло.
[2:25:56] Дмитрий: Вот это зачем нужны tombstone с точки зрения LSM-деревьев, даже если мы не говорим про какую-то распределённость. Даже в обычной локальной реализации на одном диске они тебе нужны, потому что у тебя есть процесс compaction, есть несколько источников одной и той же строки: она может быть в ряде SSTable, да и в memtable она тоже может заваляться. Поэтому тебе нужно координироваться здесь с явной записью, с метаданными. А вторая причина, почему нам это нужно, — вспоминаем про нашу историю с тремя репликами. Вот когда я говорил про «новая запись, старая запись перекрылись» — мы хотим, чтобы реплики у нас были согласованными.
[2:26:34] Дмитрий: У нас возникает классическая распределённая проблема. Допустим, ты сделал вставку по какому-то ключу, и одна из нод была выключена. Точнее, наоборот. Ты сначала сделал запись: у тебя есть три реплики, все они были включены, запись попала на каждую из них, каждая реплика запомнила, что что-то есть для нашего ключа. После этого ты хочешь сделать операцию удаления. И допустим, ты бы реализовывал не через tombstone, а просто бы удалил. Допустим, в каждой ноде Cassandra было бы не LSM-дерево, а B-tree, которое бы inline всё апдейтило. Каждая нода Cassandra — это такая маленькая реляционная база данных: она просто удалила бы данные и ничего не сказала бы. Всё, больше записи нет. И допустим, у нас одна из этих трёх реплик умерла перед тем, как мы сделали операцию удаления. То есть у нас были три реплики, в которых были записаны данные до этого. Каждая реплика говорит: у меня что-то есть. Мы одну из этих реплик выключили. Теперь ты делаешь операцию записи, которая удаляет данные. Она удаляет данные с первой реплики, со второй реплики, а до третьей достучаться не может. И допустим, ты никакие tombstone там не оставлял, а просто там ничего нет. После этого ты поднимаешь третью реплику назад, и у тебя как бы есть «ничего», есть «ничего», а в третьей реплике есть «чего-то». И если использовать дефолтную стратегию слияния данных между «совсем ничего» и «есть что-то», победит «что-то». Логично было бы. Ты не отличишь теперь ситуацию — ты не можешь отличить ситуацию, — когда у тебя никогда не было записи, и ситуацию, когда у тебя была запись, но она была удалена. Они неразличимы. Тебе не хватает метаданных. И чтобы решить эту проблему, не прибегая к другим алгоритмам, консенсусам и так далее, тебе придётся записать метаданные. И этот tombstone становится метаданными с точки зрения Cassandra как распределённой системы. Tombstone решает две проблемы: он решает проблему LSM-деревьев, и он решает проблему консистентности реплик с точки зрения Cassandra как распределённой системы.
[2:29:00] Александр: Поэтому, несмотря на то что консистентность — такая больная вещь, которая выставляет весьма много проблем на практике, она, к сожалению, такой краеугольный камень, который не выковырить.
[2:29:11] Дмитрий: Слишком много каких проблем решает тоже.
[2:29:14] Александр: Могильную плиту просто так не сдвинуть. Это уж точно.
[2:29:18] Дмитрий: Ну и вот мы, получается, всё расписали на реплике, получили от них ответ, ответили обратно. И обратно приходит уведомление о том, что запись, наконец, прошла успешно.
[2:29:27] Александр: Это мы тут обсуждали всё часа два с половиной, если не больше. В реальности это занимает единицы миллисекунд.
[2:29:36] Дмитрий: Да, не прошло и трёх часов, как мы записали одну строчку в Cassandra.
[2:29:41] Александр: Да. Но мне кажется, это было восхитительно. Мне очень понравилось. В плане, что, походу, у нас будет заход на третью часть — про чтение. Потому что, да, вот так вот получилось. Но мне от этого даже больше нравится.
[2:29:53] Александр: В общем, ребят, мы с Димой…
[2:30:00] Дмитрий: Я устал, честно. У меня голова уже подскипает.
[2:30:03] Александр: Оно тяжело, да. И мы как раз находимся в хороший момент, когда стоит, по крайней мере, взять паузу. Да, я всё понял. Но для того, чтобы понять, мне нужно было записать целый сезон подкаста про базы данных и про то, как все эти структуры данных работают. Так что, ребят, я вас призываю — если что-то недопоняли, сезон подкаста «Тысяча фичей», от корки до корки, это точно пригодится. Ну и книжка, наверное, можно посоветовать. «Designing Intensive Applications», если я не ошибаюсь.
[2:30:37] Дмитрий: «Data Intensive», да. DDIA.
[2:30:37] Александр: Да, с палтусом, по-моему.
[2:30:39] Дмитрий: Нет, не с палтусом, а вот с рыбой… Это кабанчик, это с кабаном книжка. С такими же, как на бутылке вина. Это я перепутал. С палтусом — это называется книжка «Database Internals».
[2:30:52] Александр: «Database Internals».
[2:30:53] Дмитрий: Которая написана Алексом Петровым — по объявлению коммитера в Cassandra.
[2:31:00] Александр: Я про него-то и вспомнил как раз. И в целом, в том, что мы тут обсуждали, — много чего мы обсуждали, особенно про tombstone, про устройство, — чего-то можно почитать в книжке. Если вам аудио такой формат тяжеловат для восприятия — читайте.
[2:31:16] Дмитрий: Если располагать их в порядке сложности, то я сказал бы, что Мартин Клеппман, «Designing Data-Intensive Applications», она более общая. Она, надо сказать, проходит скорее по общим идеям. В «Database Internals» даже, судя по названию, там более детально разобраны те или иные аспекты. Если мы говорим про B-деревья, там прямо вам описывается, как именно варианты — как это реализовать, как вам в памяти эти сущности представлять. Или там LSM-деревья. То есть она более детальная, но, соответственно, и более сложная для восприятия. Если вы не знаете, с чего начать, и не читали обе из них, я бы рекомендовал Клеппмана начать читать первым.
[2:32:01] Александр: Да, но прочитать стоит обе однозначно по итогу. В общем, мы тут накидали и инсайдов, и книжек, и посоветовали. Короче, спасибо большое, Дим. Слушателям спасибо, что вы дослушали. Поддержать подкаст можно репостом в Телеграме. Пошерить можно со своими друзьями-программистами в выпуске: сказать «смотрите, как классно тут про Cassandra рассказывают» — и это лучшая поддержка. Ну и лайк, если вы смотрите на YouTube. Диме большое спасибо за такую ценную информацию из первых уст, от человека, который эти структуры данных не то что знает, как работают, но, я уверен, и реализовывал собственными руками, потому что…
[2:32:37] Дмитрий: Не делал с нуля, но кое-какие подкручивал из них, да.
[2:32:40] Александр: Да, потому что так понимать тонкости можно только пройдя через все стадии коммита подобных кодов.
[2:32:50] Дмитрий: Официального принятия.
[2:32:51] Александр: Да, потому что наш мозг, как правило, отказывается понимать сложные вещи и выбирает что-нибудь попроще. А тут, видимо, была мотивация понять. Это чувствуется.
[2:33:00] Александр: В общем, да, наверное, услышимся через недельку, через две. Всем пока.
[2:33:04] Дмитрий: Спасибо.