Показаны сообщения с ярлыком Amazon's Dynamo. Показать все сообщения
Показаны сообщения с ярлыком Amazon's Dynamo. Показать все сообщения

среда, 7 апреля 2010 г.

[comp.prog] Amazon’s Dynamo: заключительная заметка

Эта (давно обещанная) заметка завершает рассказ о статье Dynamo: Amazon's Highly Available Key-value Store. Вот предыдущие части:

Amazon’s Dynamo: версионность объектов
Amazon’s Dynamo: распределение объектов по узлам системы
Amazon’s Dymano: репликация объектов и диагностирование сбоев

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

Service Level Agreements (SLA). В статье упоминается, что все службы внутри Dynamo договариваются между собой о качестве предоставления услуг посредством т.н. Service Level Agreement (SLA). Например, сервер в своем SLA может указать, что он гарантирует время отклика в 300ms для 99.9% запросов при пиковой нагрузке в 500 запросов в секунду. При этом подчеркивается, что в Dynamo делают расчет именно исходя из цифры 99.9%. Т.е. если какой-то сервер заявляет свои возможности, то он заявляет их для 99.9% запросов. Так обеспечивается качество обслуживания для подавляющего большинства пользователей Amazon-овских сервисов.

Меня, правда, больше интересует, каким образом SLA используется на программном уровне при построении компонентов. Например, подключается какая-то Amazon-овская служба к нескольким узлам Dynamo и получает от каждого узла SLA. Используются ли эти SLA для распределения запросов по узлам? Или может ли узел Dynamo динамически изменять свой SLA в зависимости от количества подключившихся к нему клиентов? Скажем, первому он обеспечивает отклик в 300ms, второму – всего лишь 450, а третьему – жалкие 800ms? К сожалению, в статье об этом не говорится.

Два способа общения с узлами Dynamo. Когда какой-нибудь сервис хочет воспользоваться услугами Dynamo, он имеет выбор: либо подключиться через Dynamo-вский балансировщик нагрузки, либо использовать специальную клиентскую библиотеку. Если работать через балансировщик, то запрос проходит дополнительную стадию – выбор узла для обработки запроса. А в случае клиентской библиотеки этой стадии нет, и запрос сразу же направляется нужному узлу. Что может уменьшать время выполнения запроса в два и более раз.

Выбор узла для чтения значения. Для того, чтобы оптимизировать распределение запросов между узлами, в Dynamo применяется простой трюк: координатором для операции read назначается узел, который до этого выполнял координацию операции write. Т.е. если клиент записал объект X и координатором выступал узел N, то при последующем чтении объекта X (которое, как правило, сразу же следует за записью), узел N будет координатором операции чтения.

Read repair. Когда координатор операции read принимает ответы от реплик, он не сразу прекращает прием ответов. Вместо этого он еще некоторое время ждет ответы от всех узлов и проверяет полученные версии. Если какой-то узел присылает устаревшую версию (т.е. узел по какой-то причине пропустил обновление объекта), то на узел отсылается новая версия объекта. Эта операция называется read repair и используется она для уменьшения энтропии в случае рассогласования реплик.

Детали реализации. Все службы Amazon’s Dynamo написаны на Java. При реализации использовался событийно-ориентированный подход и организация очередей сообщений между стадиям обработки запросов по образу SEDA. В качестве хранилищ информации используются Berkeley Database Transactional Data Store, Berkeley Database Java Edition, MySQL и in-memory буфера с периодическим сбрасыванием содержимого на диск. BDB используется для хранения объектов, не превышающих размеров более нескольких десятков килобайт, тогда как MySQL задействуется для объектов большего размера. Большинство узлов Dynamo задействуют BDB Transactional Data Store.

Конечные автоматы в рамках Dynamo. Каждый запрос от клиента приводит к запуску нового экземпляра конечного автомата (по сути агента, а это не могло оставить меня равнодушным). Этот КА отслеживает все стадии обработки запроса. Один запрос – один КА. Подозреваю, что такие КА весьма легковесны и на одном узле их можно создавать сотнями тысяч.

PS. Вот, собственно, и все. Приношу читателям свои извинения за то, что рассказ об Amazon’s Dynamo затянулся на два месяца. Надеюсь, что данная серия заметок была интересной.

четверг, 11 марта 2010 г.

[comp.prog] Amazon’s Dymano: репликация объектов и диагностирование сбоев

Продолжение рассказа о статье Dynamo: Amazon's Highly Available Key-value Store. Предыдущие части: Amazon's Dynamo: версионность объектов и Amazon's Dynamo: распределение объектов по узлам системы.

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

Поэтому в Dynamo используется т.н. sloppy quorum: все операции чтения и записи выполняются на первых N живых узлах из списка предпочтений. А эти узлы не обязательно будут первыми N узлами при последовательном обходе кольца узлов с диапазонами хеш-значений.

Например, пусть есть кольцо с узлами A, B, C, D, ..., Z и N=3 (N – это количество узлов, на которых должно быть зафиксировано значение объекта). Нужно сохранить ключ, который попадает в диапазон (Z,A]. Этот ключ должен храниться на узле A. Но если узел A сейчас недоступен, то его реплика может быть записана на узле D (поскольку N=3, то в обычное время реплика объекта попадает только на узлы A, B и C, но не D). При этом в метаданных реплики помечается, что она не принадлежит D и должна быть передана на узел A. Когда такая реплика попадает на узел D, он начинает периодически опрашивать узел A. И когда обнаруживает, что A опять вернулся в строй, то передает эту реплику узлу A, уничтожая ее у себя.

Такой подход позволяет Dynamo оставаться работоспособным даже при выходе из строя изрядного количества узлов. В некоторых случаях, для оптимизации быстродействия можно пойти даже на то, чтобы установить значение W (количество узлов, которые должны подтвердить запись объекта) равным 1. В таком случае запись будет завершаться успешно даже если в сети оказывается всего один узел, подтвердивший запись объекта.

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

Другой стороной системы репликации в Dynamo является снижение расходов на синхронизацию реплик после ввода отказавшего узла в строй. Поскольку каждый узел в Dynamo хранит большое количество объектов, то нужно быстро выбрать те из них, значение которых изменилось с момента последней репликации. Для этого в Dynamo используется механизм на основе хеш-деревьев (они же Merkle tree).

Смысл хэш-дерева в том, что значение родителя является хешем от значений его потомков. Т.е., если изменяется значение самого нижнего узла дерева (объекта), то изменяются хеши во всех родителях этого узла и, в конце-концов, в вершине дерева. Поскольку дерево состоит из нескольких независимых друг от друга поддеревьев, то это позволяет распараллеливать сравнение двух хэш-деревьев -- ведь их поддеревья можно сравнивать независимо друг от друга. Например, пусть есть два хэш-дерева у вершины которых два потомка. Если значение в вершине двух деревьев совпадают – значит деревья одинаковые (синхронизации реплик не требуется). Если не совпадает, значит можно отдельно сравнить левые поддеревья и правые поддеревья и т.д.

Этот принцип и используется в Dynamo: каждый узел хранит хэш-дерево для всех своих ключей. Когда приходит время синхронизации, узлы обмениваются сначала вершинами своих хэш-деревьев. Затем, при необходимости, вершинами поддеревьев и т.д. Такой механизм успешно работает. Однако для этого разработчикам Dynamo пришлось подумать над эффективной системой распределения ключей по узлам. Поскольку в худших случаях при перераспределении диапазонов (при вводе в строй новых узлов, например) перестройка хэш-деревьев оказывалась слишком дорогостящей операцией.

С темой репликации данных связана и еще одна тема – определение сбойных узлов во время работы Dynamo.

Ввод новых узлов и плановое отключение старых узлов выполняется в Dynamo только явным образом администратором: с помощью специального инструмента администратор подключает или отключает узел к кольцу. Получивший соответствующую команду узел использует Gossip-подобный протокол для извещения остальных узлов о своем статусе. Этот протокол позволяет узлам согласовывать свое видение текущего состояния кольца узлов. Так же каждый узел периодически связывается с любым другим случайным узлом для того, чтобы синхронизировать информацию о состоянии истории подключения/отключения узлов (почему-то эта история подлежит долговременному хранению и нуждается в синхронизации).

При обслуживании запросов пользователей узлы Dynamo используют только "локальное видение" состояния других узлов. Например, если A отсылает реплики узлам B и C, а получает ответ только от C, то A спокойно считает, что узел B вышел из строя. Даже если с узлом B все нормально и узел B успешно отвечает на запросы узла C. Для узла A все это не важно – узел B не ответил, значит он мертв. Узел A больше не будет адресовать запросы B, вместо этого он будет задействовать другие доступные узлы. Но узел A будет периодически пинговать B и, когда B ответит на пинг, узел A вновь занесет B в свой список живых узлов.

Для того, чтобы в Dynamo не произошло расщепления на несколько независимых колец (т.н. logical partition вследствие сбоев нескольких узлов, связывавших разные группы), используются т.н. seed-узлы. Это узлы с фиксированными именами, про которые знают все остальные узлы Dynamo. Периодически каждый узел должен связаться с одним из seed-ов и обменяться с ним информацией о доступных узлах. Благодаря этому узлы, оказавшиеся в независимых кольцах, узнают о существовании других колец и, тем самым, устраняется возникший разрыв.

В Dynamo используется децентрализованный протокол обнаружения сбоев, за деталями которого авторы статьи адресуют читателя к работе On Scalable and efficient distributed failure detectors.

За сим данную часть рассказа об Amazon Dynamo можно закончить. Пожалуй, затем последует еще одна маленькая заметка с несколькими моментами, которые меня в статье об Amazon Dymano зацепили.

пятница, 12 февраля 2010 г.

[comp.prog] Amazon’s Dynamo: распределение объектов по узлам системы

Продолжение рассказа о статье Dynamo: Amazon's Highly Available Key-value Store. Начало можно найти здесь.

Инфраструктура Dynamo состоит из сотен тысяч серверов, разбросанных по разным дата-центрам. Это полностью децентрализованная система, узлы которой используют основанные на gossip протоколы для установления взаимосвязей и обнаружения сбойных узлов.

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

Распределение ключей между узлами выполняется с помощью несколько модифицированной схемы consistent hashing. Все “адресное пространство” MD5 значений разбивается на диапазоны. И каждому узлу Dynamo случайным образом выделяется номер диапазона. Когда для какого-то конкретного ключа вычисляется MD5 хеш, то по значению хеша устанавливается номер диапазона, к которому принадлежит ключ. И запрос на обработку этого объекта передается соответствующему узлу.

Узлы провязываются в кольцо. Т.е. узел, обслуживающий первый диапазон, связывается с узлом, обслуживающим второй и последний диапазоны (второй узел связывается с первым и третьим узлами и т.д.). Т.е. если последовательно двигаться от одного узла к другому, то можно вернуться к самому первому узлу.

Каждый узел в кольце хранит значения объектов с (N-1) предшествующих ему узлов кольца. Например, пусть N=3 и есть часть кольца с узлами A, B, C, D. Узел D будет хранить объекты с ключами, попадающими в диапазон (C,D], а так же будет хранить копии объектов с ключами из диапазонов (A,B], (B,C]. Таким образом, значения объектов реплицируются на N узлов.

Такая схема хороша тем, что при изъятии или добавлении в кольцо нового узла, задействуются только соседи слева и справа от него (для обмена репликами). А остальные узлы остаются нетронутыми.

В Dynamo введено так же понятие виртуального узла. Т.е. в действительности каждый физический узел обслуживает не один диапазон из глобального адресного пространства, а несколько. Поэтому каждый физический узел выглядит как несколько виртуальных узлов. Такое отличие от схемы consistent hashing по мнению разработчиков более выгодно, поскольку позволяет эффективнее распределять нагрузку при добавлении или изъятии узла. И, что очень важно, количество виртуальных узлов, которые будет обслуживать физический узел, может определяться мощностью физического узла. Т.е. мощная машина может обслуживать 10 виртуальных узлов, а более слабая - всего 5.

Список узлов, отвечающих за хранение конкретного ключа, называется списком предпочтений (preference list). Благодаря тому, что по ключу можно определить номер диапазона, а узлы, обслуживающие соседние диапазоны провязаны в кольцо, каждый узел способен построить список предпочтений для своего диапазона.

Со списком предпочтений связана маленькая деталь. Список предпочтений содержит N элементов. Но из-за наличия виртуальных узлов может оказаться так, что количество физических узлов в списке предпочтений будет меньше. Поэтому список предпочтений строится так, чтобы в нем находились только разные физические узлы.

Для обеспечения согласованности и надежного хранения данных Dynamo использует схемы на основе кворума. При конфигурации системы задаются параметры N (количество реплик), R (минимальное количество узлов, которые должны участвовать в операции чтения данных) и W (минимальное количество узлов, которые должны участвовать в операции записи данных). Для обеспечения кворума значения R и W выбираются так, чтобы (R+W)>N. При этом, однако, слишком большие значения будут отрицательно влиять на отзывчивость системы.

Узел, на который Dynamo адресует запрос put или get, называется узлом-координатором. Обычно в качестве узла-координатора выбирается один из первых узлов в списке предпочтений для ключа.

Получив запрос get, узел-координатор адресует его первым живым N узлам из списка предпочтений. Затем ожидает R первых ответов. Получив их, узел-координатор возвращает ответ клиенту. Если же координатор получает несколько независимых версий объекта, то он возвращает клиенту все полученные версии, чтобы клиент сам выполнил их слияние.

Получив запрос put, узел-координатор модифицирует временную метку объекта и сохраняет ее локально. После чего передает эту версии первым живым N узлам из списка предпочтений и ожидает подтверждений от них. Получив W-1 подтверждение, координатор считает запись успешной.

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

вторник, 9 февраля 2010 г.

[comp.prog] Amazon’s Dynamo: версионность объектов

Так вот о статье Dynamo: Amazon's Highly Available Key-value Store (о которой я неделю назад говорил). Статья большая. Для меня оказалась интересной. Информации в ней много, всего не перескажешь. Так что, если кому-то эта тема интересна, то советую прочитать статью целиком. Я же у себя в блоге перескажу только то, что сам из нее запомнил.

Amazon Dynamo является быстрым, высоконадежным, распределенным хранилищем информации, представленной в виде пар ключ-значение. Это хранилище используется такими требовательными к быстродействию и надежности сервисами Amazon, как списки бестселлеров, корзины покупок, предпочтения пользователей, каталог продуктов и пр. Для этих сервисов не нужны сложные реляционные модели данных. Им вполне хватает всего двух операций, предоставляемых Dynamo: put для (пере)записи данных и get для чтения.

Главной особенностью Dynamo является отношение к целостности данных. Известная четверка свойств транзакции в БД - ACID (Atomicity, Consistency, Isolation, Durability) - в Dynamo обеспечивается своеобразно. Поскольку невозможно обеспечить высокую производительность и высокую надежность (смотрим на CAP-теорему), то подход к согласованности данных в Dynamo свой собственный.

Например, пусть приложение A выполняет обновление значения для ключа K и приложение B в тот же самый момент выполняет обновление значения для того же самого ключа. Оба эти изменения будут приняты Dynamo. Каждое изменение объекта приводит к сохранению нового, неизменяемого значения. Этому значению будет приписана временная метка в виде вектора (см. vector clock). Элементами в векторе являются номера версий и имена узлов, которые сохраняли новую копию.
Например, пусть на узле S1 зафиксировали первую версию объекта D - временная метка для него будет иметь вид [[S1,1]] - т.е. первая версия на узле S1. Затем на этом же узле объект D перезаписали, и у него метка изменилась, приняла значение [[S1,2]] - получилась вторая версия на узле S1. Затем объект D модифицировало приложение A, но модификация пошла не через узел S1, как раньше, а через узел S2. Новое значение объекта получило метку [[S2,1],[S1,2]] - т.е. первая версия на узле S2 после второй версии на узле S1. В это же время объект D перезаписало приложение B, но запись пошла через узел S3. Так получилась еще одна, независимая, копия D с временной меткой [[S3,1],[S1,2]].

Так вот главная особенность Dynamo в том, что когда приложение запросит последнюю версию объекта D, то оно (при нормальной работе Dynamo), получит сразу две копии объекта - одну с меткой [[S2,1],[S1,2]], а вторую с меткой [[S3,1],[S1,2]]. И вот тут возникает вопрос: кто и как будет делать согласование этих версий?

Dynamo позволяет ответить на этот вопрос двумя способами:

  • такое согласование выполняет само Dynamo. Это очень негибкий способ, самым разумным выбором в котором будет оставление самого последнего изменения и выбрасывание всех предыдущих (т.н. last write win);
  • такое согласование выполняет запросившее данные приложение. Т.е. приложение создается так, чтобы быть способным получить несколько версий одного и того же объекта, после чего "слить" изменения из разных версий в одну согласованную версию.

Так вот, самым удивительным для меня оказалось то, что такой сервис Amazon, как "корзина покупок" как раз умеет сливать "параллельные" версии корзинки в одну.

Но вернемся к временным меткам. У нас оказался объект D с двумя разными значениями и разными временными метками. Приложение, которое такой объект получило, должно создать и сохранить его новую согласованную версию. Пусть оно это значение сделало и записало через узел S1. Тогда временная метка нового значения примет вид [[S1,3],[S2,1],[S3,1]]. Если же новое значение было сохранено через узел S4, то временная метка, подозреваю, примет вид [[S4,1],[S2,1],[S3,1],[S1,2]] (тут я не уверен на 100%, т.к. такого примера в статье не было, но думаю, что должно быть так).

Наличие имен серверов и версий во временной метке позволяет приложениям отслеживать отношения между версиями, скажем, находить общие корни (так, в приведенном выше примере было видно, что для [[S2,1],[S1,2]] и [[S3,1],[S1,2]] был общий предок - версия [S1,2]). Но, с другой стороны, временные метки могут расти в случае сбоев в сети, когда обновление версии объекта все время выполняется разными серверами. Поэтому в Dynamo используется простая схема ограничения роста временной метки: когда какой-то узел достигает определенного порога (например, во временную метку добавляется пара [Si,10], т.е. на узле Si объект модифицировался уже 10 раз), то самый старый элемент временной метки выбрасывается вообще. Потенциально, такая схема может создать сложности при слиянии версий, но на практике эта проблема ни разу не проявилась (или о ней решили не говорить ;).

Нужно еще сказать, что одна из основных целей, которые преследуются Dynamo - это высокая доступность для записи (т.н. always writable). Т.е. если приложение желает сохранить свои данные, то оно должно это сделать. Всегда. По определению ;) Именно поэтому разрешение конфликтов Dynamo выполняет не во время записи (поскольку из-за сбоев не все узлы, хранящие реплики объектов могут быть доступны для записи), а во время чтения. Поэтому-то приложения, запросив значение объекта, могут получить несколько значений с разными временными метками.

Из-за этого, как я понимаю, может произойти следующая ситуация: приложение прочитало объект D с узла S1, обновило его и попыталось записать обратно. Но узел S1 уже недоступен. Запись пошла на узел S2. После чего приложение вновь запросило объект D, но на этот раз узел S2 уже не доступен, а доступен S3, на котором лежит старая реплика с узла S1. Т.е. приложение вычитало старое значение после того, как успешно записало новое!

В принципе, если подумать, в подобной распределенной системе это вполне нормальное дело. Все довольно логично. Что меня в этом всем удивляет, так это то, что Amazon-овская "корзина покупок" работает с Dynamo. Ведь что может получиться: положил я себе в корзину книгу по Ruby (запись на S1), потом книгу по Java (запись на S2), только-только собрался класть книгу по C# (чтение объекта с S3) - глядь, а у меня в корзине только книга по Ruby. А где Java, спрашивается? ;) Вот это для меня оказалось очень и очень неожиданным, что "корзина покупок" в Amazon допускает такие коллизии (как раз они и разрешаются приложением посредством слияния параллельных веток объекта). Если бы я его проектировал, то я бы счел подобное поведение неприемлимым (и, вероятно, был бы не прав).

Пожалуй, пока на этом стоит остановиться. В следующий раз расскажу, как в Amazon Dynamo поддерживается распределение объектов между узлами и как выполняется репликация значений объектов.

Disclaimer: все вышеизложенное может являться следствием моих искренних заблуждений из-за неверного восприятия текста оригинальной статьи ;)