Конспект 5-ой части Кабанчика - Репликация

Коротко о репликации

Словарик #

Согласованность (данных) = консистентность. Каузально = причинно-следственно.

Введение #

Репликация - хранение одних и тех же данных на нескольких машинах, соединённых между собой. Каждая такая машина называется репликой. В книге репликация разбирается на примере баз данных, но данное понятие можно обобщить.

Не следует путать репликацию с шардированием. Шардирование - распределение разных частей данных между несколькими машинами.

Существует три популярные схемы репликации:

  • single-leader - с единственным лидером
  • multi-leader - с несколькими лидерами
  • leaderless - без лидера

Кроме того репликацию можно классифицировать по принципу передачи данных с одной реплики на другую на:

  • синхронную
  • асинхронную

Синхронная репликация определяет, что операция записи новых данных не является успешной до тех пор, пока все синхронные реплики не запишут эти новые данные.

Асинхронная репликация определяет, что записанные новые данные будут когда-нибудь в будущем переданы на асинхронную реплику. Объём данных или время, на которое происходит задержка записи, называется replication lag, то есть задержка репликации. Асинхронная репликация связана с модным (или уже нет) термином eventual consistency, дословно - “когда-нибудь соответствующий”.

Single-leader #

Разберём сперва простейший вариант - репликацию с единственным лидером.

Как правило в такой схеме запись новых данных ведётся через лидера. Лидер сперва записывает данные в своё хранилище, а после передаёт сведения об изменениях на подписчиков (followers). Процесс передачи происходит обычно через лог репликации (replication log) или поток изменений (change stream). Подписчик берёт данные из лога репликации и применяет изменения на своей локальной копии. Таким образом поддерживается идентичность данных на лидере и подписчиках.

Если подписчиков много, то для каждого может применяться свой тип репликации: например синхронную репликацию поддерживают на одном подписчике, остальные подписчики настраивают в асинхронном режиме.

Для достижения ещё большей пропускной способности на чтение можно перевести всех подписчиков на асинхронную репликацию ценой ненадёжности и несогласованности системы.

Действия при поломке #

При недоступности подписчика #

Если подписчик вышел из строя, то после его восстановления он запрашивает в обычном порядке с лидера все данные, которые он пропустил во время недоступности.

При недоступности лидера #

При выходе лидера из строя необходимо передать статус лидера одному из подписчиков. Этот процесс называется failover. Failover может проходить вручную или автоматически по заранее настроенным правилам.

Как правило автоматический failover состоит из следующих шагов:

  • Определение, что лидер действительно недоступен.
  • Выбор нового лидера некоторым алгоритмом.
  • Конфигурация системы для работы с новым лидером.

При передаче лидерства может пойти не так многое, например:

  • Если настроена асинхронная репликация, то есть вероятность, что никто из подписчиков не обладает полной информацией о данных, потому что у всех накопился лаг репликации, в этом случае потеря данных неизбежна.
  • Возможны сценарии, когда из-за сетевой недоступности две машины выбираются лидерами и обе принимают запросы на запись. Это состояние называется split brain.

Простых решений для этих проблем нет. Каждое решение обладает своими плюсами и минусами. Более подробно возможные решения будут рассмотрены в частях 8 и 9.

Лог репликации #

Существуется несколько вариантов реализации лога репликации.

Лог репликации на базе выражений #

Простейший вариант. Лидер переотправляет на подписчиков выражения, которые выполняет сам. Очень ненадёжный вариант, поскольку выражения могут быть недетерминированными: содержать rand(), now(), вызов триггеров, автоинкрементирующиеся счётчики.

WAL - write-ahead log #

WAL был ранее рассмотрен в части 3.

WAL представляет собой некоторый лог изменений в специальной форме. Рассмотрим на примере PostgreSQL. Базовый вариант - бинарный (physical) WAL, когда на диск записываются бинарные изменения в очень специфичном для определённой версии БД формате. Следующая версия БД не всегда может использовать WAL предыдущей версии БД, что усложняет миграцию на новые версии БД. Переход с версии на версию возможен только с остановкой работы. Тем не менее бинарный WAL обеспечивает относительно высокую скорость работы.

Второй вариант - логический WAL. В этом случае лог дополняется сведениями менее специфичными согласно заранее заданному контракту. Изменения при этом записываются с точностью до строки БД (row state). Это упрощает переход на более новые версии PostgreSQL. Кроме того, это позволяет использовать WAL не только для репликации на БД. Можно настроить потребление данных на произвольных программах, которые могут читать формат логического WAL’а. Такой подход называется change data capture (CDC), дословно захват изменений данных.

Репликация на базе триггеров #

Наиболее гибкий и сложный подход в случае БД. БД как правило предоставляют функционал триггеров - настраиваемых пользовательских функций, вызываемых при внесении определённых изменений в данные.

Подход связан с увеличением нагрузки на систему. Опыт показывает также, что при таком подходе допускается больше ошибок при разработке.

Проблемы лага репликации #

Понятно, что в случае, если асинхронный подписчик отстаёт от лидера, и при этом отвечает на запросы на чтение от клиентов, есть ненулевая вероятность, что клиент получит неактуальную информацию в ответе.

Поэтому при работе с асинхронной репликации необходимо заранее понимать гарантии, которые вы хотите получать от системы.

Чтение собственных записей (read-after write consistency) #

Гарантия, что пользователь после записи некоторой информации увидит внесённые им изменения при чтении.

Пример, что может пойти не так: если пользователь сперва запишет данные на лидера, а после прочитает с отстающего асинхронного подписчика, то он возможно увидит, что его изменения “не применились” на базе данных.

Варианты решения:

  1. Обслуживать чтения пользователя, который только что записал данные, через лидера в течение некоторого времени. Но для этого нужно сохранять некоторую метаинформацию о пользователе.
  2. Сохранять время записи, и сравнивать при чтении, но время на разных машинах не всегда синхронизированно точно.

Всё ещё более усложняется, если консистентные чтения нужно делать с разных устройств пользователя.

Monotonic reads #

Требование, что пользователь при двух последовательных чтениях будет видеть изменения в хронологическом порядке.

Consistent prefix reads #

Требования, что пользователь при чтении видит данные согласно причинно-следственным связям.

Решением вышеуказанных проблем на уровне баз данных выступает механизм транзакций, о которых поговорим в 7 и 9 частях.

Multi-leader replication #

Другие названия этого типа репликации - master-master, active/active.

Репликация с несколькими лидерами предполагает решить проблему недостаточной скорости записи изменений.

Мы получаем:

  1. Быстродействие в смысле записи и возможности обслужить клиентов на запись с меньшей задержкой по сети при условии, что лидеры распределены территориально.
  2. Большая надёжность при выходе из строя одного из лидеров.
  3. Большая надёжность при проблемах с сетевой доступностью одного из лидеров.

Но кроме того, автоматически получаем аналог ранее упомянутого split brain, только который мы привнесли сами. Наличие двух реплик, на которые производится запись порождает конфликты записи. Например, два клиента одновременно записали конфликтующие данные на два лидера. Пока лидеры не синхронизировались, они не знают, что конфликт существует. После попытки синхронизации, необходимо разрешить этот конфликт.

Соответственно теперь система должна быть заранее построена с учётом всех возможностей рассинхронизации данных между лидерами.

В каком-то смысле системы редактирования текста одновременно несколькими пользователями (Google Docs), а также системы синхронизации данных на телефонах: календарные приложения, заметки, - тоже можно считать системами с несколькими лидерами. Например, в случае отсутствия интернета, вы всё ещё можете добавлять записи в ваш календарь на нескольких устройствах, и при подключении к сети приложению нужно будет выбрать, какие сведения с какого устройства считать актуальными.

Работа с конфликтами записи #

Далее варианты решения проблемы с конфликтами записи.

Простейший вариант - синхронное блокирование записи. Этот же вариант полностью сводит на ноль преимущества multi-leader.

Можно избегать конфликтов вынося работу с записями каждого пользователя на назначенного пользователю лидера. Но в случае недоступности назначенного лидера проблема возникает вновь.

Есть вариант выдачи каждой записи уникального идентификатора, использование потом этого идентификатора для решения, какая из записей является более приоритетной. Если в качестве идентификатора используется timestamp и при сравнении побеждает запись более свежая, то такой подход называют английским акронимом LWW (last write win).

Также можно придумать специальные варианты решения конфликтов на уровне пользовательского приложения. Например, при обнаружении конфликта пользователю предлагается решить его вручную в пользовательском интерфейсе.

Разработаны также структуры данных, которые позволяют автоматически разрешать конфликты, предлагаются к рассмотрению CRDT (Conflict-free replicated datatypes), Mergeable persistent data structures, Operational transformation. За ссылками обращайтесь к оригинальному тексту.

Топология репликации лидеров #

Топология репликации описывает структуру каналов репликации между лидерами.

Например, лидеры могут обмениваться данными между друг другом по:

  1. Кольцу;
  2. Топологии звезды, которая может быть обобщена в дерево;
  3. Все-ко-всем

Первые две проще, но связаны с риском выхода из строя узловых лидеров, через которых завязана репликация.

Последняя увеличивает возможные варианты несогласованности данных, когда из-за сетевых задержек зависимые инструкции будут выполняться в разном порядке.

Leaderless replication (репликация с кворумом) #

Если все лидеры, то никто не лидер

Идея в отказе от принципа записи на одну конкретную машину лидер, теперь клиент обращается сразу ко всем репликам, но при этом ожидает успешного ответа от некоторого подмножества реплик.

Пусть имеем n реплик, будем считать, что для записи достаточно w успешных ответов, а для чтения необходимо r успешных чтений, тогда для гарантированного получения актуальных сведений должно выполняться отношение w + r > n. Это правило называется правилом кворума.

Поскольку достаточно только подмножества успешных ответов при записи, есть вероятность, что на часть реплик не запишутся актуальные данные. Тем не менее, мы хотим чтобы когда-то они туда всё равно попали для сохранения надёжности данных. Есть два пути, которые обычно применяются одновременно:

  1. Read repair - то есть обновление записи при чтении. Если после чтения с r реплик клиент обнаруживает, что некоторые записи устарели, он проводит последующее их обновление на устаревших репликах. Но если данные читать редко, то не будет возможности узнать, что они несогласованы, поэтому существует следующий вариант.
  2. Упорядочивание данных (anti-entropy process), буквально процесс направленный на устранение энтропии. Для редко запрашиваемых данных существует фоновый процесс, который периодически сверяет данные на разных репликах и приводит к их общему виду.

Ограничения кворума #

Несмотря на кажущуюся надёжность кворума, даже при соблюдении условия кворума есть много вариантов, что что-то пойдёт не так:

  • Если две записи происходят одновременно, то возникает вопрос, какое значение считать актуальным (конфликт записи);
  • Если запись происходит одновременно с чтением, то может быть неясно, какое из значений - новое или старое было прочтено;
  • Если запись успешно произошла на части реплик, но при этом упала на другой части и при этом не удалось собрать кворум, то неопределена в общем случае судьба нового значения на первой части реплик, нужно дополнительно определять процедуру отката.
  • В общем случае наличие кворума и временные задержки могут приводить к неверным результатам.

Опять получается, что нужно поддерживать баланс между надёжностью и доступностью.

Вариант увеличения доступности - Sloppy Quorum - разрешает на некоторое время при серьёзных сетевых проблемах смягчать условие кворума, и вести запись на часть реплик, которые доступны в данный момент, чтобы после восстановления системы передать данные на отставшие реплики (такая передача называется hinted handoff). Разумеется, такой подход уменьшает надёржность и увеличивает частоту конфликтов.

Обнаружение одновременных записей #

Ранее мы опустили определение того, какое же из значений клиент будет считать актуальным.

Например, если имеем 5 реплик, то есть n = 5, а w = 3, r = 3, то существует возможность записи нового значения на 3 реплики, а потом чтения с трёх других реплик, всего одна из которых будет входить в первые три. В результате клиент получает три ответа: новое значения, старое значение, старое значение.

Очевидно, что в таком случае нельзя использовать большинство в качестве критерия актуальности, поэтому приходится опираться на другие варианты.

Можно записывать время записи значения, и такой способ действительно используется. Однако время течёт для разных реплик по-разному из-за неточностей тактовых генераторов процессоров, и такой подход не является надёжным.

Определение одновременности #

Когда мы говорим об одновременных событиях, нужно сперва определить, что такое одновременные события. В русском языке это сложно передать, тут скорее нужно переходить к термину параллельных событий.

Любые два события могут быть связаны друг с другом или не связаны. Под связью событий имеем в виду, что либо одно вытекает из другого, либо наоборот. При этом если события не связаны причинно-следственно, то можно считать из параллельными.

Поскольку параллельные события никак не влияют друг на друга, то одновременная запись двух параллельно относящихся величин не может привести к конфликтам. Поэтому имеет смысл говорить только о каузально связанных событиях в рамках разрешения конфликтов.

Для связанных записей можно ввести общий счётчик, ассоциированный с значением. При чтении клиент получает ассоциированный счётчик. При записи клиент отправляет текущее известное ему значение счётчика, а в ответ от реплики получает новое значение счётчика. На основании такой схемы при общении с одной репликой можно однозначно определить, в каком порядке изменялось значение. Алгоритм может быть обобщён на несколько реплик путём ассоциации ещё одного счётчика с репликой. Такой подход называется version vector.