Показаны сообщения с ярлыком Message-Passing. Показать все сообщения
Показаны сообщения с ярлыком Message-Passing. Показать все сообщения

суббота, 20 июля 2019 г.

[prog.c++] Идея развития темы, начатой в статье "A declarative data-processing pipeline on top of actors? Why not?"

Давеча опубликовал на Хабре большую статью на английском языке "A declarative data-processing pipeline on top of actors? Why not?", в которой описал старый, но доведенный "до ума" пример make_pipeline из состава SObjectizer-а. К самой статье Григорий Демченко задал интересные вопросы в комментарии. Плюс на Reddit-е подсунули ссылку на (как мне показалось) недоделанную и полузаброшенную библиотеку RaftLib.

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

По сути, тема создания пайплайнов (или графов в общем случае) обработки данных на C++ состоит из двух частей. Можно сказать, что есть верхняя и нижняя части проблемы.

Нижняя часть -- это то, как работа распределяется по рабочим контекстам. Т.е. это вопросы, связанные с тем, что из себя на самом деле представляют стадии пайплайна, как стадии привязываются к тем или иным рабочим нитям, как происходит передача информации между стадиями пайплайна. В принципе, нижнюю часть не обязательно делать с нуля самому. Можно использовать какой-то готовый инструмент. Скажем, Intel TBB. Или, как в моем случае, SObjectizer. С нижним уровнем связано несколько вопросов и мне хочется посмотреть, какие ответы на эти вопросы может (и может ли) предоставить SObjectizer.

Верхняя часть -- это тот DSL, который будет предоставляться пользователю для конструирования своих пайплайнов. И здесь есть ряд интересных и не исследованных для меня вопросов. Начиная от того, какую функциональность пользователь сможет получить в свои руки. И заканчивая тем, как выразить эту функциональность в C++ коде, дабы получить контроль за какими-то ошибками прямо в compile-time.

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

вторник, 29 января 2019 г.

[prog.c++] "Modern" Dining Philosophers in C++ with Actors and CSP (part 2)

In the previous post we discussed several implementations of "dining philosophers problems" based on Actor Model. In this post I will try to tell about some implementations based on ideas from CSP. Forks, philosophers and waiter will be represented as std::thread and will communicate each other via CSP-channels only.

Source code can be found in the same repository.

CSP-based Implementations

Like in the previous post we will start from implementation of Dijkstra's solution, then we will go to simple solution with putting the first fork if an attempt to take the second fork fails (without arbiter), and then we will go to solution with waiter/arbiter. It means that we will see the same solutions like in previous post, but reimplemented without Actors. Moreover the same set of messages (e.g. take_t, taken_t, busy_t, put_t) from Actor-based implementations will be reused "as is" in CSP-based implementations.

понедельник, 28 января 2019 г.

[prog.c++] "Modern" Dining Philosophers in C++ with Actors and CSP (part 1)

Some time ago a reference to an interesting article "Modern dining philosophers" was published on resources like Reddit and HackerNews. The article discusses several implementations of well-known "dining philosophers problem". All those solutions were built using task-based approach. I think the article is worth reading. Therefore, if you have not read it yet, I recommend read the article. Especially if your impressions of C++ are based on C++03 or more earlier versions of C++.

However there was something that hampered me during the reading and studying proposed solutions.

I think it was usage of task-based parallelism. There are too many tasks created and scheduled thru various executors/serializers and it's hard to understand how those tasks are related one to another, in which order they are executed and on which context.

Anyway the task-based parallelism is not the single approach to be used to solve concurrent programming problems. There are different approaches and I wanted to investigate how a solution of "dining philosophers problem" can looks like if Actors and CSP models will be used.

To do that I implemented some solutions of "dining philosopher problem" with Actors and CSP models. Source code can be found in this repository. In this post you will find description of Actors-based implementation. In the next part I will tell about CSP-based implementations.

пятница, 2 ноября 2018 г.

[prog.memories] Беглый взгляд на эволюцию SObjectizer-5.5 за прошедшие четыре года

Первый релиз SObjectizer-а в рамках ветки 5.5 состоялся чуть больше 4-х лет назад, в начале октября 2014-го. На следующей неделе планируется релиз версии 5.5.23, которая может стать финальной в рамках ветки 5.5.

Поэтому поводу есть желание написать статью для Хабра, в которой будет сделан обзор того, что появилось в SO-5.5 за это время и как это все повлияло на сам SObjectizer, и на разработку с его использованием. По ходу подготовки этой статьи сделал небольшой конспектик изменений. И сам слегка прифигел. Привожу его в текущем, еще не обработанном виде.

Хрен знает, зачем я это делаю. Может чтобы похвастаться. Может для того, чтобы те горячие головы, которые думают, что они быстренько слепят на коленке свой акторный фреймворк вместо того, чтобы взять что-то из существующего, посмотрели и задумались: а оно того стоит?

Не стоит. Берите лучше то, что есть. Не нравится вам SObjectizer -- возьмите что-нибудь другое. Тот же CAF или QP/C++. Выбор есть. Полагаю, этот выбор будет всяко лучше повторения хотя бы части пути, который мы уже прошли. И, прошу не забывать, что речь идет не только о том, чтобы придумать и запрограммировать. Но и о том, чтобы отладить, задокументировать и донести ваше творение до других людей. Которые, возможно, думают, что лучше они сами что-нибудь на коленке слепят.

Итак, вот список с ссылками на соответствующие разделы документации. Список не полный, включал в него только самые знаковые изменения/нововведения. Этот список еще предстоит переосмыслить, отранжировать и преобразовать в статью. Но, надеюсь, общее впечатление можно составить.

Да, и в этом списке нет того, что вошло в состав so_5_extra.

среда, 31 октября 2018 г.

[prog.c++] Вторая бета-версия SObjectizer-5.5.23 и so_5_extra-1.2.0

Зафиксирована вторая бета-версия для SObjectizer-5.5.23 и so_5_extra-1.2.0. Загрузить их можно отсюда: so-5.5.23-beta2.zip и so_5_extra-1.2.0-beta2.zip (либо so_5_extra-1.2.0-beta2-full.zip).

Был исправлен просчет, допущенный при подготовке первой версии: нововведение, конверты с сообщениями, не дружили с фильтрами доставки (посыпаю голову пеплом, тупо забыл про интеграцию с фильтрами доставки). Сейчас это исправлено. Но при этом поменялся интерфейс класса envelope_t. Теперь в нем вместо двух hook-методов всего один: access_hook()

virtual void access_hook(
   access_context_t context,
   handler_invoker_t & invoker) noexcept = 0;

Метод access_hook() вызывается всегда, когда нужно получить доступ к содержимому конверта. Если конверт готов предоставить доступ, то конверт должен вызвать invoker.invoke() и передать туда ссылку на содержимое (тут все осталось как и прежде).

А вот контекст, в котором вызывается access_hook(), теперь определяется перечислением access_context_t. На данный момент в нем определены следующие варианты: handler_found (содержимое конверта нужно для вызова обработчика сообщения), transformation (содержимое конверта должно быть преобразовано в другое представление, например, в результате limit_then_transform) и inspection (содержимое конверта должно быть проанализировано, например, фильтром доставки). Со временем, возможно, список вариантов будет расширен.

Собственно, это главные изменения в so-5.5.23-beta2. Соответствующим образом были изменены нужные части в so_5_extra-1.2.0 и обновлена документация в Wiki проекта.

По срокам окончательного релиза so-5.5.23 и so_5_extra-1.2.0 прогноз пока такой же: первая декада ноября. Т.е., скорее всего, на следующей неделе. На этой вряд ли получится закрыть оставшиеся вопросы и выкатить релиз без суеты и спешки.

пятница, 19 октября 2018 г.

[prog.c++] Стали доступны первые бета-версии SObjectizer-5.5.23 и so_5_extra-1.2.0

Сегодня были зафиксированы первые бета-версии наших проектов SObjectizer и so_5_extra. Загрузить их можно отсюда: so-5.5.23-beta1.zip и so_5_extra-1.2.0-beta1.zip (либо so_5_extra-1.2.0-beta1-full.zip).

Подробнее об нововведениях рассказывается в очередной статье на Хабре.

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

В общем, если кому-то интересно, что в SObjectizer-е происходит, то смотрим, делимся впечатлениями и соображениями. Пока еще есть время и возможность повлиять на то, что попадет в SO-5.5.23 и so_5_extra-1.2.0.

вторник, 16 октября 2018 г.

[prog.c++.sobjectizer] Сбылась мечта идиота: можно сделать отчеты о доставке сообщения до получателя.

Версии SO-5.5.23 и so_5_extra-1.2.0 уже начинают дышать полной грудью. В SO-5.5.23 был реализован новый механизм, который позволяет вкладывать сообщение в некий "конверт". Внутри SObjectizer-а доставка идет уже для всего конверта с вложенным в него сообщением. Сообщение из конверта достается либо когда сообщение доставлено до получателя. Либо когда сообщение по какой-то причине преобразуется из одного представления в другое (например, так происходит в случае использования limit_then_transform). Причем под "доставлено до получателя" означает не только то, что получатель извлек сообщение из очереди, но и то, что у получателя был найден обработчик для этого типа сообщения. Т.е. сообщение реально доставлено, а не просто взято из очереди.

На базе этого механизма в so_5_extra-1.2.0 добавлены средства для реализации как гарантированно отзывных таймеров, так и просто для реализации отзывных сообщений. Для отзывных сообщений, кстати говоря, сразу же придумываются сценарии, в которых их можно использовать. Так, что, наверное, эта штука может быть востребована.

Ну а в качестве примера создания собственных "конвертов" в состав so_5_extra включена самая примитивна реализация такой прикольной штуки, как отчеты о доставке отосланных сообщений до получателя.

Эти самые отчеты о доставке -- это идея фикс, которая витала в воздухе очень и очень давно. Внимание она привлекает потому, что когда агент A отсылает сообщение агенту B, то далеко не факт, что сообщение до агента B вообще дойдет. Например, сообщение может быть отвергнуто механизмом защиты агента от перегрузки (например, limit_then_drop). И временами хотелось бы знать, сообщение до B не дошло или все-таки дошло. Но вот узнать это не представлялось возможным.

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

вторник, 18 сентября 2018 г.

[prog.c++] Let's talk about hierarchical finite state machines and their support in SObjectizer-5.5

Finite state machine is a probably one of most basic and widespread thing in software development. Finite state machines (FSM) are actively used in many areas. For example in such niches as SCADA- and telecom systems FSM are used almost everywhere.

In this article we will try to speak about hierarchical FSM. And then we will try to take a look at FSM's support in SObjectizer-5. SObjectizer is one of a few OpenSource and live "actor" frameworks for C++ and actors in SObjectizer are FSMs. We will speak why SObjectizer's actors are FSMs and which capabilities they have.

A Brief Introduction Into Finite State Machines

It is hard to explain in a shot article such big topics as automata theory in general and finite state machines in particular. Because of that an basic understanding of these topics are required from a reader.

Advanced Finite State Machines And Their Features

There are several "advanced features" of FSM which significantly simplify usage of FSM. Let's talk about them briefly.

Disclaimer: if a reader has a good knowledge of UML's statechart diagrams then he or she doesn't find anything new here.

понедельник, 10 сентября 2018 г.

[prog.c++] Видео с митапа в Питере

Стало доступно видео моего выступления на митапе St. Petersburg C++ User Group с докладом "Акторы в C++: взгляд старого практикующего актородела":

Презентацию можно скачать в виде PDF-ки отсюда, либо же просмотреть на SlideShare или на Google Docs.

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

PS. Так же можно договорится, чтобы я подъехал в офис вашей компании и сделал подобный, но более подробный, доклад про SObjectizer.

пятница, 6 апреля 2018 г.

[prog.c++] Только что довелось самому воспользоваться msg_tracing-фильтрами при поиске проблемы

Делаю примеры для новой статьи про SObjectizer. В примерах имитируется работа с оборудованием и собирается небольшая статистика о задержках в обработке сообщений при разных политиках обработки этих сообщений. По сути, в примерах работает два агента. Первый, a_device_manager, имитирует работу с оборудованием. Второй, a_dashboard, раз в пять секунд выдает на экран текущие показатели.

И вот при первых запусках примеров обнаруживается, что в какой-то момент статистика перестает обновляться. Становится понятно, что это происходит потому, что не обрабатывается одно из сообщений, a_device_manager_t::reinit_device_t. Но почему не обрабатывается? Подписка на него есть:

вторник, 4 июля 2017 г.

[prog.c++] Новая большая статья про SObjectizer на Хабре

Мы сделали разбор штатного примера machine_control из дистрибутива SObjectizer в виде большой статьи на Хабре: Имитируем управление устройствами с помощью акторов. По ходу дела была мысль вместо одной статьи сделать две, а может и три поменьше. Но по опыту выходит, что каждую следующую читают меньше, чем предыдущие, поэтому решено было ограничиться всего одной.

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

Темы для следующих статей принимаются :)

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

четверг, 1 июня 2017 г.

[prog] Очередная статья про SObjectizer для Хабра: будет ли интересен разбор примера machine_control?

Обдумываю тему следующей статьи на Хабре о SObjectizer-е. Появилась идея подробно описать штатный пример из дистрибутива SObjectizer-а под названием machine_control. Этот пример имитирует управление промышленным оборудованием: контролирует температуру неких двигателей, включает вентиляторы для охлаждения или вообще выключает двигатели при перегреве. Пару лет назад я об этом примере уже писал в блоге, но там не было детального разбора агентов, их принципов работы и взаимосвязей. Оттуда же и вот эта картинка:

Собственно, почему может быть интересно подробнее описать этот пример? Потому, что в нем задействованы почти все самые важные фичи SObjectizer-а. Включая и возможность создания агентов-шаблонов. И хотя пример остается все-таки абстрактным (имитация она и есть имитация), но он совсем не игрушечный, в отличии от ping-pong-а или обедающих философов. Хотя бы чуть-чуть, но приоткрывающий завесу над тем, как на SObjectizer выглядит более-менее приближенный к реальности код.

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

В общем, если кому-то интересна такая статья, то дайте об этом знать: либо в комментариях, либо через +1 в G+, либо через лайки в FB и LinkedIn (где я размещу ссылки на этот пост).

Если не интересно, то об этом так же можно (и даже нужно) заявить прямо. А еще лучше сказать, статья на какую тему вокруг SObjectizer-а вам была бы интересна.

четверг, 25 мая 2017 г.

[prog] Интересный момент при проектировании round-robin mbox-а.

Работаю над такой штукой для SObjectizer-а как round-robin mbox. Т.е. почтовый ящик, на который могут подписаться N агентов, но доставка сообщений к ним будет осуществляться по очереди, т.е. сперва первому агенту, затем второму, затем третьему и т.д. Эта штука должна быть удобной, например, при реализации балансировки нагрузки, когда требующие обработки сообщения отсылаются в один и тот же mbox, а распределяться они будут при этом на N агентов-воркеров, прозрачным для отправителя образом.

Подобные вещи и раньше можно было делать, но вручную. А теперь есть возможность предоставить round-robin mbox прямо "из коробки".

Однако, в процессе продумывания реализации round-robin mbox-а обнаружилась неожиданная засада. Дело в том, что в mbox-ах есть такая штука, как delivery filters, т.е. фильтры, которые определяют, имеет ли смысл доставлять конкретное сообщение до конкретного получателя.

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

Чем же мешают фильтры доставки в round-robin mbox-е?

Как в принципе должен работать round-robin mbox? У него должен быть список подписчиков на каждый тип сообщения. В этом списке должен быть явно отмечен текущий элемент -- тот подписчик, которому должно быть доставлено следующее сообщение. Когда сообщение приходит, оно отсылается этому подписчику, после чего текущим элементом становится следующий подписчик (ну или первый, если мы дошли до конца). Все просто.

Но это мы пока не рассматривали фильтры доставки. В случае с фильтрами доставки мы должны для текущего элемента спросить у его фильтра "А можно ли доставлять сообщение этому подписчику?" Если можно, то все хорошо, работает привычная схема. А вот если фильтр говорит "Нет"? Что делать тогда?

Вырисовываются следующие варианты:

  1. Просто выбросить этот экземпляр сообщения. Т.е. текущему агенту в списке сообщение не доставляется (что естественно, т.к. фильтр запрещает доставку) и не делается попытка доставить сообщение следующему агенту в списке. Текущим становится следующий агент в списке (или первый, если достигли конца списка).
  2. Попытаться выбрать другого получателя. Т.е. сразу перейти к следующему агенту и спросить его фильтр доставки. Потом к следующему и т.д. до тех пор, пока получатель не будет найден. Может быть вообще ничего не найдем, но тогда данный экземпляр сообщения просто выбрасывается.
  3. Тупой запрет на использование фильтров доставки в случае round-robin mbox-а. Тогда проблемы нет в принципе.
  4. Update. Объединить все фильтры доставки в один комбо-фильтр. Когда появляется сообщение, оно сперва пропускается через этот комбо-фильтр (т.е. через все фильтры). И только если комбо-фильтр пропускает сообщение, только тогда оно доставляется текущему получателю.

Лично я пока склоняюсь к первому варианту. Логика такая: round-robin предполагает, что пытаемся доставлять по очереди. Сейчас очередь у агента X. Если X отказывается принимать сообщение, то он все равно свою очередь использовал и для следующего сообщения очередь перейдет к другому получателю.

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

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

Кстати, попутно еще один вопрос: а кроме round-robin еще какие-нибудь политики доставки для N агентов кому-нибудь нужны/интересны?

пятница, 7 апреля 2017 г.

[prog.c++] Тема movable/modifiable сообщений в SObjectizer задышала

В конце февраля, после C++ Russia 2017, возникла тема добавления в SObjectizer поддержки не константных экземпляров сообщений. Суть в том, что в SObjectizer из-за первоначальной направленности на взаимодействие в режиме 1:N, все сообщения доставлялись получателям в виде константных объектов. Получатель имел константную ссылку на пришедшее к нему сообщение и не должен был менять этот объект, т.к. в это же время то же самое сообщение мог обрабатывать еще кто-то.

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

В итоге мы начали думать, как дать возможность пересылать не константные объекты-сообщения при взаимодействии агентов в режиме 1:1. Сначала идеи крутились вокруг movable-сообщений, т.е. сообщений, содержимое которых можно перемещать между экземплярами (например, когда кто-то делает send, то содержимое сообщение сразу же конструируется где-то в заранее созданном буфере у агента-получателя). Но затем была выбрана, как мне представляется, более простая идея: мы применили к сообщениям тот же подход, который в C++ используется для динамически-созданных объектов и unique_ptr. Класс unique_ptr говорит разработчику, что есть только один "владеющий" указатель на объект. Нельзя расплодить unique_ptr, можно лишь передавать право владения так, что только один unique_ptr будет владеть динамически созданным объектом.

Мы добавили в SObjectizer такое понятие, как unique_mbox. Т.е. специальный тип mbox-а, отсылка сообщения в который приводит к передаче не константного объекта-сообщения. Для того, чтобы получить не константное сообщение из unique_mbox-а, обработчик сообщения должен использовать специальный новый тип unique_mhood_t<Msg>. Это тонкая шаблонная обертка, похожая на уже существующий mhood_t, но она позволяет модифицировать полученный объект-сообщение. Так же unique_mhood_t позволяет переслать не константное сообщение дальше.

Вот небольшой кусок реального unit-теста из SObjectizer, который проверяет работу этого механизма. Там агент отправляет самому себе сообщение типа std::string. Получив сообщение, агент меняет текст сообщения и перепосылает измененный объект. При этом проверяется то, что повторно приходит именно тот же самый std::string, но с другим содержимым.

понедельник, 27 февраля 2017 г.

[prog.thoughts] Moveable-сообщения для взаимодействия агентов в SObjectizer?

После доклада на C++ Russia 2017 возник интересный вопрос из зала. Речь о том, что сейчас в SO-5 все сообщения (но не сигналы) доставляются как динамически создаваемые объекты. Т.е. за вызовом send<Msg>(...); скрывается сперва new Msg, а затем уже передача указателя на этот экземпляр во все нужные очереди заявок. Основная причина для того, чтобы так делать в том, что у нас возможно доставка сообщения как режиме 1:1, так и в режиме 1:N. И в случае с 1:N мы либо имеем один динамически созданный экземпляр со счетчиком ссылок, либо были бы вынуждены копировать экземпляр сообщения для каждого получателя (что, имхо, гораздо хуже в общем случае).

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

И вопрос из зала состоял в том, а можно ли для случаев, когда используется только лишь взаимодействие 1:1, сделать такую оптимизацию, чтобы send приводил не к вызову new, а к выполнению move-операции. Т.е., чтобы содержимое сообщения мувилось бы куда-то в очередь получателя без дополнительного new.

В текущей реализации механизма доставки сообщений такой подход с moveable-сообщениями, вероятно, не сделать. Но я обещал подумать на эту тему.

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

Главный вопрос вот в чем: интересно ли это кому-то, кто следит за SObjectizer-ом и пытается примерить SObjectizer для решения своих задач? Если интересно, то дайте знать. Во-первых, это простимулирует дальнейшие работы в данном направлении. Во-вторых, я смогу выносить на обсуждения варианты, которые приходят в голову. Вы сможете повлиять на то, что и как в SObjectizer заработает.

В общем, если тема moveable-сообщений для взаимодействия 1:1 (а может даже и 1:N) кому-то интересно, то дайте знать. Либо в комментариях к этой заметке (можно в G+), либо по почте eao197 на gmail тчк com или info на stiffstream тчк com, либо со мной можно связаться через FB, LinkedIn или Habrhabr.

воскресенье, 19 февраля 2017 г.

[prog.thoughts] Несколько соображений на тему "А о чем же Модель Акторов?"

В конце своего доклада на C++ CoreHard Winter 2017 я сказал о том, что если кто-то ждет высокой производительности от универсальных акторных фреймворков, то это напрасно. И добавил, что SObjectizer -- это не про производительность, а про другое. Подразумевая при этом не только SObjectizer, но и вообще реализации Модели Акторов, включая Erlang и Akka. C CAF-ом ситуация, имхо, чутка другая -- они сами говорят про high performance, ну и флаг им в руки, ибо замеры говорят сами за себя ;)

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

вторник, 31 января 2017 г.

[prog.flame] Непонятое из рассказа про миграцию со Scala на Go (глазами пользователя SObjectizer)

Нашел, кажется, на RSDN-е ссылку на прекрасное: "Making the move from Scala to Go, and why we’re not going back". Там люди рассказывают про то, как они писали-писали на Scala, задолбались и с удовольствием перешли на Go. От чего сейчас буквально писают кипятком очень счастливы.

На мой взгляд, когда люди меняют высокоуровневый язык, вроде Scala, на низкоуровневую примитивщину вроде Go, то это означает, что они что-то делают неправильно. Скорее всего, неправильно сделали с самого начала: выбрали экскаватор вместо лопаты там, где нужно было выкопать яму под нужник на даче. Но это тема отдельного разговора.

В еще большем недоумении меня оставило другое: в статье описывается пример с чтением потока сообщения из Kafka и сохранения накопленных сообщений в БД. Люди попытались примерить к этой задачке Модель Акторов в реализации Akka и что-то у них не получилось. Дословно одна из причин недовольства Akka описана так: "Furthermore, the stream came from a Kafka Consumer, and in our wrapper we needed to provide a `digest` function for each consumed message that ran in a `Future`. Circumventing the issue of mixing Futures and Actors required extra head scratching time." Честно говоря я не понимаю, о чем это все.

Но если посмотреть на пример того, как они обозначают решение этой же задачи на потоках в Go, то мы видим вот что:

buffer := []kafkaMsg{}
bufferSize := 100
timeout := 100 * time.Millisecond

for {
  select {
    case kafkaMsg := <-channel:
      buffer = append(buffer, kafkaMsg)
      if len(buffer) >= bufferSize {
        persist()
      }
    case<-time.After(timeout):
      persist()
  }
}

func persist() {
      insert(buffer)
      buffer = buffer[:0]
}

Честно говоря, я не понимаю, откуда могут возникнуть затруднения в реализации этого же на акторах. Вот, скажем, у нас в SObjectizer-5 это могло бы выглядеть вот так:

class kafka_msg_persister final : public so_5::agent_t {
   std::vector< kafka_msg > buffer_;
   static constexpr std::size_t buffer_size = 100;
   so_5::timer_id_t dump_pause_;

public :
   kafka_msg_persister(context_t ctx) : so_5::agent_t(std::move(ctx)) {
      struct dump_by_timer : public so_5::signal_t {};

      so_subscribe_self()
         .event([this](const kafka_msg & msg) {
            buffer_.push_back(msg);
            if(buffer_.size() >= buffer_size)
               persist();
            else
               dump_pause_ = so_5::send_periodic<dump_by_timer>(*this, 100ms, 0ms );
         })
         .event<dump_by_timer>([this]{ persist(); });
   }

private :
   void persist() {
      insert(buffer_);
      buffer_.clear();
      dump_pause_.reset();
   }
};

Если я правильно понимаю пример кода на Go выше, то логика у него какая-то странная: сохранение сообщений в БД осуществляется либо когда буфер полностью заполняется, либо же спустя 100ms после получения последнего сообщения. Странность здесь такая. Обычно тайм-аут для сохранения попавших в буфер сообщений отсчитывают от момента получения первого сообщения. А не от момента получения последнего. Ведь если у нас, скажем, есть поток сообщений с темпом 90ms, то первое сообщение будет сохранено в БД только спустя 9000ms (т.е. 9s) после получения.

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

class kafka_msg_persister2 final : public so_5::agent_t {
   std::vector< kafka_msg > buffer_;
   static constexpr std::size_t buffer_size = 100;
   so_5::timer_id_t dump_pause_;

public :
   kafka_msg_persister2(context_t ctx) : so_5::agent_t(std::move(ctx)) {
      struct dump_by_timer : public so_5::signal_t {};

      so_subscribe_self()
         .event([this](const kafka_msg & msg) {
            if(buffer_.empty())
               dump_pause_ = so_5::send_periodic<dump_by_timer>(*this, 100ms, 0ms );

            buffer_.push_back(msg);
            if(buffer_.size() >= buffer_size)
               persist();
         })
         .event<dump_by_timer>([this]{ persist(); });
   }

private :
   void persist() {
      insert(buffer_);
      buffer_.clear();
      dump_pause_.reset();
   }
};

Либо же можно было бы использовать факт того, что агент в SO-5 -- это конечный автомат, на переходы между состояниями которого можно повесить обработчики:

class kafka_msg_persister3 final : public so_5::agent_t {
   std::vector< kafka_msg > buffer_;
   static constexpr std::size_t buffer_size = 100;

   state_t empty{this}, not_empty{this};

   so_5::timer_id_t dump_pause_;

public :
   kafka_msg_persister3(context_t ctx) : so_5::agent_t(std::move(ctx)) {
      struct dump_by_timer : public so_5::signal_t {};

      this >>= empty;

      empty
         .on_enter([this]{ dump_pause_.reset(); })
         .transfer_to_state<kafka_msg>(not_empty);

      not_empty
         .on_enter([this]{
            dump_pause_ = so_5::send_periodic<dump_by_timer>(*this, 100ms, 0ms );
         })
         .event([this](const kafka_msg & msg) {
            buffer_.push_back(msg);
            if(buffer_.size() >= buffer_size) {
               persist();
               this >>= empty;
            }
         });
         .event<dump_by_timer>([this]{ persist(); });
   }

private :
   void persist() {
      insert(buffer_);
      buffer_.clear();
   }
};

Т.е. при входе в состояние not_empty мы взводим таймер на 100ms, при входе в empty сбрасываем таймер, поскольку сейчас он нам уже не нужен.

Очевидно, что примеры на C++ и SO-5 более многословны, чем код на Go. Но задача была в том, чтобы показать, что на акторах накопление сообщений в буфер и затем сброс накопленных сообщений в БД по таймеру или по исчерпанию буфера -- это не сложно. Откуда у людей с этим возникли проблемы не понятно. Может быть дело в Akka, может быть они просто Akka готовить не умеют. Не знаю.

Очевидно, что есть люди, которым модель акторов в принципе не нравится. И которые предпочитают использовать CSP-ные каналы. Ну чтож, попробуем изобразить этот же пример на CSP-шных каналах, в SObjectizer-5:

std::vector<kafka_msg> buffer;
static constexpr std::size_t buffer_size = 100;

for(;;) {
   receive(from(chain).handle_n(buffer_size).empty_timeout(100ms),
      [&](const kafka_msg & msg) {
         buffer.push_back(msg);
      });
   if(!buffer.empty())
      persist(buffer);
}

Здесь мы выходим из receive либо после получения 100 сообщений kafka_msg, либо после того, как канал был пуст в течении 100ms. Как раз то, что было в исходном примере на Go.

А вот получить на каналах поведение, когда тайм-аут нужно отсчитывать не от последнего сообщения, а от первого полученного, с ходу не получается. Тут нужно думать. И, как мне кажется, на акторах такое делается проще, чем на каналах. Ну либо нужно использовать сразу несколько каналов. Очень грубо это может выглядеть так:

std::vector<kafka_msg> buffer;
static constexpr std::size_t buffer_size = 100;

struct dump_by_timer : public so_5::signal_t {};
auto timer_chain = create_mchain(env);

select( so_5::from_all(),
   case_(chain, [&](const kafka_msg & msg) {
         if(buffer.empty())
            so_5::send_delayed<dump_by_timer>(timer_chain, 100ms);
         buffer.push_back(msg);
         if(buffer.size() >= buffer_size)
            persist(buffer);
      }),
   case_(timer_chain, [&](mhood_t<dump_by_timer>) {
         if(!buffer.empty())
            persist(buffer);
      }) );

Но здесь возможно множественное срабатывание таймеров. Например, если идет очень большая последовательность сообщений, то на каждое 1-ое, 101-ое, 201-ое и т.д. сообщения будет отсылаться отложенный сигнал. А потом в какой-то прекрасный момент эти сигналы начнут приходить. Что не есть большая проблема, но и не есть хорошо. Поэтому придется немного поколупаться с отменой ранее отосланных отложенных сигналов. И более-менее реальный код выглядел бы как-то так:

std::vector<kafka_msg> buffer;
static constexpr std::size_t buffer_size = 100;

struct dump_by_timer : public so_5::signal_t {};

auto timer_chain = create_mchain(env,
      // Нам нужно хранить не более одного сигнала в этом канале.
      1,
      // Место под сигнал выделим сразу.
      so_5::mchain_props::memory_usage_t::preallocated,
      // При поступлении нового сигнала старый выбрасываем за ненадобностью.
      so_5::mchain_props::overflow_reaction_t::remove_oldest);

// Идентификатор таймера нам нужен дабы была возможность отменять
// доставку отложенного сигнала.
so_5::timer_id_t dump_timer;

select( so_5::from_all(),
   case_(chain, [&](const kafka_msg & msg) {
         if(buffer.empty())
            dump_timer = so_5::send_periodic<dump_by_timer>(timer_chain, 100ms, 0s);
         buffer.push_back(msg);
         if(buffer.size() >= buffer_size) {
            dump_timer.reset();
            persist(buffer);
         }
      }),
   case_(timer_chain, [&](mhood_t<dump_by_timer>) {
         if(!buffer.empty())
            persist(buffer);
      }) );

К чему я это все веду (ну, естественно, за вычетом маркетинга SO-5)? К тому, что язык высокого уровня дает разработчику возможность создавать те абстракции, которые ему нужны для решения задачи. Нужны акторы -- можно сделать акторов, нужны CSP-шные каналы -- можно сделать каналы. Если же возможность по созданию абстракций под задачу не нужна, значит задача вполне себе решается более простым и примитивным языком. Т.е. изначально людям нужен был Go, а не Scala. И нахрена было тянуть в проект Scala, дабы затем плакаться о том, что "кололись, но жрали" -- не понятно.

Впрочем, если посмотреть на бэкграунд разработчиков:

то все встает на свои места. Лишь у одного было знакомство с Java за плечами. В общем, люди изначально не видели инструментов, которые специально создавались для нормальной разработки софта, обычными программистами, а не хипстерами. Но за Scala взялись. От и результат такой вот и получился ;)

вторник, 10 января 2017 г.

[prog.thoughts] Какое тонкое различие между message-driven и event-driven...

Продавцы Akka и, заодно, главные пиарщики Reactive Manifesto, на днях опубликовали очередной white paper: Reactive Programming versus Reactive Systems. Материал наполовину технический, наполовину маркетинговый. Но, не смотря на то, что маркетингового бла-бла-бла там порядком, почитать все-таки интересно.

Мое внимание привлекли пару фрагментов, которые описывают разницу между message-driven и event-driven подходами (выделения в цитате взяты из первоисточника):

The main difference between a Message-driven system with long-lived addressable components, and an Event-driven dataflow-driven model, is that Messages are inherently directed, Events are not. Messages have a clear, single, destination; while Events are facts for others to observe. Furthermore, messaging is preferably asynchronous, with the sending and the reception decoupled from the sender and receiver respectively.

Что, в моем кривом переводе на русский язык, будет звучать как:

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

А далее приводится цитата из Reactive Manifesto, в которой эта разница раскрывается более подробно:

A message is an item of data that is sent to a specific destination. An event is a signal emitted by a component upon reaching a given state. In a message-driven system addressable recipients await the arrival of messages and react to them, otherwise lying dormant. In an event-driven system notification listeners are attached to the sources of events such that they are invoked when the event is emitted. This means that an event-driven system focuses on addressable event sources while a message-driven system concentrates on addressable recipients.

В моем пересказе это будет звучать так:

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

Вот как-то так. Вот какие проблемы людей волнуют :)

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

А у людей из Lightbend-а есть и другая головная боль: они свои Reactive Streams пытаются навесить на акторы, а это непростая задача. Поскольку акторы по природе своей асинхронны, что плохо согласуется с непрерывными потоками данных из-за отсутствия естественного back-pressure. Тут лучше подходят CSP-каналы, но это уже совсем другая парадигма получается... И, как мне думается, именно в попытках скрестить ужа с ежом разработчики из Lightbend-а и вынуждены мешать в одну кучу message-driven и event-driven, а потом старательно объяснять в чем же именно разница между первым и вторым.

У нас в SObjectizer таких проблем нет :) У нас сообщение -- это просто способ доставки информации до того, кто в этой информации заинтересован. А событие -- это факт доставки сообщения до получателя (т.е. вызов обработчика для доставленного сообщения). При этом отсылка сообщения возможна как в режиме 1:1, т.е. у сообщения есть единственный получатель, которому целенаправленно сообщение и отсылается. Так и в режиме 1:N, т.е. у сообщения будет N получателей. Так что у нас все приложения получаются event-driven, при этом унутрях у них message-driven механизм распределения информации между агентами и рабочими потоками :)

Это, конечно же, идет в разрез с Моделью Акторов, в которой сообщение отсылается конкретному актору-получателю, а массовую рассылку в режиме 1:N нужно делать ручками. Но нам пофиг ;) Publish-Subscribe, с асинхронной доставкой сообщений подписчикам рулит и бибикает :) Хотя бы потому, что с ее помощью элементарно реализуется и взаимодействие в режиме 1:1.

PS. У попыток применить publish-subscribe для dataflow-driven обработки больших потоков данных будут те же самые проблемы, что и у Модели Акторов: т.е. отсутствие естественного механизма back-pressure. Год назад в SO-5 мы добавили CSP-шные каналы. Так что dataflow на mchain-ах в SO-5 уже делать можно. А если еще придумать, как более удобным и естественным образом отобразить mchain-ы и события агентов, то можно будет удобство dataflow-driven подхода в SO-5 вывести на совсем другой уровень.

суббота, 24 сентября 2016 г.

[prog.c++] Переполенные mchains и доставка отложенных/периодических сообщений

Пока готовил очередную статью для Хабра, выяснил, что в SObjectizer при добавлении message chains (это нечто вроде CSP-шных каналов) был допущен серьезный просчет. Дело вот в чем: mchain-ы могут использоваться для отсылки отложенных и периодических сообщений. Т.е. можно вызывать send_delayed или send_periodic, а в качестве адресата указать mchain. И сообщение "упадет" в этот mchain спустя указанное время.

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

  • можно подождать какое-то время на send-е. Если за это время место в mchain-е освободилось, то просто добавить сообщение в mchain и все. А вот если мы подождали, но места не нашлось, тогда идем к следующему пункту. Впрочем, можно сконфигурировать mchain так, чтобы ожидания вообще не было. Тогда мы сразу же идем к следующему пункту;
  • т.к. места в mchain-е нет, то SObjectizer смотрит на параметр overflow_reaction для mchain и:
    • в случае drop_newest просто игнорирует новое сообщение, которые мы пытаемся добавить в mchain;
    • в случае remove_oldest выбрасывает самое старое сообщение из mchain-а, а новое -- добавляет в mchain;
    • в случае throw_exception выбрасывает самое новое сообщение и генерирует исключение;
    • в случае abort_app просто вызывает std::abort.

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

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

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

Во-вторых, на нити таймера нельзя бросать исключения. В этом нет смысла, т.к. таймер понятия не имеет, что делать с исключением о переполнении какого-то mchain-а. Любое такое исключение просто приведет к вызову std::abort.

Тем не менее, все версии SO-5 с поддержкой mchain-ов, включая последнюю стабильную 5.5.17.1, не учитывают этих ограничений для контекста таймерной нити. И, если пользователь вызывает send_delayed для ограниченного по размеру mchain-а с ожиданием на переполнении и с реакций throw_exception, то когда время доставки сообщения наступит, а mchain будет полон, то сперва нить таймера заснет на этом mchain-е, затем будет брошено исключение, от которого все приложение упадет из-за вызова std::abort.

Такой вот недосмотр.

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

  • если нить таймера обнаруживает, что отложенное/периодическое сообщение идет в переполненный mchain, то ожидание на этом mchain-е не производится, даже если такое ожидание предписано в параметрах mchain-а. Нить таймера просто сразу начнет обрабатывать overflow_reaction. Без каких-либо задержек и ожиданий;
  • вместо throw_exceptio нить таймера выполняет реакцию drop_newest, т.е. простое выбрасывание сообщения, как будто его и не было.

Эти исправления уже реализованы и протестированы. Но в основную ветку я их пока не включил. Хочу послушать другие мнения. Может есть какие-то другие подходы к решению проблемы выполнения overflow_reaction на контексте нити таймера?