Выпуск 58 · подкаст «Тысяча фичей»

#58: Apache Cassandra, часть 3: читаем данные

2:28:11
↓ скачать mp3

Заключительная, третья часть разбора Apache Cassandra: Александр Пахомов и Дмитрий Константинов (коммитер Apache Cassandra) прослеживают путь запроса на чтение от координатора до диска. По дороге — какие чтения Cassandra делает эффективно (по `partition`- и префиксу `clustering`-ключа) и почему это OLTP-, а не аналитическая база; кворумная математика `R + W > N` и `strong consistency`; выбор реплик через `Snitch`, спекулятивные ретраи и оптимизация «данные с одной реплики, digest с остальных»; `read-repair` и пагинация без `OFFSET`; а на нижнем уровне — `Bloom filter`, `primary index` с `index summary`, `key cache`, блочное сжатие с `compression metadata`, вынос метаданных в off-heap и, наконец, боль tombstones при чтении «очередей».

Главное

  • API драйвера для чтения и записи одинаковый (те же `PreparedStatement`/`BoundStatement`); разница лишь в обработке результата — при чтении данные текут обратно от сервера, и почти всегда работает пагинация.
  • Cassandra — OLTP-, а не аналитическая база: эффективны чтения по полному `partition`-ключу плюс диапазону или префиксу `clustering`-ключа; запрос без полного `partition`-ключа вырождается в full scan по всем нодам, потому что хэш от части ключа не связан с хэшом целого.
  • Кворумное чтение опирается на неравенство `R + W > N`: если множество ответивших на запись реплик пересекается с множеством читаемых хотя бы в одной ноде, чтение увидит последнюю запись (в терминах Cassandra — `strong consistency`); `replication factor` при этом отвечает за долговечность, а не за консистентность.
  • Координатор выбирает реплики через `Snitch` (по умолчанию dynamic snitch — по медиане времени ответа), но такая обратная связь может входить в автоколебания; поэтому иногда выгоднее слать запрос сразу во все реплики ради предсказуемой нагрузки.
  • Голосования между репликами нет: ответы сливаются по «последний по таймстемпу побеждает»; в оптимизации Cassandra запрашивает полные данные лишь с одной (обычно локальной) реплики, а с остальных — `MD5`-digest, и сравнивает хэши (это ~3% CPU из-за неудачно выбранного `MD5`).
  • `read-repair`: если при чтении данные на репликах разошлись, координатор не только вернёт клиенту актуальную версию, но и запишет её обратно на устаревшие реплики — то есть чтение может породить запись.
  • Пагинация в Cassandra — только «вперёд» через закладку (курсор в партиции), без `OFFSET`/прыжков на произвольную страницу; консистентность гарантируется лишь в пределах одной страницы, а не всего запроса.
  • На чтении Cassandra отбрасывает лишние `SSTable` по min/max ключам и `Bloom filter`, находит позицию через `primary index` + `index summary` в памяти и `key cache`, а данные хранит блоками под сжатием (`LZ4` по умолчанию, `ZSTD`), что требует отдельного `compression metadata` для маппинга смещений.

В выпуске

  • Дмитрий КонстантиновСистемный архитектор и Java-разработчик, коммитер Apache Cassandra; специализируется на распределённых системах, производительности и отказоустойчивости. Регулярный спикер JPoint/Joker, на Хабре — @netudima. Habr ↗ jpoint.ru ↗
Расшифровка

[00:00] Александр: Здарова! Это 58-й выпуск подкаста «Тысяча фичей». Вы слушаете третью и заключительную часть про Cassandra. Ранее мы подробно рассмотрели клиент и сервер на примере записи данных, а сегодня погружаемся в процесс чтения. В гостях Дима Константинов, коммитер в Apache Cassandra. Поехали!

[00:33] Александр: Так, ну что, давай верхнеуровнево проговорим в третьем выпуске то, что было в первом и втором. Мне кажется, освежить память слушателям стоит. И лишним не будет напомнить: если вы, дорогие подслушатели, не слушали первую и вторую части про Cassandra, где мы до болтиков разобрали, как устроен клиент и как работает запись в Apache Cassandra, — то обязательно сначала прослушайте их по порядку, первую и вторую, и только потом переходите к этой, к третьей. Потому что многие детали мы будем опускать, подразумевая, что вы уже знакомы с материалом. А мы переходим к тому, о чём говорили верхнеуровнево в первых двух. Сначала мы разобрали клиента: как он устроен, прошли всякие интересные алгоритмы — retry, backpressure и так далее, — разобрались с сетью. Это было интересно. Во второй части мы глубоко посмотрели на то, как вообще устроены LSM-деревья в Apache Cassandra и какие там вообще деревья. Я помню, я представил это себе как трёхуровневую штуку: хэш-таблица, хэш-таблица, хэш-таблица — вот это я помню. И хэш-таблица там только на первом уровне, а дальше B-деревья.

[01:46] Дмитрий: Даже на первом уровне, на самом деле, там не совсем хэш-таблица — там тоже нужна упорядоченность. Потому что мы должны записывать данные из памяти на диск в отсортированном виде, поэтому нам выгодно, чтобы они сортировались сразу при вставке. Поэтому там либо ConcurrentSkipListMap — сортированная concurrent-коллекция, единственная сортированная concurrent-коллекция в стандартной поставке Java, если не считать совсем медленного для наших случаев варианта — синхронизированной обёртки вокруг TreeMap. А на следующих уровнях у нас, соответственно, B-дерево и B-дерево, на втором и третьем уровне.

[02:24] Александр: Точно. И эти структуры данных мы рассматривали в контексте того, как происходит запись в Apache Cassandra. А сегодня мы будем эти данные читать. И начнём мы, наверное, по классике, с клиента. Как мы читаем данные? Дим, расскажи.

[02:41] Дмитрий: Да там вся та же самая логика, что мы обсуждали для записи, разницы практически никакой.

[02:52] Дмитрий: Был китайский фильм про объединение Китая, «Герой», по-моему, назывался, где была идеология, что все под небесами, всё везде одинаково. А тут то же самое: если ты понял принцип какой-то одной системы, то тебе открывается довольно большой класс систем, они становятся понятными, даже если ты не знаешь деталей.

[03:14] Александр: Ну, кстати, да. Давайте всё-таки повторим — напомнить тем, кто не слушал, и, может быть, дать чуть лучшее понимание. Итак, мы каким-то способом подключились к кластеру, открыли наши TCP-коннекшены. Дальше мы хотим послать запрос. Либо мы его на ходу склеиваем и посылаем, либо — что чаще всего вы будете делать в реальном приложении под большой нагрузкой (а Cassandra всё-таки рассчитана на случай, когда у вас большая нагрузка и большие объёмы данных) — вы подготавливаете запрос заранее, и дальше вам нужен только id этого запроса плюс изменяемые параметры. Эти параметры вы формируете в виде объекта Statement. В нашем случае это PreparedStatement, объект, который содержит id подготовленного запроса. И на его основе вы порождаете BoundStatement, в котором линкуете этот id запроса с конкретными параметрами, которые хотите выполнить. Дальше всё это каким-то способом кодируется в набор байтиков и с помощью библиотеки Netty отправляется на сервер. Слушай, эти стейтменты, которые я писал — ну, мы писали вместе — в записи, где делали, грубо говоря, INSERT, — это те же самые сущности в клиенте, что я использую для чтения? То есть те же самые стейтменты, просто SQL другой? Или это всё-таки другой API?

[04:44] Дмитрий: Нет, это те же самые стейтменты. Разница только в том, как ты работаешь с результатами. В случае записи результат — успешно/неуспешно, и по сути дела больше тебе ничего не возвращается.

[04:56] Александр: А что там, «успешно/неуспешно» — это как? Completable future, какой-то boolean?

[05:00] Дмитрий: Это CompletableFuture, который скажет либо success, либо бросит исключение. Либо, если API синхронный, — либо успех, и операция просто вернётся как обычно, либо бросит exception: «извини, не получилось». И там довольно разнообразные иерархии исключений, которые как раз позволяют понять, какие у тебя ошибки, и, может быть, сделать более точечную обработку. Вот то, что мы обсуждали про retry, — там логика такая, что по типу исключения ты можешь понять, что конкретно случилось на сервере и стоит ли делать retry.

[05:38] Александр: Да, я помню, мы про это говорили, там жёстко. Но API точно такой же.

[05:41] Дмитрий: Нет отдельных prepared statement для чтения и prepared statement для записи, они все одинаковые. Разница лишь в том, что когда ты читаешь, то, в отличие от записи, основные данные, с которыми ты работаешь, — это то, что ты получил от сервера. То есть поток данных, скорее, идёт в обратную сторону: в случае записи ты больше отправляешь, а в случае чтения отправляешь немного, а получаешь немного или много — в зависимости от того, какой тип запроса послал.

[06:13] Александр: Давай тут немножко проговорим, какого типа чтения вообще можно делать в Cassandra.

[06:22] Дмитрий: Самый типичный случай, под который Cassandra рассчитана… Она всё-таки OLTP-база, а не аналитическая, поэтому каких-то сложных аналитических запросов ты в ней не напишешь. Простенькие написать можно, но не сказать, что они прямо будут очень эффективно работать. Витрины, агрегированные с каунтерами, суммами, скользящими средними, — не стоит писать.

[06:49] Александр: Да, возьмите ClickHouse какой-нибудь, и всё.

[06:52] Дмитрий: У тебя, как я помню, уже были рассказы от авторов ClickHouse, где они как раз рассказывали, как всякое такое оптимизировали.

[07:01] Александр: Безусловно. У меня есть серия выпусков про ClickHouse. В телеграм-канале «Тысяча фичей» есть подборка подкастов на эту тему. Так что гуглите, ребят, «Тысяча фичей ClickHouse», или заходите на YouTube — там у нас есть целый плейлист про то, как оптимизировать ClickHouse. Мы прямо непосредственно это делаем.

[07:20] Дмитрий: Да. Но при этом, несмотря на то, что базы заточены под разные сценарии использования, написаны на разных языках и имеют разную архитектуру, тем не менее в них есть одинаковые части. Например, и там, и там используются LSM-деревья. Казалось бы, очень разные вещи, но отдельные алгоритмы, как мы уже обсудили, могут встречаться тут и там.

[07:44] Дмитрий: Итак, наиболее типичный случай. Вы можете прочитать по полному ключу конкретную запись. Вспоминаем, что в Cassandra у таблицы всегда есть primary-ключ — не может быть таблицы без primary-ключа, как в некоторых других базах. Этот primary-ключ состоит из двух частей. Есть partition-ключ, который отвечает за то, на каких репликах ваши данные по факту находятся, и позволяет делать горизонтальное масштабирование за счёт шардирования по этому ключу. И внутри партиции есть clustering-ключ, который позволяет выбирать данные в отсортированном виде. Поэтому вы можете либо указать partition-ключ плюс clustering-ключ полностью — и тогда вернётся либо одна, либо ноль записей, поскольку мы указали полный идентификатор строчки. Либо вы можете сделать запрос по partition-ключу плюс диапазону, например, clustering-ключей. Допустим, у нас partition-ключ — это id пользователя, а clustering-ключ — дата. И таблица у нас — это какие-то события, связанные с пользователем. Я хочу выбрать события за диапазон времени для конкретного пользователя: передаю id пользователя как partition-ключ, реплики становятся известны, и я могу послать запросы понятно куда. А для clustering-ключа передаю диапазон строк, которые, как мы посмотрим позже, за счёт того, что они хранятся в отсортированном виде, я могу быстро найти и вывести, не сильно утруждая Cassandra-сервер.

[09:27] Дмитрий: В этом плане идеология Cassandra отличается от реляционных баз данных. Хотя в реляционных базах, если вы сильно задумываетесь о производительности и ходите уже по границе того, что база может, вы тоже будете про это думать. Но всё-таки классическая модель в реляционных базах: давайте сначала разложим данные в какую-то нормальную форму, разделим на кусочки — таблицы сюда, сущности сюда, связи туда, всё нормализовано, всё понятно, — а дальше начинаем по ним строить запросы, где надо докидываем индексы. В Cassandra же путь идёт от паттернов чтения. То есть вы сначала должны понять, как будете читать данные, и под эти паттерны строите таблицы. Сейчас, конечно, появились вторичные индексы в Cassandra, уже более эффективные, но в основном случае вы должны подстраивать структуру данных, ваши таблицы, под то, как вы будете читать. Поэтому выбор primary-ключей — какие у вас ключи, по каким идёт сортировка (это clustering-ключи) — очень важен. Если вы сделаете его неправильно, будете сильно мучиться. И про аналитику, кстати, то же самое: в аналитике вы заранее не знаете паттерны чтения, поэтому сделать эффективную структуру хранения не очень получается. Плюс Cassandra не лучший вариант для выполнения аналитики поверх неё — вы можете использовать её как хранилище, откуда берёте данные, но непосредственно логику анализа на ней самой вряд ли сделаете.

[11:12] Дмитрий: Значит, мы можем выбрать по диапазону: partition-ключ конкретный, clustering-ключ — диапазон. Можете просто выбрать все ключи для конкретной партиции — говорим, вот такой-то конкретный id, а события все, от минус бесконечности до плюс бесконечности; частный случай, когда мы clustering-ключ никак не ограничиваем. И есть интересные варианты посередине. Вспоминаем, что и partition-, и clustering-ключ могут состоять из нескольких колонок. И для clustering-ключа порядок важен. Мы можем выбирать по префиксу clustering-ключа: сказать «дай мне данные, где первая колонка clustering-ключа конкретная, а для второй колонки — диапазон значений». Нечто похожее есть и в реляционных базах: когда у вас есть индекс по нескольким колонкам, вы можете эффективно искать по нему, указывая не все колонки, а первые несколько. Но с конца не получится: если у вас три колонки, вы можете указать первую; первую и вторую; первую, вторую и третью. А вот только вторую и третью — это эффективно уже не работает.

[12:33] Александр: Ага.

[12:33] Дмитрий: Потому что они отсортированы, грубо говоря, иерархически: сначала по первой части составного ключа, внутри — по второй, внутри — по третьей, по четвёртой. Чтобы зароутить, найти нужную строчку, локализовать её, нужно искать именно в этом порядке. И никакого другого порядка нет — у нас нет inverted-индексов, которые позволяли бы смотреть и так, и так.

[13:03] Александр: Это тоже важно для понимания того, что можно и чего нельзя с Apache Cassandra делать, когда у нас реально много данных: просто не будет работать, если не укажем в нужном порядке.

[13:14] Дмитрий: Надо понимать: для clustering-ключа можно задавать его кусочно, начальными несколькими колонками. Но для partition-ключа это не работает. Вспоминаем, что от него мы считаем хэш. Если вы посчитаете хэш только по первой колонке и не учтёте значение второй, вы попадёте совершенно в другое место. Значения хэша двух колонок и хэша одной колонки никак между собой не связаны, не скоррелированы — эта информация вам ничего не даст. Если вы не дадите Cassandra полный partition-ключ, ей придётся пойти на все ноды в кластере, которых может быть очень много, и в каждой ноде просто делать полный перебор, искать по вашему кусочку partition-ключа нужную партицию. Такие запросы эффективно не работают. Это, можно сказать, full-scan-запросы, которых в Cassandra я бы рекомендовал избегать.

[14:16] Дмитрий: Иногда их всё-таки имеет смысл выполнять — это скорее третья категория, когда вы хотите сделать полную выборку, просто просканировать таблицу целиком. В этом случае в Cassandra есть определённые механизмы, есть даже специальные утилиты. Например, во встроенном консольном клиенте cqlsh есть команда COPY, которая умеет выгружать CSV-файлик. И есть отдельная тула, бесплатная, — хотя её сделал DataStax, она под лицензией Apache, — называется DSBulk (DataStax Bulk Loader), которая позволяет…

[14:52] Александр: Как, DSBulk?

[14:54] Дмитрий: DSBulk. У меня с буквой «эль» проблемы.

[15:00] Александр: Ты вполне отчётливо её произнёс.

[15:02] Дмитрий: DSBulk. Под этим ключевым словом её можно найти. Там как раз написана логика, оптимизированная под загрузку и выгрузку большого объёма данных. Это когда мы хотим таблицу из базы сдампить наружу или, наоборот, сделать импорт — загрузить в базу таблицу из внешнего источника, того же CSV-файла. Во время экспорта нам придётся пробежаться по всей таблице и выгрузить записи наружу. Partition-ключей мы здесь не знаем, поэтому бежим по всем нодам и вытаскиваем данные кусочек за кусочком. Там специальный тип запросов.

[15:42] Александр: А какой API она использует? Она как бы через файловую систему?

[15:45] Дмитрий: Нет, тоже через CQL и cqlsh. То есть через тот же клиент, те же самые запросы. Но запросы специальные. Чтобы это эффективно работало, мы можем в запросах указывать так называемые токены. Вот когда мы разговаривали про consistent hashing, мы говорили, что каждой ноде назначается набор виртуальных нод, которые соответствуют неким числам на кольце хэширования. И вот эти числа в Cassandra называются токенами. Каждой ноде принадлежит диапазон этих чисел — token range в терминах Cassandra. И, перебирая эти точки на кольце, мы можем запрашивать данные с конкретных реплик. То есть мы явно указываем Cassandra, за счёт этих диапазонов, с какими нодами хотим работать, чтобы не делать broadcast на все ноды, не вытаскивать данные сразу с кучи нод, когда координация будет дорогой. За счёт указания этих диапазонов ключей как токенов мы попадаем на конкретные ноды и получаем данные более эффективно. Собственно, cqlsh и DSBulk этим и занимаются: служебную работу по вытаскиванию диапазонов ключей из системных таблиц они берут на себя и делают полный перебор более оптимальным способом.

[16:13] Александр: Ну, понятно, их ценность теперь приобретает смысл. То есть почему я не могу быстренько написать свой CLI-tool — потому что мне придётся всем этим заниматься?

[16:25] Дмитрий: Да, технически ты можешь: API ровно тот же самый, не какой-то скрытый. Но тебе придётся повозиться, запустить руки в эти кишки, а это не очень тривиально. Хотя у меня была практика, что приходилось похожим заниматься, — когда хочешь что-то кастомизированное, а готовый tool тебе не подходит. Такая возможность есть, но вряд ли те, кто начинал работать с Cassandra недавно, будут чем-то таким заниматься.

[17:55] Дмитрий: Наверное, стоит упомянуть — ты как раз спросил, через какой интерфейс она работает: всё, что я описал, работает через CQL-интерфейс, через этот протокол. Но есть ещё случаи, когда мы хотим выгружать большие диапазоны данных более эффективно — например, переносить что-нибудь из нашего онлайн-хранилища в какую-нибудь аналитическую базу. Недавно, не помню, в какой конкретно версии — возможно, пока только в транке, — в Cassandra появилась фича, которая позволяет делать такие выгрузки и загрузки для аналитических целей непосредственно с диска. То есть она идёт не через CQL-интерфейс, а это Java API, которое позволяет работать напрямую с SSTable-файликами, которые мы обсудили при записи. Накладные расходы на все эти коммуникации — сокеты, кодирование, декодирование — уходят: ты работаешь по факту с некими итераторами, которые двигаются по файлам напрямую. Это, мне кажется, более эффективный способ, если хочется очень много данных быстро куда-то экспортировать.

[18:39] Александр: Да, но тут надо понимать, что раз база распределённая, то на каждой ноде тебе придётся этим заниматься по отдельности.

[18:44] Дмитрий: Да. Тут же нет одной файловой системы, которую ты взял, проитерировал; у тебя по сути набор файловых систем, каждая из которых расположена на разных машинах. В каких-то крайних случаях они могут быть и на одной, но просто разделены. По факту это усложняет работу: всё нужно оркестрировать, синхронизировать, оно может где-то отвалиться и так далее. Короче, нетривиальная штука.

[19:45] Александр: Но, кстати, идея для стартапа: написать нормальную синхронизировалку вокруг Cassandra, которая заливает данные в ClickHouse, — и вот вам HTAP-база данных готова.

[19:57] Дмитрий: Да. И тут ещё такой аспект: вспоминаем, что данные у нас хранятся в нескольких экземплярах, у нас несколько реплик для каждой записи ради отказоустойчивости. Поэтому после выгрузки ты должен ещё сделать дедупликацию. Если ты всё это загрузишь as-is, не дедуплицировав, то получишь несколько копий одной и той же записи. Поэтому где-то на этапе загрузки нужно будет их склеить.

[20:23] Александр: Да ладно, действительно, указал в аналитическом запросе — нормально.

[20:27] Дмитрий: Да, может быть, ты готов будешь делать такую дедупликацию на чтении, это тоже в каких-то случаях вариант. Такого рода логику, я знаю, делали разные компании у себя внутри. Не знаю, есть ли готовые движки для выгрузки, о чём ты говоришь про стартап; я, по крайней мере, точно видел, что Netflix или кто-то подобный делали такие движки для себя. Ну, это довольно распространённый use case для больших компаний, когда у тебя много данных, — особенно сейчас, понятно, что ты хочешь потом навернуть какую-нибудь аналитику или machine learning.

[21:07] Александр: Так, ну давай теперь пойдём обсуждать, как у нас работает процесс чтения. Погнали. Итак, мы отправили запрос по сети, он пролетел сквозь сокеты, закодировался, прилетел на сторону Cassandra-ноды, прочитали его из сети — и все те же самые истории про аутентификацию, авторизацию, rate limiting здесь применимы. То есть тут всё по аналогии. На этом этапе, мне кажется, даже эта логика авторизации, возможно, ещё разделяет read/write, а аутентификация даже не знает про тип пришедшего запроса, то есть она универсальна.

[21:48] Дмитрий: Аутентификацию ты сделал на этапе подключения, когда только открыл connection. Ты не аутентифицируешь каждый запрос — аутентификация происходит при подключении к серверу, когда ты открываешь TCP-connection. Дальше на стороне сервера существует сессия, ассоциированная с connection, и там написано, что этот TCP-connection связан с таким-то пользователем; есть некий контекст, от имени которого выполняется запрос. Соответственно, мы можем вытащить из этого контекста информацию о пользователе и во время авторизации, когда прилетел конкретный запрос, понять, что мы можем читать, а что нет.

[22:30] Дмитрий: Кстати, в Cassandra недавно появилась такая фича, как маскирование запросов. Мы можем на уровне схемы указать, что вот есть колонка с конфиденциальными данными, и мы хотим, чтобы её содержимое нельзя было посмотреть. Это фича на границе авторизации и того, как формируются результаты чтения. Можно сделать так, чтобы данные были замаскированы по определённым паттернам: например, тебе покажут не полностью номер твоей кредитной карты, а только последние четыре цифры.

[23:11] Александр: А вот эта авторизация — это называется… Токенизацией иногда это ещё называют. Есть RBAC, Role-Based Access Control. И есть, кстати, ещё Row-Based Access Control, где row — это не «роль», а строка. На уровне строки, даже на уровне каждой колонки внутри строки, мы можем накрутить правила: с этими ролями пользователи могут читать эту колонку, а с этими — нет. Я не про Cassandra говорю, а про идею.

[23:48] Дмитрий: Да, в Cassandra в бесплатной версии такого нет. Такая фича была сделана в Enterprise-версии Cassandra, это DataStax, — там как раз были row-level access policy. А здесь мы говорим именно про маскирование на уровне колонки: ты на уровне схемы задаёшь, что к этим колонкам могут иметь доступ только определённые роли. Но это касается всех строк в этой таблице, а не конкретных.

[24:17] Александр: Ага. То есть мы берём не для каждой строчки и её колонки. Грубо говоря, в той продвинутой платной системе, которую ты говоришь, я могу в одной и той же таблице сделать SELECT *, увидеть данные своей карты, а твои, рядом в запросе, — нет, потому что у меня нет прав их читать, хотя колонка одна и та же. А сейчас мы разбираем именно тот момент, когда я либо у всех колонок увижу, либо не увижу, то есть не разделяю на уровне строк.

[24:49] Дмитрий: Да. По строчкам типичный случай — это когда у тебя есть финансовые данные в таблице, есть разные департаменты, и каждый департамент может смотреть данные только по своим сотрудникам. Но ты не хочешь делать на каждый департамент, у которого динамическое количество и меняется достаточно часто, отдельные таблицы. Поэтому ты начинаешь делать логику ограничения доступа к конкретным наборам строк, относящихся к департаментам, чтобы руководитель каждого департамента мог смотреть только свои строчки, но таблица при этом одна. Логика полезная, но в бесплатной версии, по крайней мере сейчас, такой функциональности я не видел.

[25:30] Дмитрий: Прилетел у нас запрос, авторизовали и так далее, и дальше, как и в случае с записью, у нас есть логика координации. То есть кто-то должен пообщаться с несколькими репликами, где хранятся данные, вообще найти реплики и как-то сформировать результат, чтобы ответить клиенту. Эта логика в Cassandra называется координатор. Вспоминаем, что любая нода может быть координатором: запрос можно прислать любой, она разберётся, откуда брать данные. Но выгоднее, конечно, прислать такой запрос той ноде, в которой эти данные локально есть, чтобы она просто меньше по сети бегала. Опять-таки, это логика, которая есть и в клиенте: она общая для чтения и записи. Если мы знаем partition-ключ, на который вставляем строчку или для которого хотим прочитать набор строчек, мы можем вычислить, какие ноды содержат эти данные, и как клиент послать запрос в одну из таких нод.

[26:32] Дмитрий: Так что у нас начинает работать логика координатора, и она должна поговорить с репликами, чтобы получить данные. И тут мы снова вспоминаем такую штуку, как consistency level. Когда мы делали записи на ноды, мы указывали параметр consistency level, который говорил, сколько нод должны нам успешно ответить, чтобы мы сказали клиенту, что операция записи end-to-end прошла успешно. И наиболее частый параметр — consistency level LOCAL_QUORUM, то есть большинство нод в локальном дата-центре. Обычно вы будете иметь три реплики: одна — слишком мало, две — не очень удобно в плане отсутствия отказоустойчивости, три — в самый раз. Следующий уровень — пять, но пять — это уже достаточно большие накладные расходы на хранение данных, поэтому так делают не очень многие. Типичное значение — три реплики для каждой записи в дата-центре. Но если дата-центров два, то у вас получится шесть реплик — более чем достаточно.

[27:38] Дмитрий: Для записи мы определяли, сколько реплик отвечает. Для чтения этот параметр означает немножко другое: сколько реплик мы хотим опросить. И значения такие же — например, ONE, TWO, ALL (типа «все ноды хочу расспросить») или опять-таки LOCAL_QUORUM. И тут начинается интересная consistency-математика. Попытайтесь теперь представить себе такие диаграммы Венна, по-моему, это так называлось.

[28:09] Александр: Подожди, Сань. Ты вот некоторое время назад мучился с теорией множеств для консенсуса. Это такой алгоритм консенсуса для детского сада. Сейчас мы с тобой ещё раз пройдёмся в очень простом, простейшем варианте.

[28:28] Александр: Так, внимание, сейчас мы разбираем консенсус в подкасте. Немножко опять придётся про множества в голове подумать — про теорию множеств и как они там разными углами пересекаются.

[28:42] Дмитрий: Вот у нас есть полный набор нод. Давай возьмём число 11. Полный набор, весь кластер. Но поскольку мы сейчас читаем данные по конкретному partition-ключу, ограничимся для простоты случаями, когда partition-ключ известен. Это основной кейс, под который мы должны рассчитывать нашу логику. Поэтому всё ограничивается replication factor — сколько у нас реплик. Забудем пока про другой дата-центр, пусть у нас только один. Как я сказал, наиболее типичный случай — три реплики в дата-центре: каждая строчка представлена в трёх экземплярах на трёх разных нодах. Поэтому общее количество реплик — три. Я говорил «11», имея в виду весь кластер, а теперь из этих 11 отбрасываем 8 и оставляем 3 — это те ноды, где реально лежат данные для того partition-ключа, который мы указываем в запросе. Наш координатор об этом знает.

[29:53] Дмитрий: И вот у нас была операция записи, которой мы указали какой-то параметр. Обозначим его буквой W — обычно я так пишу. Например, 2, это был LOCAL_QUORUM. Две из трёх ответили, что запись прошла успешно. То есть отправили на все три, но что ответила третья — неизвестно, мы её не ждали. Может, она вообще была выключена; может, запрос прилетел к ней слишком поздно, и она сказала: «извини, он уже затаймаутился, обслуживать не буду». Но мы точно знаем, что две ноды ответили. То есть из этих трёх кружков на нашей диаграмме два попали во множество. Запись подтверждена. А дальше мы начинаем операцию чтения. И читать мы можем с разного количества нод. И тут прямо один в один пересечение с идеями distributed consensus: мы хотим, чтобы множество нод, в которые мы пишем, и множество нод, из которых читаем, пересеклись хотя бы в одной точке.

[30:54] Александр: Так, вот тут важно. Получается, что для того, чтобы прочитать последние, консистентные данные, нам необходимо выполнить условие, при котором те ноды, на которые мы записали данные…

[31:11] Дмитрий: …которые сказали нам, что запись прошла успешно.

[31:14] Александр: Да, мы пишем на все, но есть некое подмножество нод, которые подтвердили нам успех. Хотим, чтобы это множество хоть в одной точке пересеклось с множеством нод, из которых будем читать.

[31:28] Дмитрий: Да, потому что тогда хотя бы одна нода из тех, из которых мы прочитали, будет обладать последними данными.

[31:40] Александр: Да, всё именно так. То есть мы хотим добиться консистентности уровня «всегда вижу последнюю запись».

[31:46] Дмитрий: Да. Это не обязательно — вы в принципе можете это неравенство нарушать, но тогда потенциально будете иметь ситуацию, что читаете старые, устаревшие данные. В каких-то ситуациях вас это может устраивать, и вы можете специально читать поменьше нод — это будет эффективнее, но консистентность такого уровня вы потеряете. Чаще всего в большинстве сценариев люди пытаются добиться этой консистентности, и поэтому эти два множества должны хотя бы в одной точке пересечься. Итак, у нас есть W — write, сколько нод должны подтвердить запись. В нашем случае — 2, поскольку LOCAL_QUORUM от 3 это 2 (большинство из трёх — двойка). Теперь количество чтений. Если у нас 3 ноды и мы возьмём одну ноду для чтения, может получиться, что эта одна нода оказалась не из тех двух, что подтвердили запись, — пересечения не произошло, и мы прочитали старые данные. То есть одной ноды недостаточно. Но если мы возьмём для чтения две ноды — то есть опять LOCAL_QUORUM, — то пересечение будет всегда. Это принцип Дирихле, или через теорию множеств: если для трёх точек есть два множества по две точки каждое, хотя бы одно пересечение они между собой иметь будут.

[32:44] Александр: Да нарисуйте два кружка — и все эти теории и умные слова. Просто два кружка, они пересекутся в одном месте.

[33:30] Дмитрий: Вот, и получается, что у нас есть количество нод, которые мы читаем, — обозначим буквой R, — и получилось три числа. N — число реплик, R — число чтений, W — число записей. Можно выписать такое неравенство, верное для любого количества реплик: read + write > N. То есть количество нод, которые мы записали с успехом, плюс количество нод, которые прочитали, в сумме больше числа реплик. Это означает, что хотя бы одна точка при пересечении всегда будет. Вы можете, например, записать во все три ноды — сделать write ALL, — тогда для чтения вам достаточно одной, любой ноды: если записали во все, то какую ноду ни возьми, пересечение есть. Записали в три, прочитали из одной: 3 плюс 1 равно 4, 4 больше 3 — общего количества реплик. Вот это неравенство, наверное, основное в понимании consistency, которую даёт Cassandra.

[35:02] Дмитрий: В литературе по Cassandra это называется strong consistency. Тут, конечно, надо быть осторожным с терминами, потому что под этими словами могут понимать много чего. Но в Cassandra это означает, что чтение видит последнюю запись. И чаще всего вы такое хотите. И чаще всего для простоты мы берём, что чтение и запись — это LOCAL_QUORUM. Чтение — большинство, запись — большинство: в пересечении всегда дадут хотя бы одну точку, и мы будем читать то, что только что записали. Для любого replication factor — возьмите вы 5, 4, 6, неважно.

[35:20] Александр: Ну да, то есть replication factor — это по факту не про консистентность, а про долговечность: если мы данные записали, то не потеряем их даже при выходе одной-двух нод из топологии — смотря какой у нас, конечно, replication factor. То есть это просто сохранность данных. А вот когда мы говорим про так называемую strong consistency в рамках Cassandra, то это в рамках одной партиции. Потому что вообще говоря можно написать такой запрос в SQL — в стандартном, допустим, каком-нибудь Postgres, — который делает JOIN из пяти таблиц по разным данным, и в одну таблицу кто-то пишет с транзакцией уровня dirty read, а мы читаем с таким же уровнем изоляции. Естественно, мы можем те данные не увидеть. Это уже совсем другая консистентность — тоже консистентность, но в рамках всех данных, всей модели, которую мы смоделировали. А тут мы говорим конкретно про партицию: что происходит в других таблицах, keyspace’ах, партициях — мы вообще не знаем. В рамках одной конкретной записи мы сейчас разговариваем.

[36:34] Александр: И ты очень хороший рецепт дал этим read + write > N. То есть множество чтений плюс множество записей больше, чем общее множество реплик. Таким образом, если у нас реплик 3, то читать из 2 и писать в 2 достаточно, чтобы достичь такого уровня консистентности. Но не обязательно выбирать по 2. Можно, например, читать из 3 и писать в 1 — тогда мы тоже получим. То есть мы, как инженеры, настраивающие эти системы, можем этот слайдер двигать, и у нас несколько слайдеров. Что нам важнее? Мы много читаем? Тогда, наверное, читать стоит с меньшего количества нод, а пишем редко — ну тогда давайте будем писать во все, а читать из одной; так мы на чтение немножко оптимизируемся. Или наоборот. Понимание этой формулы — это не только готовый рецепт, но и уровень понимания системы, экспертизы, которая позволяет строить и управлять нагрузкой. Мне кажется, это важно было подчеркнуть.

[37:58] Дмитрий: Да, и из опыта: чаще всего вы всё-таки делаете read и write как LOCAL_QUORUM. Это наиболее удобный случай, когда вы пишете OLTP-нагрузку — приложение в онлайн-режиме должно быстро делать и чтение, и запись, чтобы ответить клиенту в рамках одной бизнес-операции. Поэтому нужно, чтобы и чтение, и запись относительно быстро работали и переживали падения нод. В этом плане вот этот дефолтный режим — и чтение, и запись через LOCAL_QUORUM — околооптимален для таких сценариев. Но если у вас есть, например, batch-обработка, где вы что-то читаете из Kafka, из какого-нибудь персистентного хранилища, и пишете в Cassandra, а потом кто-то активно всё это вычитывает, — то вам выгоднее сделать так, чтобы чтения были дешёвыми: делаете для чтения consistency level ONE, а запись идёт в фоне. Если что-то отвалится, всегда можно повторить фоновую операцию, перезаписать, поэтому запись можно сделать consistency level ALL: пусть пишут все реплики, а если какие-то отвалятся — не страшно, повторим позже, и рано или поздно всё запишем в кластер, зато чтения, если их много, будут дешёвыми. Тут можно играть с параметрами, но надо понимать, какие гарантии — в том числе гарантии надёжности для чтений и записей — вы хотите иметь. То есть это про надёжность и про производительность.

[39:38] Дмитрий: Итак, мы находимся в координаторе, и нам прилетел запрос с каким-то consistency level, который говорит, сколько нод мы хотим читать. Допустим, наш типичный случай — LOCAL_QUORUM, то есть надо прочитать с двух реплик, а всего реплик три. И вопрос: окей, мы знаем, что надо прочитать два из трёх, а какие реплики выбрать? Казалось бы, простой вопрос, но весьма нетривиальный, потому что тут могут быть различные факторы — примерно как в клиенте, когда мы обсуждали балансировку запроса между разными нодами, там было много-много аспектов. Мы хотим, наверное, чтобы отвечала та нода, которая содержит данные, но у нас одна нода локальная, а другие удалённые. Наверное, локальная нода быстрее ответит сама себе, чем удалённая, поэтому её, возможно, стоит предпочесть. С другой стороны — а вдруг какая-то из других реплик сейчас перегружена по какой-то причине? Может, там тяжёлые запросы выполняются, или с виртуальной машиной что-то не так — плохой сосед оказался и мешает работать, или с JVM что-то не так — недонастроили garbage collector, и он ушёл в stop-the-world надолго? То есть мы хотим выбрать ещё и реплики, которые сейчас отвечают быстро.

[41:17] Дмитрий: И вот в Cassandra есть логика с очень интересным названием. Не знаю, кстати, почему, но эта штука называется snitch. Как в Гарри Поттере — вот этот шарик с крылышками, за которым Гарри Поттер гонялся, снитч. Логика, которая при чтении говорит, какие ноды нам стоит выбирать, называется snitch. Там есть разные алгоритмы, и по умолчанию включён так называемый dynamic snitch — не помню точное название, dynamic network snitch, — который отслеживает среднее время ответа других реплик на ваши запросы. С точки зрения координатора я посылаю в них запросы, и они за какое-то время отвечают. И я хочу с большей вероятностью выбирать те ноды, которые отвечают мне быстрее. Я сказал «среднее», на самом деле там перцентиль — пятидесятый, то есть медиана. Это не то же самое, что среднее.

[42:24] Александр: Медиана — это точка на графике, а среднее — значение, являющееся суммой, разделённой на количество, арифметическое среднее этих точек. И они часто не совпадают. Это очень легко понять: берёшь среднюю зарплату в любом исследовании — средняя зарплата это сумма, и когда у тебя люди получают очень мало и очень много, средняя будет нормальная. А медианная зарплата — это как раз сколько получает реальный человек: половина меньше, половина больше. И это всегда очень разные цифры, забавно наблюдать.

[43:20] Дмитрий: Да, тут можно прямо конкретный пример представить, чтобы совсем стало понятно. У вас есть 10 человек. Девять получают зарплату 10 рублей, а последний — условно, миллион. Медиана — когда вы сортируете числа по порядку, как в массиве, — это средняя точка массива, там окажется вот эта десятка. А если вы посчитаете среднее арифметическое: 10 умножить на 9 плюс миллион, и поделить на 10, — то будет условно 100 тысяч с копеечкой. Так что разница: либо 10 рублей, либо 100 тысяч.

[44:12] Дмитрий: Так что в этом плане время отклика серверов — не очень равномерная величина, там распределение обычно достаточно хитрое, многомодальное. Если вы представите график вероятности времени ответа сервера: равномерное распределение — это когда вы можете получить любое время с равной вероятностью, такого обычно не бывает. Обычно бывает так, что довольно много запросов отвечают достаточно быстро — есть горб в начале распределения, — а потом есть горб ближе к концу, когда случаются всякие неприятности, запросы застревают в очередях, вылезают какие-нибудь внешние эффекты, и возникает второй горб очень медленных ответов. Называется бимодальное, или многомодальное распределение, когда на графике вероятности несколько таких горбов. Поэтому среднее вам тут совсем ничего не скажет — будет среднее этих двух горбов. Лучше смотреть на более сложные метрики, хотя бы перцентили. Они тоже не панацея, но чуть более правдивую картину дают.

[45:28] Дмитрий: Итак, у нас есть в Cassandra эта гека — snitch, — которая говорит нам, какие ноды стоит выбирать. И вариант по умолчанию учитывает, как эти ноды отвечали в прошлом. Гипотеза: если нода отвечала быстро в прошлом, согласно нашим метрикам, то, наверное, и в будущем будет отвечать быстро. Понятно, что не всегда эта гипотеза работает, но работаем с чем есть. В этой области были попытки сделать более точные предсказания, использовать дополнительные статистические данные. К сожалению, в Cassandra эти алгоритмы пока не реализовали. Но интересный факт: была работа в начале 2000-х, где на примере Cassandra нарисовали новый алгоритм, как выбирать эти реплики так, чтобы даже в условиях, когда случаются garbage-collection-паузы и прочее, всё работало эффективнее. Статистически там получалось лучше. Алгоритм я рассказывать не буду, он достаточно сложный — не архисложный, чтобы его не понять, но сложный, чтобы объяснить без слайдов. Его нет в Cassandra на текущий момент, но эту статью реализовали в Elasticsearch. Можно пойти в Elasticsearch и найти там этот алгоритм.

[46:54] Александр: То есть статья, написанная как бы для Cassandra, где ресёрчеры использовали Cassandra в качестве модельного примера.

[47:06] Дмитрий: Да. Они показали: вот смотрите, на примере Cassandra мы тут свою флоу сделали, покрутили, и получился такой интересный академический результат, но саму Cassandra это не законтрибьютили, и оно так бы и не попало в итоге. Но эту работу подхватил Elasticsearch и внёс в промышленную эксплуатацию.

[47:32] Александр: Круто, прикольно такое. Переопылили, к сожалению, не в ту сторону — Cassandra это пока не затронуло.

[47:42] Дмитрий: Но, может быть, всё-таки и туда мы это засунем. Чтобы у кого-то дошли руки.

[47:47] Александр: Да, это, кстати, разница между исследовательским программированием и реальным промышленным. Одно дело — провести исследование и показать, что в большинстве случаев это работает эффективнее, а другое — довести до продакшена. Вот вам яркий пример.

[48:06] Дмитрий: Да, очень много идей на самом деле пропадает посередине. Часто бывает, что те, кто занимается академическими исследованиями, — им задача написать статью…

[48:15] Александр: Не умеют пользоваться git.

[48:19] Дмитрий: Нет, скорее умеют. Скорее у них нет долгосрочного интереса заниматься поддержкой какого-то open-source-проекта. Их цель — опубликовать статью, показать результат, и дальше интерес потерян. А промышленное программирование — это когда ты довёл фичу до продакшена и потом её ещё поддерживаешь. Это марафон. Одно — спринт, другое — марафон. Спринтеры не любят бегать марафоны.

[48:51] Дмитрий: Так что есть разные способы это делать, но в Cassandra какой есть, такой и есть. Он не всегда идеален, у него есть иногда не очень приятные эффекты. Он может, например, входить в автоколебания. Система, которая выбирает по прошлому значению наиболее быстрые ноды, приведёт к тому, что в эти быстрые ноды вы будете посылать больше запросов — ну они же лучше, значит, надо их побольше использовать. Но если ты их больше используешь, ты их сильнее нагружаешь. А если сильнее нагружаешь, они, наверное, начинают отвечать медленнее — просто от того, что нагрузка больше. И эти ноды уже не такие привлекательные, ты выбираешь другие. С них нагрузка спадает, теперь любимчики — другие, и они снова становятся лучше. И ты можешь войти в такой синусоидальный ритм, когда нода то лучше, то хуже, и нагрузка то лучше, то хуже, и эти два эффекта в противофазе бьются — получается автоколебательная система.

[49:56] Александр: Главное, чтобы в резонанс не вошла.

[49:58] Дмитрий: Да.

[50:02] Дмитрий: Так. Мы выбрали какие-то две реплики и послали в них запросы на чтение. Дальше пока спускаться не будем. Они как-то эти запросы обработали, ответили нам результатами. И тут снова может быть ситуация: а вдруг какая-то из этих реплик не ответила достаточно быстро? В начале выполнения запроса она была жива, а потом затормозила, сдохла, ещё что-то с ней случилось. И мы не хотим снова отвечать клиенту ошибкой про недобор реплик, таймаутом и так далее. А у нас есть ещё третья реплика — вспоминаем про неё. Мы же тоже можем её попросить что-нибудь прочитать. Поэтому здесь мы можем снова сделать ретрай на эту третью реплику. Вспоминаем наш разговор про драйвер, где мы говорили, что просто ретрай неинтересен, потому что ждать очень долго. И можем сделать спекулятивный ретрай, когда мы подождали некое не очень большое время, за которое проходит большинство запросов, — например, 99% запросов, — а всё, что выше, считаем: с большой вероятностью ответов от реплик мы уже не получим, надо спрашивать у третьей оставшейся реплики.

[51:18] Дмитрий: Вот этот спекулятивный ретрай в Cassandra-сервере, в координаторе, реализован. У нас есть два места, где делается спекулятивный ретрай. Вы можете сделать его на стороне клиента — чтобы делать спекулятивный ретрай на разные координаторы, если координатор затормозил. А сам координатор может сделать спекулятивный ретрай, когда общается со своими репликами. За счёт того, что и координаторов много, и реплик много, мы можем себе позволить такие ретраи. Это настраивается per-table. Если вы сделаете DESCRIBE TABLE для вашей таблицы, там, кроме основной информации о колонках, есть суффикс, который описывает разные параметры настройки. Например, каким compaction-алгоритмом собирать эту таблицу. И в том числе есть параметр speculative retry — как делать спекулятивный ретрай. По умолчанию, даже если вы ничего не настраивали, вы увидите там 99percentile. Это дефолтный режим работы. Если мы не получили достаточного количества ответов от реплик и прошло время, равное 99-му перцентилю, то координатор Cassandra сделает дополнительный запрос в оставшуюся реплику и попробует за счёт этого ответить вам побыстрее. Вы можете этот параметр поменять, можете вообще выключить, если вам это неинтересно, но такое бывает редко.

[52:54] Дмитрий: Там появился, по-моему, в четвёртой версии полезный вариант, который позволяет задать этот параметр более сложным образом. У вас может случиться, что оборудование деградировало — что-то пошло не так, и нода начала отвечать всё медленнее, или данных стало больше, — и этот 99-й перцентиль тоже начинает расти. А у вас ожидание, что система должна отвечать за константное время: есть SLA на ваш сервис, он как был, так и остался. И мы хотим не так сильно зависеть от этого контекста, поэтому на уровне настройки сейчас можно написать минимум или максимум от двух параметров. Например: делай спекулятивный ретрай либо когда мы пересекли 99-й перцентиль, либо по константному значению. Если мы знаем, что наш сервис ожидает выполнения запроса не больше чем за 20 миллисекунд, то можем написать, что спекулятивный ретрай выполняется по минимуму из 99-го перцентиля и 20 миллисекунд.

[53:47] Александр: Слушай, а можно настроить так, чтобы сразу отправлять запрос во все реплики?

[53:55] Дмитрий: Что нам стоит — можешь сказать «сразу». У этого варианта есть свои плюсы и минусы. Минусы понятны: мы загружаем систему побольше. Мы так-то две реплики опрашивали, а теперь будем всё время три — в целом больше работы делаем. Но в зависимости от того, какая нагрузка на систему и как реплики себя чувствуют… Если мы отправляем запрос сразу во все, то мы отправляем фиксированное количество запросов, то есть нагрузка пропорциональна всегда одинаково. Это более предсказуемое поведение системы.

[54:54] Александр: Да, я про эту идею слышал в двух местах. Первое — ребята из Одноклассников в своё время рассказывали. Олег и Анастасия, когда делали доклады про Cassandra лет пять назад или чуть больше, — у них была как раз такая идея, что лучше мы будем читать со всех реплик, чтобы, что бы ни происходило, у нас нагрузка на ноды не плавала. В случае каких-то нехороших изменений в системе мы не грузим её ещё сильнее и тем самым не получаем негативную обратную связь: пошло что-то не так, а мы стали ноды сильнее грузить, и им совсем плохо стало. И похожая идея встречается в докладах Amazon, когда они говорят про отказоустойчивость. Системы, в которых уровень нагрузки от запроса динамически не меняется от характеристик системы, не связанных с самой нагрузкой, — такие системы более отказоустойчивы, потому что у вас нет дополнительного фактора обратной связи.

[56:07] Дмитрий: Это, конечно, не защищает вас от всех проблем — вспоминаем недавний outage Amazon.

[56:12] Александр: Да, эту улыбку на твоём лице после слова «Amazon» я подметил.

[56:17] Дмитрий: Да, теперь, когда мы говорим про Amazon и отказоустойчивость, все внутри улыбаются, а потом делают вид, что об этом не подумали. Но там история всё-таки частично связана с тем, что все очень понадеялись на то, что один регион всегда доступен. Это плохая надежда. Есть разговоры, есть реальные устройства систем, и они, к сожалению, частенько расходятся.

[56:44] Дмитрий: Есть случаи, когда мы действительно хотим все реплики сразу опросить, и какие бы две ни ответили быстрее — мы получим результат. Это будет работать быстрее и более предсказуемо в плане нагрузки на серверы. Поскольку мы опрашиваем не две реплики, а всегда три — на одну треть больше работы в целом на кластер.

[57:09] Дмитрий: Вот эти спекулятивные ретраи мы все сделали, и нам либо в итоге пришли ответы от реплик, и мы будем формировать результат, либо они не пришли, и мы отвечаем клиенту ошибкой. Тут опять-таки по ошибке можно понять, что случилось: Cassandra пошлёт вам разный код ошибок, и это будут разные исключения на стороне драйвера. Вы можете увидеть ошибку в стиле: «ты попросил consistency level LOCAL_QUORUM, а я вот смотрю — у меня живых реплик всего одна, а две остальные точно мёртвые». А как я это узнал — так gossip-протокол, о котором мы говорили в прошлый раз про запись, точно так же и тут применим.

[58:01] Александр: Одна бабка сказала.

[58:03] Дмитрий: Да, одна бабка сказала, что у меня в кластере одна реплика живая, а ты просишь читать с двух. Не могу, извини, вот тебе ошибка. Не читал и даже пытаться не буду. А может быть ошибка в стиле: говорят, ноды-то все живые, мы в них запросы послали, а они что-то не отвечают, обманули нас. И тогда мы вернём клиенту ошибку: «извини, ты попросил LOCAL_QUORUM, а ответила мне только ноль или одна реплика из трёх». Недобор. И в ошибке это будет видно — там, по-моему, даже в дополнительных данных передаётся, какие ноды ответили, а какие нет. Иногда полезно для троублшутинга. Но это надо явно из исключения вытаскивать, оно в логи на клиенте, по-моему, само не пишется.

[58:49] Дмитрий: Но допустим, мы разбираем успешный случай, когда реплики нам что-нибудь ответили. И тут есть очень частое заблуждение. Раз мы говорим про всякие LOCAL_QUORUM, кворумы, голосование, большинство, — наверное, тут у нас… Вот есть какие-то реплики, у одной реплики может быть одно значение, у другой другое. Кто прав?

[59:18] Александр: Да, часто заблуждение в том, что люди думают: большинство, голосовать тут кто-то сейчас будет, как в протоколах с консенсусом, — какое значение правильное?

[59:32] Дмитрий: Нет, голосов тут нет. Вспоминаем: в основном случае мы опрашиваем две реплики, но одна скажет значение А, другая — значение Б. У каждого ответа голос один — некому тут голосовать. Поэтому здесь работает другая логика. Золотое правило Cassandra. Я его, кстати, упоминал при операции записи. Помнишь? Правило Cassandra, золотое, при операциях записи, которое ты упоминал.

[1:00:01] Александр: Вот это ты меня, конечно, застал врасплох. Не помню.

[1:00:05] Дмитрий: Насколько ты хорошо слушаешь. Помнишь, история была, когда мы записываем данные, и записи могут лететь в разном порядке, и нам надо их как-то смёржить между собой? А тут тоже похожая история: нам прилетели ответы от разных реплик, а нам нужно их смёржить и сформировать финальный результат — это и будет результат для клиента. Поэтому мёрж при записи и мёрж при чтении весьма похожи.

[1:00:34] Дмитрий: Да. И ты можешь всё это смёржить по таймстемпу: у кого таймстемп выше, тот и выиграл. Можно было бы всё так и реализовать, и в некоторых случаях в Cassandra так и работает. Но мы тут захотели сделать дополнительную оптимизацию. Это не недавняя вещь, она существует в Cassandra давно. Мы можем немножко попытаться пойти более оптимальным путём. А именно: в нашем базовом алгоритме мы запросили у каждой реплики данные, они прилетели на координатор, и мы начинаем их мёржить. Идея заключается в том, что в большинстве случаев, в здоровой системе, где всё нормально, наши реплики содержат одинаковые данные. Там мёржить особо нечего — любой ответ подойдёт. Поэтому на самом деле в базовом сценарии нам нужно не мёрж делать, а просто убедиться, что ответы действительно совпадают.

[1:01:28] Дмитрий: Поэтому от одной реплики — обычно от локальной — мы запросим полные данные. То есть мы, координатор, локально на свою собственную ноду сходим (поскольку запрос обычно прислали на эту локальную ноду) и запросим полные данные. Это называется data request в Cassandra. А от второй ноды мы будем просить не данные, которые надо гонять по сети, сериализовать, десериализовать, — мы запросим хэш от данных. От результата запроса, который мы сейчас выполняем. Мне как координатору прилетел хэш. Я могу от своих локальных данных тоже посчитать хэш и сравнить. Если хэши совпадают, значит данные на моей реплике и на соседней одинаковые. Я могу свои данные просто дать клиенту. И мне не надо со второй реплики вычитывать данные, мёржить, сравнивать — я могу сравнить просто хэши. То есть я уменьшаю накладные расходы.

[1:02:28] Александр: Слушай, а как этот хэш считается? Ведь данных-то может быть много, и посчитать хэш вот так в лоб, возможно, будет очень дорого.

[1:02:34] Дмитрий: Не очень дорого, но и не очень дёшево. К сожалению, в Cassandra для этой задачи был выбран не самый удачный алгоритм. Можешь даже угадать какой. Ну, мы тут уже упоминали MD5.

[1:02:50] Александр: Всё правильно, да. Чего изобретать много разных алгоритмов, если можно всё через один делать.

[1:02:54] Дмитрий: К сожалению, взяли MD5. И прямо на все данные берётся… Ну не на все, на которые ты запрашиваешь. Ты же запрашиваешь всё-таки не супергигантские объёмы, а результат запроса — вспоминаем, что у нас OLTP. Или там 10 строк вернулось, может, 1000 строк. Ты для них это дерево результата обходишь и считаешь хэш. Сейчас там MD5. И в текущей версии Cassandra это немножко пытались оптимизировать: сейчас, если в потрохах покопаться, можно найти, что используется альтернативный провайдер реализации MD5. В Java для криптографических алгоритмов всё сделано через точки расширения — есть специальный а-ля SPI для криптографии, когда можно реализовывать разные алгоритмы не только в самой JDK, но и во внешних библиотеках. Bouncy Castle, например, самая известная библиотека в этой области. И для хэшей то же самое. С относительно недавних пор в Cassandra стали использовать библиотеку от Amazon.

[1:03:56] Александр: Мы второй раз уже упоминаем Amazon в этом выпуске.

[1:04:04] Дмитрий: Amazon Corretto, в виде отдельной джарки, — доступная альтернативная реализация некоторых криптографических алгоритмов, более вылизанная, оптимизированная по сравнению с JDK-версией. Cassandra эту библиотеку к себе затянула, и там получилось чуть получше. Но всё равно, понятное дело, всё сильно зависит от того, какие у вас данные, сколько колонок, какие размеры полей, — у меня по флеймграфам получалось где-то процента три CPU на этот хэш. То есть если снять флеймграф для работающего Cassandra-процесса, где чтение и запись примерно 50 на 50, то где-нибудь процента три — условно, чтобы понимать оверхеды, — уходит в вычисление этих хэшей. Не слишком много, но заметно.

[1:05:05] Александр: Слушай, я согласен, что это довольно значимая штука. Ведь это же необязательная вещь — мы могли бы возвращать данные, нам не нужны криптографические алгоритмы. На другой чаше весов что? Альтернативное решение — давайте не будем считать хэш и будем возвращать эти 10 000 строчек координатору, и пускай он дальше принимает решение.

[1:05:31] Дмитрий: А оверхеды на сериализацию этих строчек на стороне реплики, десериализацию потом на стороне координатора и мёрж результата, скорее всего, будут дороже. То есть ты всё равно здесь выигрываешь. Тут если говорить про альтернативы, то альтернатива — скорее другой хэш использовать. Нам не нужен криптографический хэш, нам нужен просто хэш с достаточно низкой вероятностью коллизий. Есть всякие Murmur-хэши, о которых мы уже говорили, когда обсуждали consistent hashing ring; есть CityHash, xxHash — много разных алгоритмов, существенно лучших в плане производительности. И заход в эту сторону в Cassandra был: эту логику абстрагировали, вынесли в некие интерфейсы, чтобы можно было всё подменить.

[1:06:30] Дмитрий: Основная проблема в том, что этот MD5-хэш — часть протокола между нодами. Реплика возвращает тебе хэш, и ты, как координатор, должен посчитать хэш тем же алгоритмом, иначе они не совпадут. И как только ты введёшь новый хэш, у тебя появляется проблема нарушения обратной совместимости. В смешанном кластере, когда ты обновляешь Cassandra, эти хэши начнут расходиться, координатор начнёт страдать и думать, что всё не так. И ты должен что-то с этим делать — выдумать какие-нибудь алгоритмы хендшейка между нодами, чтобы они договорились, какой алгоритм использовать конкретно в этом запросе. Это всё сложно. По сути, тебе нужно выпустить промежуточный релиз, предназначение которого — подружить ежа с ужом, а потом забыть про ежа и пойти с ужиком. То есть это сложная инженерная практическая задача, которую, блин, неохота, честно говоря, решать.

[1:07:34] Александр: Да, то есть понятно, как это делать, это не составляет сверхбольших усилий с точки зрения сложной математики.

[1:07:46] Дмитрий: Всё ясно, как делать. Но вот ради этих 3% пока никто не засучил рукава, не нашёл время, чтобы этим заниматься. Так что пока оно, к сожалению, в таком состоянии. Может, я когда-нибудь до этого доберусь, может, нет — не буду пока обещать.

[1:08:02] Дмитрий: Работает это сейчас так, что одна реплика отвечает данными — обычно это локальная реплика, — а все остальные реплики отвечают хэшами. Мы считаем хэш от наших данных, сравниваем хэши между собой. Если всё совпало, значит, все реплики в синке, можно ответить клиенту результатом. Если consistency level единичка, ONE, то никакой второй реплики нет — мы просто с одной реплики читаем данные и отвечаем клиенту, консистентность проверять не с кем. Если это consistency level ALL, например, в нашем случае с replication factor 3, то две реплики дадут хэш, а одна — данные, от которых мы тоже посчитаем хэш, и эти три хэша сравним между собой. Если всё совпало — хорошо.

[1:08:52] Дмитрий: А вот если не совпало — тут интересный механизм. Мы запрашиваем с каждой реплики полные данные. Нам теперь хэша уже не хватает — нужно выстроить последнюю, наиправильнейшую версию строчки, которую нас попросили. Мы с каждой реплики запросили реальные данные и в координаторе их смёржили тем же способом, по таймстемпу. И получили эту актуальную строчку. То есть, допустим, запись по какой-то причине не долетела до всех реплик, и на одной реплике были старые данные, а на других — новые. Новые данные при мёрже победят, и мы можем ответить клиенту новыми данными. Но это не всё. На самом деле в этот момент в Cassandra включается механизм автовосстановления. Cassandra будет пытаться не только ответить клиенту актуальными данными, но и сама починить свои устаревшие данные на репликах. Это называется механизм read-repair. После мёржа мы эти данные обратно отошлём во все реплики, чтобы они их записали. Вот интересный факт: мы, казалось бы, делали чтение, но если в момент чтения прочитанные строчки были не засинхронизированы, мы попытаемся их засинхронизировать обратно и сделаем запись. То есть операция чтения с точки зрения клиента может породить на стороне Cassandra запись.

[1:10:29] Александр: Прикол. Но, на самом деле, такой, знаешь, трюк. А насколько часто он используется? Мне кажется, крайне редко.

[1:10:39] Дмитрий: Такие ситуации в целом редко возникают. Расчёт на то, что такие ситуации должны быть редки.

[1:10:45] Александр: Да, они слишком дорогие.

[1:10:49] Дмитрий: Был интересный случай, когда была бага в расчёте хэша. Когда мы для результирующих данных обходили это дерево, мы там, по-моему, то ли что-то обходили неправильно, то ли какую-то колонку с таймстемпом неправильно посчитали. В общем, была бага в Cassandra, которая возвращала неправильный хэш от результата. И получилось так, что всё время надо реперить данные. Из-за баги люди поставили новую версию Cassandra — не помню, в какой это версии было, то ли в какой-то из четвёрок, то ли в начальных пятёрках, — и заметили, что стало очень много этих read-repair’ов. А это ты замечаешь, потому что read-repair — часть операции чтения. У тебя был один раунд: прочитал с реплик, ответил, — а теперь стало три раунда. Ты прочитал, хэш сказал, что не совпал; читаешь с каждой реплики реальные данные — это второй раунд общения с репликами; а потом ещё записываешь результат, и только после этого отвечаешь клиенту. То есть три раунда вместо одного, и latency будет раза в два-три выше. Если это случается очень часто, ты видишь, что для большинства запросов время операции чтения возросло в три раза. Но это была бага, её пофиксили уже давно, вряд ли те, кто использует актуальные версии, с ней встретятся. Но надо понимать, что такие эффекты могут быть.

[1:12:22] Александр: Это мы с вами проверили и рассказали про алгоритм восстановления данных на случай, если мы читаем записи с реплик, и они по какой-то причине разъехались. И в итоге результат возвращаем клиенту, клиент его обрабатывает.

[1:12:44] Александр: И тут, прежде чем мы перейдём к уровню хранения данных на конкретной ноде, надо подчеркнуть ещё один аспект. Когда мы читаем данные с сервера, может вернуться одна строчка, может десяток. А вдруг наш запрос вернёт 50 тысяч строчек? Мы что, должны их все сразу отправить клиенту, потратить кучу памяти на то, чтобы всё это закодировать, переслать между репликами, а потом ещё драйвер будет всё это декодировать? Это не очень безопасный способ в плане потребления памяти, можно и с out of memory упасть.

[1:13:16] Дмитрий: С одной стороны, мы хотим, чтобы эти запросы имели ограниченный лимит по строчкам. Но с другой стороны, мы хотим, чтобы, если данных много, клиент всё-таки мог их обрабатывать. Поэтому для этого есть механизм, по-английски он называется paging.

[1:13:34] Александр: Пагинация.

[1:13:35] Дмитрий: Да, это калька просто. «Страничность» не звучит.

[1:13:42] Дмитрий: Мы хотим вычитывать данные с сервера некими чанками, кусочками строк относительно небольших размеров. И эта функциональность есть. Более того, она по умолчанию включена — надо ещё постараться её выключить. Когда вы делаете запрос в драйвере, в реальности у этого запроса есть скрытый параметр, сколько строчек вы хотите прочитать за раз. Запрос летит в координатор, координатор видит, что стоит параметр, по-моему, 5000 строк (в разных версиях драйвера может быть по-разному). А там, допустим, лежит 10 000 или 50 000 строк. В этом случае мы, координатор, вычитываем с реплик эти 5000 строк и возвращаем клиенту. А дальше стоит задача: клиент эти 5000 строк обработал, хочет читать дальше. Но мы не хотим, чтобы сервер помнил, где он остановился, где в этот момент читал. У нас любой координатор может отвечать на этот вопрос, и мы не хотим держать данные о состоянии процесса чтения на сервере. Поэтому что мы сделаем? Мы вернём их клиенту. Вместе с данными мы вернём некие метаданные, которые помогут понять, с какого места надо продолжать читать. Указатель.

[1:15:10] Дмитрий: Мы читали из книжки некий текст, прочитали страницу, положили закладку. В этой закладке написан номер страницы, на которой мы остановились. И когда клиент в следующий раз придёт и скажет, что хочет ещё данных, он вернёт нам обратно эти метаданные. Я не помню, как этот объект в клиенте называется — paging metadata или paging info, — в общем, некая сущность, которая содержит данные о том, где мы остановились. По сути там закодировано что-то типа clustering-ключа или позиции в партиции, докуда мы дочитать успели. Мы читали-читали, и вот есть некая последняя запись, которую мы увидели, и её, очень упрощённо говоря, primary-ключ мы запомнили в качестве закладки.

[1:16:07] Александр: А теперь интересный вопрос. Что будет, если данные ушли вперёд в этот момент, и там что-то поменялось? То есть результат отображения запроса будет консистентный или всё-таки консистентность как-то теряется?

[1:16:21] Дмитрий: Нет, не будет. Вот в этом как раз тонкое место, известная в Cassandra проблема: консистентность мы обеспечиваем только на уровне конкретного запроса. Когда у вас запрос развалился на несколько страниц — это разные запросы. Если, грубо говоря, в тех старых данных, в диапазоне, который вы только что прочитали, произошли изменения… это тяжело реализовать. Поэтому тут, к сожалению, только держать это в уме, такие особенности.

[1:16:56] Александр: Ну, если бы была какая-то синхронизация на уровне глобального, гибридного таймстемпа… Эту проблему ты можешь решить только честным multi-version concurrency control. То есть когда у тебя хранится прямо срез данных на транзакцию или что-то подобное. Иначе тебе сложно этот консистентный view долго хранить на сервере и отвечать так, чтобы клиент мог в его рамках работать. То есть тут без транзакций уже не обойтись. Это будет такая read-only-транзакция, причём интерактивная: когда много раз ходишь к серверу и говоришь «я всё ещё читаю вот в той транзакции».

[1:17:43] Дмитрий: На самом деле даже в обычном SQL-сервере у тебя будет похожая проблема с дефолтными уровнями изоляции. Потому что почти всегда по умолчанию — это тоже компромисс между производительностью и консистентностью — у тебя стоит read committed. И ты читаешь с какой-то большущей таблицы данные от начала до конца — PostgreSQL, там, whatever, — бежишь, по кусочкам вычитываешь эти данные. Условно говоря, строчки пронумерованы от 1 до 10, ты этот диапазон прочитал, идёшь дальше. А прибежала транзакция в середине и поменяла какую-нибудь пятую строчку. Если ты перечитаешь снова эти диапазоны, то увидишь эти изменения — потому что read committed: транзакция закоммичена, значит, другие транзакции, которые не в процессе, эти изменения увидят. И тебе, на самом деле, нужна более строгая консистентность — то, что называется snapshot read, уровень изоляции, который не даёт тебе получить фантомные чтения.

[1:18:41] Александр: Да, то есть тебе нужно сказать: возвращай мне данные, как будто я заморозил базу на момент начала моей транзакции. То view, которое было в начале моей личной транзакции, я хочу видеть всё время, несмотря на то, что мир вокруг меняется.

[1:19:22] Дмитрий: Так что в Cassandra такое не получится сделать, тут просто приходится держать это в уме. Это одна особенность пейджинга, но другая, более важная, — этот пейджинг немножко другой по сравнению с тем, который мы привыкли видеть в реляционных базах данных. Обычно в реляционных базах пейджинг сделан по принципу OFFSET и LIMIT: мы хотим читать данные с такой-то позиции и столько-то. И можем рисовать странички — первая, вторая, третья, — и прыгать по разным страницам случайно. Можно представить навигационный интерфейс в браузере: вы открываете большой список, там появляются номера страниц, 1, 2, 3, 4, 5, и вы можете по ним тыкать. В Cassandra же механизм вида «туда-сюда»: я прочитал текущую страницу, могу прочитать следующую, могу нажать «далее, далее, далее». Но я не могу сказать «не далее, а на страницу номер 4, пожалуйста, перекинь». Cassandra пытается давать только тот механизм, который можно эффективно реализовать. А механизм с навигацией по страницам очень дорогой. Если данных очень много, бежать куда-то в конец весьма дорого: вам приходится на стороне SQL-сервера, например, при миллионе строчек, чтобы попасть на последнюю страницу, все эти строчки пропустить почти до самого конца и только конец, который вы попросили, реально прочитать. Можно найти статьи, что когда вы делаете пагинацию на очень больших объёмах данных, у вас будут проблемы в SQL-сервере. В Cassandra это решено радикально: там просто такой механизм не даётся. Даётся механизм «прочитал страницу — можешь получить следующую», и всё. Других вариантов нет. Это даёт гарантированную сложность, но может быть неудобно клиентам, привыкшим делать как хотят.

[1:21:49] Александр: Хотя, если тебе миллион данных, то от того, что ты перепрыгнешь на конкретную страницу, вероятность, что ты там найдёшь нужные данные, очень мала. То есть большой вопрос, зачем тебе такой интерфейс.

[1:22:02] Дмитрий: Значит, у нас есть механизм пагинации. Он, более того, встроен в драйвер CQL с точки зрения API, так что ты можешь его даже не заметить. Вот у тебя есть statement, который ты выполнил, — он вернул набор строк в виде некоторого объекта, ResultSet, по-моему, он называется. И ты дальше можешь использовать его как итератор, написать просто цикл for (row : resultSet), и каждый элемент — это строчка. То есть ResultSet имитирует интерфейс Iterable как дженерик — кстати, отличный интерфейс. И когда ты итерируешься, сервер реально вернул первые 5000 строчек, они лежат в памяти в драйвере, и ты итерируешься по ним в памяти. Добежал до конца этой страницы — драйвер послал следующий запрос в сервер, вернул следующую страницу, и ты снова итерируешься по этим строчкам в памяти.

[1:23:15] Дмитрий: Это такая немножко нетривиальная вещь, которая может иногда вводить в заблуждение. Вроде как со стороны клиента в коде ты сделал один запрос, выполнил execute один раз, а смотришь на метрики сервера — а там запросов несколько. Почему так? Потому что во время итерации ты, возможно, неявно запрашиваешь несколько страниц. По коду в клиенте одна операция, а по факту, когда ты итерировался по результатам, ты запрашивал страницы и запросов сделал несколько. Вот такая вещь, которая может немножко сломать мозг.

[1:23:54] Александр: Ну, это то, как работают абстракции. Собственно, на уровне ResultSet мы абстрагируемся от этих деталей, а в деталях может быть и много походов по сети, потому что реализация такая. Но интересно.

[1:24:09] Александр: Ну и вот мы проговорили логику координации, а теперь давай пойдём в конкретные реплики смотреть, как там реально эти данные вытаскиваются. Сейчас у нас была абстрактная логика, которая как-то из чёрных ящиков наших реплик доставала данные. Но мы же говорим про базу данных — там, наверное, как-то всё это работает под капотом, интересно. И тут нам надо сдуть пыль с нашего объекта под названием LSM-дерево и снова на него посмотреть, потому что данные реально хранятся в нашем LSM-дереве.

[1:24:42] Дмитрий: Вот у нас прилетел запрос на реплику. Вспоминаем, что у нас есть некая структура в памяти, промежуточный буфер memtable, который накапливает недавние записи, — там могут быть данные. И есть ряд файликов на диске, SSTable, которые содержат наши данные. И строчки, которые мы хотим вернуть, могут быть как в memtable, так и в каком-то количестве SSTable. Нам надо эффективно оттуда их достать и смёржить, потому что там может быть несколько вариантов: строчку сначала записали одной версией, потом ещё пять раз обновили, и эти изменения могли попасть в разные файлики. Так что нам опять-таки нужен этот мёрж, конечно же, по таймстемпу, и мы хотим сделать это эффективно.

[1:25:34] Дмитрий: Начнём мы с memtable. Вспоминаем, что это мапа, мапа, мапа. Partition-ключ у нас известен, поэтому первый уровень, первую мапу, мы проходим, просто запрашивая данные по ключу. И попадаем на второй уровень, где хранится набор данных по clustering-ключу, а дальше ниже — строчки. И тут вспоминаем, что это не hash map. Hash map дала бы нам только возможность попросить по конкретному partition-ключу, а иначе пришлось бы перебирать все значения. А у нас B-дерево, в частном случае сортированные мапы. В сортированной мапе мы можем эффективно делать запросы вида «дай мне диапазон ключей», «дай мне строчки по диапазону ключей» либо по их префиксу. Вот почему нам здесь нужно сортированное дерево — и для записи, и для чтения, чтобы эффективно обрабатывать такие запросы. Здесь особой другой магии нет. Понятное дело, что мы можем вычитывать не все колонки, а подмножество: в SELECT могли написать не звёздочку, а конкретный набор колонок. Соответственно, когда вытаскиваем строчки из третьего уровня, какие-то колонки надо отбросить, отфильтровать.

[1:26:54] Дмитрий: И второй аспект — хотя бы в контексте одной страницы мы хотим консистентность, не хотим, чтобы на этом уровне была частичная запись. Поэтому это B-дерево, как я говорил раньше, когда описывали механизм записи, — copy-on-write. Когда мы читаем, мы из этого дерева прочитаем некую консистентную версию: если мы вставляли batch, то либо увидим целиком новый набор строк, либо увидим старую строчку в рамках этой партиции. То есть если у вас есть возможность написать SELECT, в котором partition-ключ либо A, либо B — синтаксис partition_key IN (A, B) в Cassandra существует, — то это уже работать не будет. Весь механизм консистентности ограничен конкретной партицией: у каждой партиции своё независимое дерево, а сквозного version control нет — тут опять-таки мы должны уйти на уровень транзакций.

[1:28:04] Дмитрий: Итак, с памяти мы прочитали значения, но у нас ещё есть набор файликов на диске, где лежат другие значения. Нам нужно прочитать их с диска, при этом хотелось бы сделать это максимально эффективно, поменьше диск трогать. Что мы здесь можем сделать? Во-первых, такую оптимизацию. Если мы читаем конкретную строчку — указали полностью partition-ключ и clustering-ключ целиком, то есть ответ либо одна строчка, либо ничего, — и допустим, мы уже нашли строчку в памяти. То есть у нас уже есть кандидат на ответ. И у нас есть набор SSTable на диске. Мы можем — как помнишь, была история с tombstone’ами — использовать метаданные: в каждой SSTable на диске есть минимальный таймстемп и максимальный таймстемп. Если мы знаем, что ничего странного не происходило и в памяти наиболее свежие значения, а метаданные каждой SSTable показывают, что эти данные старее, то они уже не перетрут моё текущее значение, которое я нашёл в памяти. Я могу их все вообще пропустить. Это прямо отдельная ветка в коде Cassandra — query in timestamp order, — которая позволяет все остальные файлики вообще не читать. Или то же самое: мы нашли строчку в каком-то из файликов на диске, а все остальные файлики просто старее по записи, — можем их проигнорировать, отфильтровать и вернуть результат.

[1:29:40] Дмитрий: Но этот механизм не работает, если ты запрашиваешь несколько строчек, поскольку мы не знаем, сколько конкретно строчек там должно быть. Ты говоришь «дай мне записи по диапазону дат для этого пользователя», нашёл какие-то записи в памяти, — а какая гарантия, что это полный набор и между ними посерёдке не вклиниваются другие записи? Они где угодно могут быть расположены. Поэтому нам придётся полноценно искать всё это по диску. Если вы можете задать primary-ключ целиком, это будет работать эффективнее, чем если зададите только partition-ключ или диапазон clustering-ключей, — просто потому, что мы можем применить эту оптимизацию.

[1:30:32] Дмитрий: Допустим, таймстемпы не совпали, либо нам надо прочитать набор строчек, поэтому придётся читать набор файликов. Первое, что мы можем сделать, — попытаться отбросить некоторые файлики как неинтересные. Для каждого файлика у нас есть минимальный ключ и максимальный ключ, который в нём находится. Если наш искомый ключ не попадает в этот диапазон — в отрезке, относящемся к этому файлику, этого ключа нет, — значит, файлик можно выкинуть. Это первое, слой отбрасывания файликов. А дальше, допустим, в диапазон мы всё-таки попали. Тогда появляется новая структура данных, которая позволяет сказать: с одной стороны, возможно, в файлике наш ключ есть и придётся его читать, а с другой — она может сказать, что в файлике ключа точно нет, можешь не тратить ресурсы и просто его проигнорировать. Такая асимметричная вероятностная структура данных: ответ «да» она говорит с какой-то вероятностью ошибки, но её прелесть в том, что для работы ей требуется небольшой объём памяти. Мы можем эти данные загрузить в память и на диск вообще не ходить. Эта структура данных называется Bloom filter, фильтр Блума. Блум — это фамилия.

[1:31:31] Дмитрий: И давай попытаемся представить, как он работает на базовом уровне.

[1:32:01] Александр: Да, сейчас разбираемся, как работает Bloom filter. Господа, надо тоже немного воображения тут подключить.

[1:32:09] Дмитрий: Я для себя нашёл хорошую аналогию. Вспоминаем, как устроена хэш-таблица в Java, например HashMap. Мы берём ключ, считаем от него хэш, и на первом уровне есть массив с набором bucket’ов, и хэш нам говорит, в какую ячейку этого массива, в какой bucket мы попали. Вот первый этап работы Bloom filter ровно такой же: мы посчитали хэш от ключа, у нас есть массив, и этот хэш говорит, в какую ячейку мы попали. Нам надо узнать ответ «да»/«нет» от Bloom filter: он должен сказать либо что запись в структуре на диске есть, либо что её нет. То есть нам нужен ответ 0 или 1. У нас value в этой хэш-мапе нет, есть только хэш-сет, скорее. Хэш-сет, но вероятностный, где один ответ стопроцентный, а другой — «может быть».

[1:33:06] Дмитрий: Да. И соответственно, в этом массиве bucket’ов мы вместо того, чтобы писать какие-то значения, будем хранить битики — 0 либо 1. 0 означает, что ничего нет; 1 означает, что, наверное, есть. Если мы увидели 0 в результате такой операции — всё, можем сказать, никто ничего не писал. А вот если увидели 1 — это то же самое, что в хэш-мапе была бы коллизия: либо ты реально по этому ключу что-то вставлял, и в этот bucket попало значение, либо был другой ключ, у которого хэш совпал с твоим, и попал тоже в эту же ячейку. В реальной хэш-мапе это решилось бы через механизм разрешения коллизий — там была бы цепочка bucket’ов или ещё какой-то способ. Но в нашем случае это только 0 и 1. Поэтому если там написано 1 — ещё не гарантия, что 1 записали ровно для нашего ключа.

[1:33:57] Александр: Угу. Нам надо что-то ещё сделать, чтобы на 100% узнать, есть там наш ключ или нет.

[1:33:59] Дмитрий: Чтобы узнать на 100%, нам нужно только на диск ходить, других вариантов нет, но мы хотим снизить… На самом деле вот, всё, Bloom filter. Если очень обобщённо, это и есть Bloom filter.

[1:34:30] Дмитрий: Да. И его прелесть в чём? Что он покрывает огромное количество данных, но по сути эти данные не хранит, поэтому и применим. И он очень лёгкий, битсет такой. При этом размер этого битсета ты можешь делать любой. Можешь сказать, хоть два элемента у меня будет в этом массиве, 0 и 1, — понятно, что коллизий тогда будет очень много. Дальше ты хочешь уменьшать вероятность ошибки, хотелось бы контролировать, чтобы эта штука не всегда тебя обманывала и хотя бы иногда говорила правду. И у тебя есть два способа это контролировать. Первый — ты можешь увеличивать размер массива, который хранит эти бакеты. Чем больше массив, тем меньше вероятность, что два разных ключа по хэшу попадут в одну точку и будет коллизия. Но чем больше массив, тем больше памяти под это надо — трейд-офф памяти и точности.

[1:35:32] Дмитрий: А второй механизм опять-таки похож на хэш-мапу. В алгоритмах хэш-таблицы есть два способа разрешения коллизий. Первый — когда ты делаешь цепочку: к bucket’у прицеплена какая-то цепочка, куда ты добавляешь элементы (как в Java реализовано; она потом превращается в деревья, если элементов много, но это детали). А альтернативный механизм — открытая адресация: если ты посчитал хэш от своего элемента, пошёл в массив, а там уже занято, — альтернативная идея пересчитать хэш на какой-нибудь другой. И там толпа алгоритмов. Можешь сказать «значение хэша плюс единичка» — текущая ячейка занята, посмотри в соседнюю; или «значение хэша возвести в квадрат», перепрыгнуть через какой-то ключ. И в Bloom filter похожая идея. Вместо одного хэша ты считаешь несколько хэшей. Пришёл ключ, и для него ты посчитал хэш 1, хэш 2, хэш 3, — n хэшей. Каждый хэш указал тебе позицию в этом массиве с бакетами, и там либо нолик, либо единичка. Если все из этих хэшей сказали единичку, значит, с большой вероятностью это не коллизия, а реально такой ключ был записан. Если какой-то из них сказал нолик, значит, для твоего ключа пометку точно раньше не делали, и твоего ключа нет. То есть, увеличивая количество этих функций, ты тоже можешь уменьшать количество ошибок, но увеличиваешь объём вычислений.

[1:37:28] Дмитрий: То есть у тебя получается структура данных, в которую ты можешь вкрутить два параметра: размер массива с бакетами (в нашем случае вырожденного, с ноликами и единичками) и количество хэшей, которые адресуют этот массив. И там дальше есть хитрая математика, которая позволяет вычислить вероятность ошибки в зависимости от того, какой у тебя объём данных, сколько возможных ключей. Ты говоришь: хочу, чтобы с вероятностью 99% эта структура отвечала мне правильно, — и можно посчитать, какой должен быть размер массива и сколько взять функций. Эта логика в Cassandra прямо реализована. Когда ты настраиваешь таблицу — мы говорили уже про speculative retry, параметр на таблице, — там также есть параметр вероятности bloom_filter_fp_chance, false positive probability. Там ты задаёшь какой-то процентик, и Cassandra автоматом считает, сколько нужно данных. То есть чем точнее механизм ты делаешь, тем больше платишь памятью.

[1:38:55] Александр: Круто, что это можно настроить. И, конечно, я не представляю, сколько всяких комбинаций разных настроек существует, и каждая комбинация уникальна ещё тем, что она на конкретных данных под конкретным workload используется. Тут золотого рецепта не существует, надо каждый раз смотреть свою нагрузку и экспериментально подкручивать под то, чтобы вас устраивало.

[1:39:23] Дмитрий: Могу сказать, что на практике реально ты не очень часто крутишь этот параметр. Я не припомню за свой опыт, чтобы прямо приходилось его подкручивать. Дефолт здесь достаточно хорош. И надо понимать: в чём ещё прелесть этой структуры данных, Bloom filter? У тебя реально это битсет, набор битов памяти, и он никак не зависит от размера твоих ключей. Ключ может быть гигантским, может быть очень маленьким — ты всегда считаешь от него хэш и в этом битсете ставишь нолики-единички. Размер структуры никак не зависит от размера ключей. А уж от размера значений подавно — они даже не участвуют в вычислении. То есть единственное, что влияет на размер, кроме параметров настройки, — это сколько различных ключей у тебя в принципе в данных есть. То есть какая кардинальность твоих данных.

[1:40:23] Дмитрий: И такая структура данных в Cassandra строится для каждого SSTable-файлика. Когда мы файлик из памяти сбрасываем на диск, в этот момент мы Bloom filter составляем: итерируемся по структуре в памяти, пишем её на диск, и параллельно в памяти проставляем эти биты. Получившийся результат рядышком тоже кладём на диск, чтобы после рестарта его можно было загрузить. Если посмотрите, как устроен Cassandra SSTable-файл на диске, это набор файликов: один из них основной файл с данными, несколько других файлов, о которых мы ещё поговорим, и среди них есть файл с Bloom filter, где этот битсет прямо сериализован на диск. Но ещё раз: Bloom filter — это, я представляю, такой обрубок хэш-мапы, у которого просто бакеты с ноликами и единичками вместо разрешения коллизий. Мне такая аналогия лично мне помогает. А заполнение идёт точно так же: мы считаем хэши, только сейчас не читаем из ячеек, а туда пишем единички. Получили очередной ключ, посчитали от него n хэшей, эти n хэшей указали на n позиций в массиве, и во все эти позиции мы поставили единичный бит.

[1:41:50] Александр: Окей, допустим, получили мы ответ от Bloom filter, что данные по этому ключу, скорее всего, есть. И что мы делаем дальше?

[1:42:01] Дмитрий: Да. То есть мы надеемся, что большинство таблиц всё-таки отбросило, если нам повезло, и их вообще читать не надо. Но есть набор файликов, где данные, наверное, есть, и нам придётся их читать. И вот у нас на диске: мы знаем, что в этом файлике по partition-ключу данные отсортированы, а внутри каждого partition-ключа секция по clustering-ключу тоже отсортирована. То есть есть такой гигантский отсортированный массив ключей, с которыми как-то привязаны данные. И нам надо в этом гигантском отсортированном массиве найти позицию на диске, откуда читать данные. А мы знаем только partition-ключ и clustering-ключ. Что делать?

[1:42:49] Дмитрий: Тут первая мысль: хочешь что-то эффективно искать в структуре данных — сделай индекс. Как в базе данных: хочешь искать по какому-то набору данных — надо индекс. И здесь такой индекс есть. На верхнем уровне, опять-таки, сначала по partition-ключу: для partition-ключей мы строим отдельный индекс. Когда мы пишем memtable на диск, мы составляем дополнительный индекс, называется primary index. Он состоит всего из двух частей: первая часть — partition-ключ, вторая — offset. То есть у тебя есть файл с данными, а рядом второй файл, где написано: «ключ такой-то — ищи по такой-то позиции в этом файле». И это для каждого partition-ключа так написано. Упрощённо говоря, структура данных такая. Ты знаешь partition-ключ, можешь в этой структуре, в первичном индексе, найти свой partition-ключ, найти offset, для которого в файле с данными лежат твои данные.

[1:43:57] Дмитрий: Но в чём тут проблема? Если у тебя данных очень много, то этот индекс в память целиком ты уже не загрузишь — он будет слишком большой. Поэтому его придётся держать на диске, а дальше как-то делать в нём поиск, и хотелось бы делать это эффективно. Ты можешь загрузить в память не весь индекс целиком, а его прореженную версию: бежишь по этому индексу и каждую n-ную строчку кладёшь в память. Эта штука называется index summary. То есть у тебя есть такой а-ля словарь: все слова отсортированы по порядку, и есть большой перечень всех слов без описания — слово, номер страницы; слово, номер страницы. Но этот перечень тоже очень большой, слов в словаре дофига. И рядышком у тебя есть такое супероглавление, где написано: буква «А» — страница такая-то, буква «Б» — страница такая-то. А дальше ты уже можешь найти детальный перечень, зайти туда и искать своё слово. То есть нам надо быстро найти позицию в файле, которую имеет смысл читать, где, скорее всего, располагается интересующая нас информация. Для этого мы сделали индекс, но поскольку целиком в память он не влезает, загрузили прореженную версию — index summary. И в этой прореженной версии, где только каждая n-ная строчка индекса, мы можем сделать поиск.

[1:45:32] Дмитрий: Поскольку изначальный файл отсортирован, индексы тоже отсортированы по partition-ключу, и у нас есть значение ключа. Задача поиска в сортированном массиве того, что наиболее близко к нашему ключу, — это снова бинарный поиск. В этом index summary мы делаем бинарный поиск по нашему partition-ключу и находим диапазон: наш ключ лежит между вот этой строчкой summary и этой. А эти две строчки summary соответствуют какому-то кусочку индекса на диске. Мы идём в этот кусочек индекса на диске и просто линейно по нему пробегаем. Он не очень большой, скорее всего, мы прочитаем его за одну операцию ввода-вывода с диска и вытащим в память.

[1:46:23] Дмитрий: Насколько он большой — как настроишь. Там есть прямо параметр, который позволяет это настроить. Но опять-таки это трейд-офф между памятью и диском: ты можешь делать index summary более гранулярным, но он будет кушать больше места в памяти, зато ты будешь очень быстро искать на диске нужный кусочек. Либо наоборот, сделать его очень прореженным, но тогда будешь больше читать с диска, чтобы найти конкретную позицию в индексе. В этом индексе ты, пробегая по порядку, находишь нужную строчку, и там есть offset — попадаешь в основной файл.

[1:46:54] Александр: Да. Тут важно, что offset нам нужен, потому что это по сути указатель на место в файле, где лежат интересующие нас данные. А мы-то изначально пришли с ключом или с хэшом. И это не одно и то же, поэтому нам нужно одно перевести в другое. И вот этот перевод ключа в offset уже требует некоторого индекса, который большой, который требует index summary, хранящийся в памяти. Это уже задача нетривиальная: мы не можем себе позволить всё это загрузить в оперативную память. Мы проделали эту работу, и теперь у нас есть offset — конкретный указатель на конкретную строчку в конкретном файле.

[1:47:39] Дмитрий: Байтовую позицию, в байтах. 17 байт от начала. Мы прямо, когда работаем с этими файлами данных, можем сделать seek туда и начать что-то читать.

[1:47:52] Дмитрий: Тут опять-таки можно делать seek и оптимизировать. И одна из оптимизаций такая. Получается, что чтобы прочитать наши данные, мы должны сделать два чтения: чтение из индекса (из summary в памяти мы попадаем в основной индекс на диске), а потом ещё на диске лежат сами данные. То есть как минимум два чтения. Мы можем попытаться первое чтение выкинуть. За счёт чего? Мы можем сказать: наверное, ключи читают неравномерно, есть какие-то ключи, которые читают часто, а какие-то — нет. Давайте сделаем кэш в памяти, где ключом будет тот ключ, который мы ищем, а значением — сразу этот offset. Там будут не все данные, но это кэш. Называется key cache в Cassandra. Это ограниченный по объёму кусочек памяти, где мы держим по сути дела хэш-мапу «ключи → offset’ы», которые вытесняются. Если мы забили память, надо что-то освободить. И опять-таки, тут мы не изобретаем ничего нового, никаких своих алгоритмов кэширования не делаем. Это Caffeine снова. Библиотека Caffeine. Она у нас уже встретилась три раза, на самом деле.

[1:49:16] Дмитрий: У нас есть Caffeine-кэш для prepared statement, где мы по id prepared statement храним текст и разобранный синтаксический запрос — наш CQL-запрос. У нас есть Caffeine-кэш для авторизации и аутентификации, где хранится информация, что какие-то пользователи, роли имеют права на какие-то ресурсы — таблицы, keyspace’ы и так далее. И третье место, где у нас есть Caffeine, — это key cache, где мы для SSTable-файликов храним маппинг partition-ключа на позицию в data-файле, где этот partition-ключ держит данные. Все три кэша вытесняемые, у которых задаётся метрика — размер. Мы прикидываем, сколько в памяти, в Java heap, мы готовы под эти структуры данных выделить, и просто возвращаем функцию веса: вот хранение одной строчки нам нужно, 15 байт. И суммарно под эту структуру данных мы готовы выделить 10% heap, 5% heap, 60 мегабайт — я не помню конкретное число. И если мы превысили этот размер, то менее популярные строчки вытесняются за счёт хитрого алгоритма, который Caffeine реализует, достаточно продвинутого — TinyLFU и подобных.

[1:50:56] Дмитрий: Это как мы конкретно файлики нашли. Но если ты следил внимательно, я везде говорил слово «partition-ключ, partition-ключ». То есть мы нашли в файлике не нашу конкретную строчку, а нашли партицию — какой-то кусок данных, отвечающий за эту партицию. А у нас ещё clustering-ключ есть.

[1:51:20] Дмитрий: А дальше тебе нужно для clustering-ключа либо пробежаться по диапазону ключей, чтобы вернуть данные, — может, ты всю партицию целиком попросил, тогда читаем с самого начала. Но возможно, тебе нужно вернуть данные по конкретному clustering-ключу. И тут немножко опять похоже на Java-хэш-мапу — логика динамическая. Если данных для этой партиции из SSTable мало, то там clustering-ключи просто записаны по порядку, и ты можешь быстренько их перебрать: прочитаешь весь блок данных, переберёшь, ненужное откинешь, найдёшь нужное и вернёшь. А вот если по одной партиции записали очень много данных, то ты можешь очень долго читать — переберёшь с диска несколько мегабайтов данных, это невыгодно, хотелось бы побыстрее. И тут снова сделана индексация, примерно по такой же идее, как мы только что проходили. Ты можешь взять уже сформированные clustering-ключи и сделать для них такой summary-индекс: проредить их, каждую n-ную строчку положить в некое summary, которое будет лежать в этом data-файле в самом начале.

[1:52:35] Александр: То есть та же самая идея, применимая просто в другом… Да, такая рекурсивная идея. Интересно. И там, наверное, даже код один и тот же?

[1:52:46] Дмитрий: Не совсем, но похож. Ты прочитаешь этот первый блок, и там будет прореженное начало диапазона clustering-ключей. Можешь понять, в какой из них тебе интересно смотреть. Это, по-моему, называется shadow index — некая структура, которую мы тоже on-demand подгружаем, вместо того чтобы грузить сразу на лету. Это классическая структура, которая в основном используется в Cassandra. В пятой версии появились ещё новые вариации, про которые я сейчас рассказывать не буду, просто упомяну. Так же, как для memtable придумали, что можно использовать префиксное дерево, trie, — то же самое можно попытаться реализовать и на диске. То есть можно сделать префиксное дерево в персистентном виде, и это была отдельная фича в Cassandra, которую реализовали; она доступна, в версии 5.0 уже есть. Когда механизм хранения данных не тот, что я описывал, а использует префиксную структуру данных в персистентном виде, там уже не нужны дополнительные индексные файлы и key cache — там мы путешествуем по этим префиксам, прямо используя блоки на диске. Чем-то это напоминает B-дерево, чем-то напоминает то, как это реализовано напрямую как префиксное дерево. Можно отдельно про это долго говорить, просто упомяну, что такая штука есть. Реализовал её тот же человек, что реализовал trie, и, на самом деле, тот же самый человек, что реализовал Unified Compaction Strategy, который я упоминал в прошлый раз.

[1:54:25] Дмитрий: Казалось бы, вот уже вся история: мы нашли, откуда данные из файлика читать, читаем их, передаём на координатор, он там всё мёржит или хэши считает. Вот бы казалось, всё. Но есть ещё один важный аспект, который потребует дополнительных метаданных, дополнительной памяти, чтобы делать это эффективно.

[1:54:54] Дмитрий: Мы, когда данные на диске храним, хотим хранить их эффективно, а именно — кушать меньше места на диске. Один из способов этого достичь — использовать какой-то механизм сжатия, то есть я хочу применить компрессию для моих данных. Есть разные алгоритмы, а Cassandra использует блоковые. Мы берём большой data-файлик, где есть строчки, и режем его на фиксированные кубики, фиксированные блоки по 4 килобайта: от 0 до 4 килобайт, от 4 до 8 килобайт, вот такие фиксированные блоки. Этот размер блока можно задавать на уровне конкретной таблицы. Если вы знаете, что будете читать много данных всё время, имеет смысл делать размер блока побольше: чем больше блок, тем эффективнее его сжимать — больше повторяющихся записей, словарик составится эффективнее, в алгоритме сжатия всё будет веселее работать. Но если вам из этих записей нужно одна-две, то вы будете читать с диска большой блок. Чтобы разжать блок, надо прочитать его целиком. Поэтому если вам нужна только одна строчка, невыгодно делать большие блоки: вы будете читать с диска много, разжимать это, а потом почти всё выбрасывать и оставлять парочку строк. Тут обратный трейд-офф — хотелось бы сделать блок поменьше. То есть трейд-офф между эффективностью ввода-вывода и эффективностью сжатия: чем больше блок, тем лучше сжимается, но тем больше читается.

[1:56:25] Дмитрий: Вот мы нарезали данные на эти блоки, каждый из них сжали. Есть несколько алгоритмов, которые Cassandra поддерживает. Основные — либо LZ4 по умолчанию, хороший компромисс между скоростью и силой сжатия, либо ZSTD. Он подороже в плане CPU — дольше сжимает и разжимает, но делает это эффективнее. И сейчас, кстати, буквально неделю или две назад, в транке появился продвинутый вариант этого механизма сжатия — ZSTD со словарями. Ты можешь заранее натренировать базу на данных, чтобы она составила словарь и сжимала ещё эффективнее, используя этот ZSTD-алгоритм.

[1:57:15] Александр: То есть статистику подключить.

[1:57:17] Дмитрий: Да, появляется статистика, которая позволяет сжать эти данные ещё лучше. Понятно, что ценой дополнительных манипуляций: тебе нужно сначала статистику подготовить. По умолчанию у вас LZ4, даже если вы ничего не настраиваете. Теоретически вы можете вообще выключить эту компрессию и держать данные как есть, но на практике такое никто не делает.

[1:57:43] Дмитрий: Сжимается всё по-разному. У меня были случаи, когда, если исходный объём данных брать за 100%, после сжатия мы занимали 15–20% от изначального объёма. То есть сжимаем данные в 5–7 раз.

[1:58:03] Александр: Ничего себе! Для OLTP-базы данных, где у нас не поколоночное хранение…

[1:58:09] Дмитрий: Да, это не колоночное, где можно всякие дельта-кодирования хорошие сделать по повторяющимся данным. Это отличный результат. Понятно, что всё зависит от того, какие у вас данные: у меня были довольно повторяющиеся данные, которые хорошо сжимались. Если у вас какие-то совершенно случайные значения, понятно, что они будут сжиматься хуже. И опять-таки размер блока: чем больше объём блока, тем больше повторений возможно, тем лучше сожмётся. Если будете сжимать маленькие блоки, цифра может быть хуже.

[1:58:43] Александр: А это происходит непосредственно прямо перед записью на диск и после чтения, или где-то промежуточно сжатые данные у нас ещё гуляют?

[1:58:51] Дмитрий: Так, тут несколько вопросов даже. Давай по порядку. Мы храним данные на диске в сжатом виде — этот data-файлик у нас сжат. Соответственно, data-файлик пишется в момент флеша memtable на диск: когда мы перебираем данные и сбрасываем их из памяти на диск, в этот момент мы их и жмём. Прямо на нижнем уровне, когда работаем с байтиками: накапливаем блок, фиксируем размер, жмём, флешим на диск — ничего сложного. Аналогично во время компакшена: когда делаем компакшен, читаем старые файлики, записываем новые, — понятно, что новые тоже жмём. А когда данные читаем с диска, из SSTable, — вот в этот момент мы их разжимаем. У нас есть позиция, которую мы через key cache или напрямую по индексам посчитали. Позиция, например… для простоты скажем, что блок размера 100 (понятно, что там кратное число килобайт удобно делать, но допустим, для нас проще десятичная система). Блок мы хотим прочитать… ой, не блок, а строчку 215, позицию байта. А блоки у нас размера 100, соответственно, мы пойдём во второй блок, прочитаем его с диска. И в этот момент, когда читаем блок с диска, мы должны его целиком разжать, потому что по частям разжать я не умею — оно будет весь кусочек разжимать целиком.

[2:00:25] Дмитрий: Вот. Но кроме сжатия на уровне дисков, аналогично можно пытаться использовать сжатие, чтобы сэкономить объём пересылаемых байтов по сети. Ты можешь настроить сжатие при общении между клиентом и сервером: сказать, что клиент поддерживает сжатие, сервер поддерживает сжатие, они обменяются вначале хендшейком, договорятся, что можно данные сжимать. И ты там можешь обменять дополнительное CPU на стороне клиента и сервера на более компактное представление по сети. То же самое можно включить между общением реплик между собой. Причём ты можешь даже более хитро сделать: сказать, что у меня реплики в локальном дата-центре, между ними хорошая сеть, и там сжатие не нужно — наоборот, пускай побыстрее работает, потому что сжатие-разжатие требует времени, CPU — это ресурсы. А вот когда реплики общаются между дата-центрами, у меня межцентровый трафик, канал более ограниченный и может быть вообще платный — какой-нибудь Amazon, который берёт дополнительные деньги за cross-datacenter-трафик. В этом случае лучше сжимать. Я могу на уровне настроек Cassandra задать сжатие между DC, но не сжимать в локальном DC.

[2:01:47] Александр: Круто, очень продвинутая фича.

[2:01:53] Александр: Так, это что касается сжатия. Вопрос про сжатие здесь. Когда мы сжали один раз эти блоки… ты говорил, можно в разных ситуациях использовать сжатие, каждый раз мы будем сжимать-разжимать. То есть эти данные — это не один и тот же физически сжатый набор байтов?

[2:02:11] Дмитрий: К сожалению, это разные наборы байтов. Ты не можешь их переиспользовать, как с чего-то. То, что я упоминал для commit-лога и для механизма сериализации данных по сети, — там один формат. А здесь, к сожалению, это немножко разные форматы, раз. А во-вторых, ты из каждой SSTable читаешь данные, у тебя несколько наборов данных, которые ещё надо смёржить. Ты из одного файлика SSTable прочитал, из второго прочитал, — там могут быть половинки, пересекающиеся данные, — тебе надо смёржить. А у тебя операция сжатия не гомоморфная, чтобы ты мог сделать преобразование над сжатыми данными без разжатия.

[2:02:57] Александр: Да, тебе надо разжать, сделать… смёржить и сжать.

[2:03:01] Дмитрий: Смёржить и сжать, да. На самом деле, ещё одно место, где есть сжатие, — это commit-лог. Ты можешь сжать ещё и commit-лог. Иногда это бывает полезно: если у тебя CPU много, а диск почему-то слабоватенький вдруг, то можно включить сжатие для commit-лога — по-моему, сейчас оно выключено, — и тогда там будет просто объём меньше.

[2:03:23] Дмитрий: Но, как я сказал, принеся оптимизацию, мы принесли себе новую арифметическую проблему. Вспоминаем, что у нас в индексе была позиция в файлике: говорили, что для того, чтобы прочитать данные для такого-то partition-ключа, сходи по такому-то offset’у. Но это был offset в несжатой структуре данных, не в сжатом файлике, а в разжатом. А ты представь себе: у тебя есть файлик изначально несжатый, ты его режешь на кусочки фиксированного размера, блоки. Каждый блок жмётся по-разному, то есть каждый блок усохнет уникальным образом. И этот равномерно нарезанный файлик из равных сегментов схлопнется в файлик с непредсказуемыми offset’ами.

[2:04:18] Александр: Беда.

[2:04:19] Дмитрий: Да. И ты не сможешь показать… вот у меня в индексе позиция в несжатом файлике 75 — а какой это блок в сжатом? Непонятно. И без дополнительных метаданных это эффективно не узнать. Поэтому тебе нужно хранить дополнительный маппинг — маппинг между разжатым блоком и сжатым блоком. Чтобы попасть в такой-то разжатый блок, в какой сжатый блок мне нужно зайти? И это отдельный индекс. По сути дела, два числа, которые надо парами записать.

[2:05:01] Дмитрий: Этот индекс потом записывается в файл, он называется compression metadata. Вспоминаем опять набор файликов, которые соответствуют SSTable: вот там есть ещё файлик, который называется compression metadata. Теперь вы знаете, что это за файлик — это файлик с этими offset’ами. Индексы разжатых блоков — это просто порядковые числа, 1, 2, 3, 4, 5, 6 и так далее, их на диск можно не писать. По факту у нас остаётся только вторая половинка — это просто offset’ы в сжатом файлике. То есть набор чисел, которые мы сохраняем на диск, но мы не хотим делать ещё одну операцию ввода-вывода. Этот индекс достаточно маленький, мы его тоже можем положить в память. И это просто массив, где разжатый индекс — это индекс массива, а значение массива — куда в сжатом файлике нужно с каким offset’ом сходить, чтобы найти блок.

[2:05:53] Дмитрий: Вот этот, если говорить про реализацию, массив, — и мы вспоминаем, что у нас ещё был битсет для Bloom filter, — и то, и другое достаточно большие сущности. И тут в какой-то момент Cassandra натолкнулась на проблему, что garbage collection не очень хорошо переваривает такие большие объекты.

[2:06:11] Александр: Да, вспоминаем, что мы на Java пишем.

[2:06:15] Дмитрий: Да, это называется humongous-объекты, для них специальные алгоритмы аллокации. И их как бы копировать — вспоминаем, что у нас есть копирующий сборщик мусора, который туда-сюда объекты копирует. Когда у тебя объект весит десятки мегабайт, а то и сотни, копировать его весьма дорогое удовольствие — двигать такую махину. Итого в Cassandre это решили за несколько итераций и в конце концов остановились на варианте, что это всё выпихнули в off-heap. То есть эти две сущности — Bloom filter и compression metadata index — хранятся в off-heap-памяти.

[2:06:57] Дмитрий: И тут, можно сказать, сейчас есть технический долг в Cassandra, потому что этот off-heap реализован очень быстрым, хорошим способом, но который, к сожалению, не считается public API в Java, а именно — Unsafe. То есть прямо в Cassandra есть куски кода, написанные на прямом Unsafe. Оно очень быстрое, да, небезопасное, но если ты понимаешь, что делаешь, то достаточно безопасно. К сожалению, Unsafe в самых последних версиях постепенно всё худеет и худеет, и операции работы с памятью уже депрекейтят и скоро совсем выпилят. Поэтому одна из задач в Cassandra, чтобы поддержать наипоследнейшие версии Java, — начинать заменять этот механизм на какой-то другой. Наиболее очевидный вариант — Foreign Memory API в Java, который позволяет нам работать с памятью. Но, к сожалению, он немножко подороже. Я слушал последние доклады от тех, кто его делал, от товарища из Oracle: по производительности они по-прежнему не до конца сравнялись. Если мы говорим про случаи с векторизацией, когда мы хотим много ячеек прочитать за раз, накладные расходы размазываются между многими ячейками и становятся незаметными. А вот когда у нас единичный доступ в эти ячейки — это такие hash-map-like-доступы, где надо прочитать буквально пару байтов, — эти накладные расходы всё-таки видны. Насколько я помню, процентов на 10 оно станет медленнее. Но, видимо, других вариантов нет, придётся переходить. Пока это не так, пока в Java это Unsafe, и можно будет местами увидеть ворнинги, где JVM говорит «ай-яй-яй, меня тут использовать не надо».

[2:08:55] Александр: Да, частенько, знаешь, базы данных или всякие большие распределённые системы, которые хранят кучу данных в памяти и кучу метаданных в памяти, и метаданных над метаданными, — на Java они сталкиваются с этой проблемой, что, блин, эти объекты большие. Это большие структуры данных, и garbage-коллекторы просто замедляют работу: копировать становится дорого, да и обходить эти графы тоже. Короче, это становится нерабочая схема. То есть Java, вообще говоря, не самый лучший язык, чтобы писать на нём базы данных, если уж совсем откровенно. Потому что вот такие структуры…

[2:09:38] Александр: То есть garbage-коллектор, грубо говоря, вот он мешает. Когда ты хочешь голыми руками байтики сериализовать, десериализовать и аккуратненько с ними работать, потому что хочешь достигнуть максимальной производительности, — такие штуки, как garbage-коллектор, тебе мешают.

[2:09:51] Дмитрий: Да, они ускоряют разработку, потому что база данных, как и любой другой программный продукт, — это вещь, которую разрабатывают программисты. И скорость разработки, количество библиотек — это всё выше на Java, чем на тех же C++, если бы мы писали. Это да. Но в конечном итоге, когда программа уже написана, у неё возникают вот такие проблемы. И чаще всего в Java просто выгружают это в off-heap, и есть другие базы данных, которые точно так же делают.

[2:10:22] Дмитрий: Соответственно, как я сказал, Bloom filter — в off-heap, и offset’ы для compression metadata — в off-heap, и memtable, на самом деле основная жирная структура данных у нас в памяти — это memtable, где данные копятся. Сейчас она в состоянии частично off-heap. Её можно настроить, чтобы данные частично хранились off-heap. Там есть разные режимы работы, и один из них называется, по-моему, native objects. Есть параметр memtable type, или как-то так, конфигурационный параметр, который позволяет вам сказать, где хранить эти данные. Там определённый trade-off: иногда эти off-heap-механизмы подороже, потому что тебе приходится во время процессинга копировать данные туда-сюда. Когда ты процессишь, удобнее, чтобы они были в heap, а когда хранишь — чтобы были off-heap. Приходится гонять их туда-сюда, из-за этого есть накладные расходы. Поэтому это по-прежнему конфигурационный параметр, и разные люди ставят его в разные значения в зависимости от того, что лучше для них работает. Но есть возможность частично memtable вытеснить off-heap.

[2:11:40] Дмитрий: В этом плане, кстати, trie-дерево как раз пошло дальше: кусок, который непосредственно связан с индексом, — это префиксное дерево для partition-ключа, — оно прямо целиком реализовано off-heap. Причём там реализован вариант не на Unsafe, а альтернативный подход через direct byte buffer. Мы выделяем достаточно большие direct byte buffer’ы и потом в них индексируем память — как в таком off-heap-массиве. Есть, по-моему, идеи ещё больше вытащить, то есть, возможно, eventually прогресс в Cassandra дойдёт до того, что эта структура будет полностью off-heap, но там ещё много работы. В целом тренд такой есть.

[2:12:30] Дмитрий: С одной стороны — вот это, а с другой стороны, JVM тоже в этом плане эволюционирует, становится лучше. Во-первых, новые garbage-коллекторы, которые появились, особенно их версии с поколениями — ZGC, а теперь и G1 с поколениями, — существенно снижают stop-the-world-паузы. У меня есть опыт использования Shenandoah в продакшене: там действительно, если базу не перегружать, единичные миллисекунды паузы уже вполне реальны. Так что эта головная боль постепенно уходит. С другой стороны, в последних версиях Java появилась такая штука, как компактные заголовки, compact headers. Каждый объект в памяти в Cassandra имеет накладные расходы — некий заголовок, который нужен JVM для работы с объектом, и он по умолчанию был достаточно жирненький, по-моему, 12 байт, а теперь стал 8 байт. Есть надежда, что и до 4 его уменьшат. Чем меньше этот overhead, тем больше объектов вы можете хранить в памяти, тем лучше с ними работают процессоры, потому что в кэш попадает больше объектов, и становится всё лучше.

[2:13:49] Александр: Обратный плюс Java в том, что если просто ничего не делать и ждать, становится лучше.

[2:13:56] Дмитрий: Да, это факт, прикольно.

[2:13:58] Александр: Обновляешь просто Java — и становится просто лучше.

[2:14:01] Дмитрий: И мы, получается, прочитали данные, и дальше мы их мёржим. И тут снова в каком-то смысле получается, как в compaction, сортировка слиянием. Когда мы читаем, у нас есть набор итераторов — на диске по данным и в памяти, — и везде данные отсортированы. И мы можем просто эти потоки данных слить в такой итератор, который их мёржит и на лету склеивает, получая финальные строчки, которые мы либо можем захэшировать, либо отправить по сети как есть. И вот так работает чтение.

[2:14:44] Дмитрий: Основной тут, наверное, подводный камень, который стоит упомянуть уже под конец, — это то, что в механизме чтения наибольшую боль вам будут доставлять данные, которые вы удаляете. Если у вас read-only-данные, то всё работает прекрасно. Но если вы делаете данные, которые очень часто удаляются, — например, вы хотите сделать таблицу-очередь, очередь ордеров, куда один процесс постоянно накидывает строчки, а другой их обрабатывает и удаляет, — у вас получается, что большую часть времени в таблице ничего нет: там вставили, почти сразу удалили. Но вспоминаем, что всё, что мы удалили в Cassandra, удаляется не безвозвратно, а превращается в tombstone’ы. И когда вы начинаете читать такую якобы пустую таблицу, вы в реальности начинаете читать кучу SSTable-файликов, где лежат эти tombstone’ы: базе придётся их поднять с диска и дальше отбросить уже на этапе ответа клиенту. Там, на самом деле, два этапа отброса: совсем старые данные она отбросит на этапе чтения на репликах, а финальные, даже свежие tombstone’ы отбросит, когда будет отвечать клиенту. Но если там миллионы строк, которые были удалены, то база данных будет их поднимать и тратить на это время. Поэтому в Cassandra это считается классическим антипаттерном — queue-like-таблица.

[2:16:11] Дмитрий: Есть способы это обойти, если совсем уж никак. У меня такие примеры были. Основной механизм — надо нарезать эту таблицу как можно мельче на секции. Вы в качестве partition-ключа используете, например, дату: типа каждый день — это такая мини-очередь. Соответственно, вы не читаете всё, что удалили в прошлые дни, вы читаете только то, что попадало в таблицу в текущий день. Ну или вы можете гранулировать эту очередь как угодно — по минутам или ещё как-то. И просто на уровне бизнес-логики явно эти мёртвые диапазоны не читать, тем самым избежав оверхедов на чтение tombstone’ов.

[2:16:58] Александр: Да, хороший рецепт.

[2:17:00] Дмитрий: Ну это да, одна из основных болей Cassandra — это tombstone’ы. Но, как мы обсудили в прошлый раз, от них, к сожалению, крайне тяжело избавиться: слишком много задач они решают. И в плане локального хранения для LSM-деревьев они нужны, и в плане координации между серверами — сделать так, чтобы не появились зомби после удаления данных.

[2:17:26] Александр: Ну, один из вариантов — не иметь TTL и не удалять данные.

[2:17:30] Дмитрий: Но с этим сложно, со всякими GDPR’ами и так далее.

[2:17:37] Александр: Ну, в плане GDPR’ов это вещь о двух концах. С одной стороны, есть требование удалять данные пользователя, а с другой — регуляторы в итоге требуют хранить данные в стиле: пожалуйста, данные по всем финансовым транзакциям вы должны, согласно государственным регуляциям, хранить два года. В зависимости от страны — или три года. Так что тут ещё такая хитрая вещь.

[2:18:06] Александр: Да, ну мы очень хорошо поговорили про чтение. Подключая прошлые выпуски про запись, про взаимодействие с сервером и так далее, мне кажется, мы покрыли процентов 80 того, как работает Cassandra. И покрыли довольно глубоко. То есть мы точно не поговорили про транзакции, про lightweight-транзакции — тут мы их прямо обошли. И, наверное, ещё кратко перечислим вещи, которые не затронули.

[2:18:36] Дмитрий: Да, давай. Как ты сказал, мы обсудили только базовые механизмы чтения и записи. Многое на них построено, но это не единственный способ. Иногда нам хочется, даже путём просадки в производительности, получить какие-то более надёжные в плане consistency вещи. Есть две штуки. Первая — lightweight-транзакции, которые существуют уже с третьей или четвёртой версии Cassandra, и их можно использовать в продакшене. Правда, про них есть шутка, что lightweight-транзакции — это не lightweight и не транзакции. Во-первых, они достаточно дорогие: это по сути дела полноценно реализованный протокол Paxos, для того чтобы вы могли делать операции compare-and-set — в стиле «добавь строчку, если такой строчки нет» или «обнови строку, если предыдущее значение было такое». Такие операции они поддерживают. Но, во-первых, всё ограничивается только одной партицией. Во-вторых, чтобы это реализовать, нужно несколько раундов общения между координаторами и нодами, и вы за это заплатите производительностью и меньшей отказоустойчивостью. Поэтому лезьте туда, только если совсем никак не реализовать на базовом уровне.

[2:19:54] Дмитрий: И сейчас, на самом деле, на смену или как дополнение этим lightweight-транзакциям в транке уже закоммичен, но пока в релизную версию ещё не попал, новый протокол распределённых транзакций в Cassandra. Называется Accord. Там прямо честные распределённые транзакции, которые позволяют вам работать с несколькими партициями одновременно и дают уровень изоляции serializable — самый сильный уровень, как будто вы всё последовательно выполняете. При этом там нет выделенного координатора. Можно по ключевым словам «Accord транзакции Cassandra» найти — есть статьи и несколько докладов на эту штуку. Я ссылки, наверное, тоже скину на всякий случай. Все очень ждут. Я надеюсь, что производительность там будет приемлемая. Понятно, что это всё тоже не бесплатно, и базовые операции чтения-записи оно не победит, но если оно работает неплохо, то решит многие проблемы, которые раньше на Cassandra было очень тяжело решить.

[2:21:02] Дмитрий: Мы не поговорили про вторичные индексы. Опять-таки, довольно свежая вещь. В Cassandra было несколько заходов на вторичные индексы. Сейчас очень хорошо полетели так называемые SAI-индексы, Storage Attached Index, когда мы строим индекс для каждой SSTable- и memtable-структуры. В паре к тем структурам данных, которые мы описали, есть ещё рядышком такой вторичный индекс по ним. И мы можем делать эффективный поиск, декомпозировав задачу на эффективный поиск в каждой такой сущности, сделать broadcast-запрос на кластер и собрать результаты. Есть про это доклады, в интернете можно посмотреть.

[2:21:45] Дмитрий: Есть всякие дополнительные вещи, которые обеспечивают в Cassandra консистентность данных. Мы обсудили два из них. Это хинты — когда мы делаем операцию записи, а реплика недоступна, мы откладываем запись на потом. И это read-repair — то, что мы обсуждали сегодня: когда мы читаем данные с реплик, данные разошлись, и мы, обнаружив это, записываем результат обратно на реплики, тем самым восстанавливая консистентность между нодами. Третий столп консистентности Cassandra, который мы не затронули, — это background repair: когда мы в фоне запускаем некую операцию, которая пробегает по всем или по множеству данных в Cassandra-нодах, сравнивает эти реплики между собой и тоже пытается в таком батчевом режиме их выровнять.

[2:22:39] Дмитрий: И сейчас есть ещё всякие проекты, связанные с операционной частью Cassandra. Активно пилят такую штуку, как Sidecar, когда рядом с Cassandra крутится некий мониторинг, некий управляющий процесс, который позволяет всякие вещи с ней автоматизировать. Можно координировать процесс репейра, снимать метрики, профилировать ноду Cassandra, в том числе делать эти балковые загрузки-выгрузки данных, которые я упоминал сегодня раньше, через операции чтения с диска. Очень интересный, свежий проект, который сейчас активно пилят. Скорее всего, я что-нибудь ещё забыл, но вот это, наверное, основные вещи, которые можно упомянуть.

[2:23:30] Дмитрий: В операционном плане сейчас есть ещё активность по тому, чтобы побольше вынести в CQL. Довольно многие операции в Cassandra сделаны через JMX API, и одна из идей в том, чтобы CQL-интерфейс сделать не только интерфейсом для работы с данными, но и — как в других реляционных базах или в Postgres, где вы многие операции можете сделать как вызов функций в консоли, — здесь то же самое. Идёт попытка сделать для Cassandra тоже, чтобы поменьше использовать JMX и сделать Cassandra более кроссплатформенной. JMX — всё-таки такая специфичная для Java технология со своими проблемами и оверхедами.

[2:24:16] Александр: Да. Короче, проект развивается, open-source. Дима, видно, глубоко погружён в процесс. Интересно было очень послушать. Мне все три выпуска про Cassandra безумно зашли. Наверное, это будет одна из моих любимых трилогий.

[2:24:36] Дмитрий: Как «Властелин колец». Такой же длительности, я чувствую, мы получили.

[2:24:42] Александр: Да, там если посчитать…

[2:24:43] Дмитрий: Почти режиссёрскую версию можно получить.

[2:24:46] Дмитрий: Да, там около 10 часов, вообще говоря, если без монтажа, без вырезания каких-то моментов. Но и в финальной версии тоже около 9 часов, наверное, выйдет.

[2:24:59] Александр: 8 плюс точно. Но зато, послушав 8 часов подкаста «Тысяча фичей» про Apache Cassandra, вы, придя потом на интервью, где будет секция про базы данных, просто завалите любого интервьюера.

[2:25:14] Дмитрий: Да, то есть ни один интервьюер этого всего в таком объёме, который мы тут обсудили, не знает, если это не Дима Константинов. Ну, есть ещё куча народу, которые контрибьютят, но вероятность, что вы попадётесь на них, наверное, невелика.

[2:25:32] Александр: Да, но я бы хотел отметить, что всё-таки контрибьютор контрибьютору рознь. Некоторые люди сконцентрированы на своей части и про другое просто не думают или нет времени. А тебя отличает, что у тебя есть широта знаний. И глубина присутствует, безусловно, но и широта — ты про многое знаешь. И в этом плане уникальность контента, который у нас получился, достигнута, как мне кажется.

[2:25:56] Дмитрий: У меня, наверное, основной драйвер, который привёл к такому состоянию — и широко, и достаточно глубоко, — это то, что мои основные контрибьюции в Cassandra касаются тюнинга производительности. А производительность — это такой фактор, который не потюнишь локально. Приходится смотреть все этапы обработки, знать процесс end-to-end и тюнить его end-to-end. Поэтому вот отсюда… Это побочный эффект, на самом деле. Я так много знаю, потому что я все эти места в каком-то виде тюнил.

[2:26:36] Александр: Очень круто. У меня раньше была традиция в подкасте: я в конце какие-то мощные питчи задвигал, какие-то советы, именно такие капитанские, возможно; иногда рекомендовал полезные ресурсы и часто об этом спрашивал гостей. Хотел бы вернуться к этой традиции и спросить у тебя, Дим: что бы ты хотел сказать напоследок слушателям подкаста «Тысяча фичей»?

[2:27:02] Дмитрий: А вот, на самом деле, только что прозвучавшая мысль, которую я озвучил, мне кажется, реально полезная. Если вы хотите в чём-то реально попробовать разобраться, попробуйте сделать это быстрее. В рамках этого процесса вы точно поймёте, как что-то работает, на порядок лучше. Мы знаем, что это вообще интересный процесс — performance tuning и troubleshooting. Troubleshooting я вообще сравниваю с… Для меня troubleshooting подобен детективу: ты не знаешь, кто убийца, свидетели врут, и ты начинаешь искать правду в этом процессе. А performance — это тоже такой процесс, когда ты пробираешься сквозь слои абстракции, понимаешь, как что-то работает. А когда оно ещё начинает работать быстрее — это такое положительное эмоциональное подкрепление. Поэтому вот мой совет: хотите понять что-нибудь — затюньте это.

[2:27:56] Александр: Супер. Спасибо большое. Было мега интересно. И пока.

[2:28:01] Дмитрий: Спасибо.