пятница, 4 ноября 2022 г.

Отказоустойчивость в Greenplum

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

А вот на уровне сегментов в Greenplum автоматическая отказоустойчивость есть. Сегменты - это наборов экземпляров Postgres. В каждом таком сегменте есть своя главная реплика, управляющая сегментом. В случае её падения главной становится одна из подчинённых реплик.

Примечание. Greenplum - широкоиспользуемая аналитическая СУБД с открытым кодом. Один из главных конкурентов Vertica. Используется, например, в Тинькове 


среда, 9 февраля 2022 г.

Шардирование

 На Хабре появилась интересная статья А.Комягина (ссылка) о масштабировании СУБД с помощью шардов. Приведу несколько цитат 

Определение шардирования:

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

 Выбор ключа шардирования - диапазонный (range) или хэшированный?

в большинстве случаев я бы рекомендовал использовать хешированное распределение (hashed shard key). Причина проста — hash-функция позволяет даже плохой ключ, с точки зрения формальных признаков, превратить в хороший. Если ещё проще, то задача hash-функции — равномерно распределить сущности по шардам. Я мог бы углубиться в суровую математику и рассказать про критерий Пирсона и другие интересные вещи, но в этом нет необходимости, поскольку инженеры компании MongoDb Inc. уже за нас всё продумали и выбрали хорошую hash-функцию для задачи шардирования. Поэтому нам осталось просто этим всем пользоваться и наслаждаться.

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

в шардированных коллекциях не всегда работает операция count. Она может возвращать большее количество документов, чем в действительности. Причина — балансировщик, который в фоне «переливает» документы с одного шарда на другой. В какой-то момент времени возникает ситуация, когда документы уже записались на целевой шард, а на исходном ещё не удалились — count посчитает их дважды.  

Ну и интересная картинка - кластер с шардами СУБД, распределённый на два ЦОД (кликните, чтобы увеличить):


  

вторник, 21 сентября 2021 г.

Проектирование высоконагруженной системы на примере отделения банка

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

На Дзене нашлась остроумная статья (в блоге компании OTUS, где еще много интересных статей), в которой построение распр системы объясняют с помощью аллегории - вот есть отделение банка, чтоб увеличить пропускную способность делаем несколько окошек, при этом принтер будет общим (разделяемый ресурс) ну и так далее.

В итоге: 

получили на выходе сложную систему, включающую в себя:

— распараллеливание;

— предобработку;

— очередь;

— балансировку;

— конвейер;

— отложенные вычисления;

— кэширование;

— толстого клиента.

Заметка называется "Как думать при проектировании высоконагруженной системы?", читать здесь

среда, 25 марта 2020 г.

Определение завершения распределённых вычислений

Несколько ссылок


  • Слайды-книжка Кшемкальяни и Сингала Termination detection
  • Раджив Мисра Termination detection
  • Обзор статьи Мисры (видимо, другого) Detecting Termination of Distributed Computations Using Markers (1983)
А вообще как-то негусто лекционного материала на эту тему

суббота, 14 марта 2020 г.

Снимки глобального состояния - ссылки

Самый известный алгоритм для случая FIFO каналов - алгоритм Чанди-Лэмпорта, для неFIFO - алгоритм Лай-Янга. Есть ещё парочка алгоритмов для случая каузальной доставки (causal ordering).

Теперь ссылки. Сначала классические статьи

  • K. Chandy, L. Lamport Distributed Snapshots: Determining Global States of Distributed Systems (ссылка)
  • Ten H. Lai and Tao H. Yang On Distributed Snapshots (ссылка) 
  • F. Mattern Efficient Algorithms for Distributed Snapshots and Global Virtual Time Approximation (ссылка)

  • Материалы университета Принстон (курс COS 418: Distributed Systems):
    • Themis Melissaris and Daniel Suo Chandy-Lamport Snapshotting (ссылка) 
    • Kyle Jamieson Vector Clocks and Distributed Snapshots (ссылка)
    • Лабораторка Chandy-Lamport Distributed Snapshots (ссылка)
    Материалы университета МакМастера (курс CAS 769):
    • Dr. Borzoo Bonakdarpour  Introduction, Logical clocks, Snapshots (ссылка)
    Материалы университета Айовы:
    • Ghosh Distributed Snapshot (ссылка) - есть граф достижимости состояний

    понедельник, 9 марта 2020 г.

    Разбор статьи Чанди-Лэмпорта (ссылка)

    Нашёл тут интересный блог The morning paper, где разбирают разные статьи по информатике. Понравился подзаголовок журнала: A random walk through Computer Science research. Блог интересный, располагается здесь.

    Поскольку сейчас вникаю в тему распределённых снимков глобального состояния, просмотрел заметку о статье Чанди и Лэмпорта Distributed Snapshots: Determining Global States of Distributed Systems (1985). Хороший пересказ, есть пример о разноцветных шариках.

    Из другого поста того же блога (Asynchronous Distributed Snapshots for Distributed Dataflows) можно узнать, что алгоритм Чанди-Лэмпорта используется в Apache Flink (это такая штука, которая позволяет проводить вычисления на потоках данных).

    суббота, 25 января 2020 г.

    Снимки: предварительные сведения

    Глобальное состояние - совокупность состояний всех локальных состояний процессов и каналов.

    Глобальное состояние согласовано, если выполнены два условия:
    C1. если посылка сообщения по каналу от P_i к P_j входит в локальное состояние P_i, то либо это сообщение входит в состояние канала C_ij, либо событие приёма этого сообщения входит в локальное состояние процесса P_j. Имеются ввиду состояния канала C_ij и локальные состояния процессов P_i и P_j, входящие в данное глобальное состояние. 
    C2. если посылка сообщения по каналу от P_i к P_j не входит в локальное состояние P_i, то это сообщение не входит в состояние канала C_ij и не входит в локальное состояние процесса P_j.

    Две основные проблемы при записи состояний:
    П1. Как отличить сообщения, которые нужно записать в снимок от тех, которые записывать не нужно. С1 и C2 говорят, что нужно записывать те сообщения от P_i, которые были посланы процессом P_i до фиксации своего локального состояния.
    П2. Как определить момент, когда процессу нужно сделать снимок. С2 говорит, что процесс P_j должен зафиксировать своё состояние до обработки сообщения от P_i, которое было послано  P_i после фиксации своего состояния.

    Снимки глобального состояния

    Непростая тема. И вроде не такая уж большая . Нужно только понять, что такое согласованное состояние и изучить два алгоритма - Чанди-Лэмпорта (требует FIFO каналов) и Лая-Янга (нет требования FIFO). А что-то трудно идёт.

    По книге Фоккинга разбираться сложно. Гораздо лучше написано у Кшемкальяни и Сингала. Вот ссылка на их лекцию (на самом деле, это просто кусок их книги, перенесённый на слайды).

    У Фоккинга нашёл только хороший перевод понятия FIFO-канала : каналы с обработкой сообщений в порядке очереди.

    Дополнение. Почитав Кшемкальяни (упомянутую лекцию и их же статью An introduction to snapshot algorithms in distributed computing), понял, что весьма упрощённо представлял себе область. На самом деле, там три группы алгоритмов:

    1. для FIFO-каналов. Здесь алгоритм Чанди-Лэмпорта и его модификации ( Spezialetti and Kearns, Venkatesan, Helary). 
    2. для не-FIFO.  Алгоритм Лая-Янга, Маттерна
    3. для систем с каузальной доставкой. Алгоритмы Acharya-Badrinath, Alagar-Venkatesan.

    вторник, 21 января 2020 г.

    Алгоритм банкира

    Здесь я пишу о своём нынешнем понимании алгоритма банкира, которое пока является неполным.

    Этот алгоритм - стратегия, которая позволяет ресурсы потребителям, избегая тупиковых ситуаций (deadlock avoidance). Такой алгоритм пригодится, например, операционной системе, которая распределяет процессам память, даёт доступ к принтеру и т.д.
    Тупиком здесь является ситуация, когда у ОС (или у банкира) нет достаточно ресурсов, чтобы удовлетворить запрос хотя бы одного потребителя (это моё понимание). Получается, что потребители ждут выделения ресурсов  от ОС, а ОС в свою очередь ждёт, пока ему вернут достаточно ресурсов, чтобы можно было выдать их кому-нибудь. 

    Нам дано: количество ресурсов и количество потребителей (процессов). Важное ограничение алгоритма: необходимо знать максимальную потребность в каждом ресурсе каждого потребителя. 

    В двух словах алгоритм выглядит следующим образом. У нас есть Available - вектор наличных ресурсов (один на всю систему) и Need - вектор потребностей каждого процесса. Просматриваем потребности процессов, начиная с первого. Если у i-го процесса Need <= Available, то система выдаст ему все нужные ресурсы. Процесс отработает, вернёт ресурсы и будет исключен из списка работающих процессов. Далее опять происходит просмотр обновлённого списка работающих процессов, определяется следующий процесс, который получит ресурсы.  

    Ссылки:
    Dijkstra’s Banker’s algorithm detailed explanation [ссылка] - лучшая
    Wiki: Banker’s algorithm [ссылка]
    What is Banker's Algorithm? [ссылка]


    четверг, 28 ноября 2019 г.

    Репликация и модели согласованности

    Нашёл несколько замечательных материалов по теме.

    Consistency Models and Protocols in Distributed System [ссылка] в блоге некоего Qing-а
    А еще хорошие  статей Майкла Уиттекера (молод, но уже много написал): сборник Consistency in Distributed Systems [ссылка] и его блог [ссылка]

    среда, 23 октября 2019 г.

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

    Наткнулся на статью со смелым названием Distributed Data Processing 101 – The Only Guide You’ll Ever Need (ссылка). И знаете, автор не подкачал. Кратко пробежался по основным уровням обработки (сбор, хранение, обработка, безопасность, визуализация). И чуть-чуть об основных архитектурах - Каппа и Лямбда. Вполне себе неплохая отправная точка для новичка в теме.

    И вообще, на 8bitmen.com много интересных статей об ИТ-архитектуре банков и др продвинутых в этом плане организаций

    Смарт-контракты с нуля

    Как научиться писать смарт-контракты на Solidity, не имея вообще никакого представления о предмете?

    Начать можно с онлайн-компилятора, который называется Remix (remix.ethereum.org). Уроков по написанию контрактов на Solidity в этом компиляторе полно в интернете. Проблема только в том, что многие из них написаны для старой версии Remix (а иногда для старой версии Solidity). Руководство, использующее новую версию, можно найти здесь.

    Документация по языку Solidity - по адресу solidity.readthedocs.io, введение с примерами самых простых контрактов - здесь.


    воскресенье, 13 октября 2019 г.

    Прямо как я

    Учиться чему-то и вести онлайн-дневник о процессе обучения - хорошая идея, так делаю не только я. Вот пост некой Фло из Нью-Йорка Live notetaking as I learn about distributed computing
    (ссылка). Человек разработал для себя некий план (плюс список литературы), который должен помочь в изучении распределённых вычислений. Есть, что позаимствовать.

    Несколько хороших статей об блокчейне

    Сравнение Ethereum и Bitcoin на blockgeeks.com - статья Bitcoin VS Ethereum: [The Ultimate Step-by-Step Comparison Guide] (ссылка). Много интересных деталей:

    1. Основные даты в истории развития этих двух криптовалют
    2. Краткое описание proof-of-stake (протокол Casper). Майнеры делают некоторую ставку (из своих средств) при попытке сформировать новый блок. Если именно этот блок будет включён в блокчейн - майнер получит награду, пропорционально своей ставке. Если же будет замечена какая-то мошенническая активность с его стороны- его монеты будут списаны
    3. Пользователь в своей транзакции может сам указывать размер награды майнеру. Чем больше награда - тем быстрее транзакция попадёт в блок
    4. В Эфире различные операции стоят по-разному в единицах газа - приведена таблица "цен"
    5. Эфир, как и биткоин, начался со статьи. Статьи Виталика, разумеется 😊
    6. Аргументы за и против увеличения размера блока. Этот спор расколол Биткоин (произошёл хардфорк) на Bitcoin (который не стал увеличивать размер блока, вместо этого включил технологию) SegWit и Bitcoin Cash (без SegWit, размер блока стал 8 MБ).
    7. В Эфире нет максимального размера блока, есть предел газа. Можно добавлять транзакции в блок, пока суммарное количество нужного для этих транзакций газа не превысит порогового значения.

    среда, 9 октября 2019 г.

    Матричные часы

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

    Если у процесса Pi есть матричные часы Mi, то мы можем сказать, что Mi[j][k] - это представление i-го процесса о том, что знает j-й процесс про локальное время k-го процесса.

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

    Допустим процесс номер 2 хочет переслать процессу номер 5 первые 30 записей из лога процесса 10. Сначала этот процесс посмотрит на соответствующую компоненту своих матричных часов M2[5][10]. Пусть M2[5][10] =23. Это означает, что по сведениям процесса 2 процесс 5 получил первые 23 записи из лога процесса 10. Значит, имеет смысл 2му процессу выслать 5му записи из лога 10го с 24й до 30й. Итог - экономия на пересылках.

    воскресенье, 6 октября 2019 г.

    Скрытые каналы

    Что, если компоненты системы шлют сообщения по разным каналам? Отслеживать причинно-следственные связи станет очень сложно. Нечто подобное было в знаменитой статье Лэмпорта про время, часы и упорядочивание событий. Там люди кроме посылки сообщений ещё и по телефону говорили.

    В науке это называется скрытые каналы (hidden channels). Пара примеров таких каналов есть в книге Distributed Systems for System Architects на стр 54 (есть на гугло-книгах, ссылка).

    суббота, 5 октября 2019 г.

    Книги по распределённым системам

    Захотелось составить список литературы по распределенным системам. Сегодня книги.
  • A. S. Tanenbaum, M. van Steen Distributed Systems: Principles and Paradigms 
  • G. Coulouris et al  Distributed Systems: Concepts and Design
  • S. Ghosh Distributed Systems: An Algorithmic Approach
  • A.D. Kshemkalyani, M. Singhal Distributed Computing: Principles, Algorithms, and Systems
  • N. A. Lynch Distributed Algorithms 
  • G. Tel Introduction to Distributed Algorithms  
  • R. Sharp Principles of Protocol Design
  • суббота, 28 сентября 2019 г.

    Пара блогов о распределенных системах

    Есть люди, которые пишут много. В том числе и о распределенных системах.

    Блог Уиттекера [ссылка] - Single-Decree Paxos, Visualizing Linearizability, Lamport's Logical Clocks,  Two Generals and Time Machines, An Illustrated Proof of the CAP Theorem

    Блог Маркуса Чиу [ссылка]  - о распределенных алгоритмах (плюс миллион статей на другие темы).

    Блог compiosition.al - здесь интересная статья про алгоритм снимков распределенной системы Чанди-Лэмпорта  [ссылка]

    Страница семинара по распределенным системах (посты пишут сами студенты) [ссылка]

    Гибридный консенсус

    Proof of Stake and the History of Distributed Consensus: Part 1, Nakamoto Consensus, Byzantine Fault Tolerance, Hybrid Consensus, Thunderella [ссылка]

    Линеаризуемость (ссылки)

    A beginner’s guide to Linearizability [ссылка]
    What is “Linearizability”? [ссылка]