| 1 | 1 родительский процесс и N дочерних процессов. | В данном примере имеется в виду, что дочерние процессы могут выполняться параллельно другу и независимо друг от друга, но в конце должны оповестить родительский процесс о необходимости продолжения обработки. Если речь идет о каких-либо зависимостях порядка выполнения в дочерних процессах, то это может контролировать дочерний процесс (выделяя группу, которую сейчас можно запустить и ожидая окончания). | Вариант 1: CounterTrigger. | - Родительский процесс создает триггер со счетчиком N, создает и запускает дочерние процессы, засыпает.
- Дочерние процесс при завершении публикует TriggerEvent.
- TriggerConsumerRunner периодически считывает батч TriggerEvent, уменьшает считчик триггера и делает запись в БД. За счет агрегации событий завершения процессов мы уменьшаем нагрузку на БД.
- Когда все дочерние процессы отработали TriggerConsumerRunner получает значение счетчика 0 и взводит триггер.
- Триггер пробуждает родительский процесс для дальнейшего выполнения.
| | TriggerEvent публикуются без использования TransactionOutbox напрямую в брокер после коммита транзакции (иначе мы бы нагружали БД). | Предполагаем, что основную часть времени система работает стабильно, но допускается ситуация, что транзакция закоммитилась, но TriggerEvent не смогли опубликоваться (остановка сервиса без graceful shutdown, проблемы соединения или работы с брокером сообщений). Для таких случаев создается страхующий триггер (1 общий на тип процесса). Этот триггер запускается периодически и проходится по всем ожидающим процессам, проверяя условие (в реализации можно использовать keyset пагинацию) (в реализации можно использовать join для проверки условия). Этот триггер выполняется периодически с более крупной временной задержкой. В случае обнаружения потери TriggerEvent, он поднимет заклинивший родительский процесс и он будет обработан (но позже). Можно установить этому триггеру низкий приоритет. |  |
| Вариант №2: Мы просто ставим TimerTrigger на условно 1-5-10 минут (насколько важна задержка) и перепроверяем условие завершения. В этом случае будет - Из минус: что родительский процесс узнает о завершении дочерних процессов с задержкой.
Если дочерний процесс падает в ошибку, TimerTrigger все равно будет крутиться и создавать пустую нагрузку. - Из плюсов: будет меньше пишущей нагрузки на БД чем в варианте 1 (но больше читающей - на проверку) т.к. у нас не будет CounterTrigger, но будет периодический запрос на проверку завершения всех дочерних процессов (аналогично страхующему триггер).
- [Расширенный]: Дочерние процессы в блоке wakeup condition проверяют наличие незавершенных процессов.
- Если все процессы завершены или (незавершенных процессов мало и нет процессов с ошибкой), то можно опубликовать событие на TimerTrigger на сброс или установку минимальной задержки.
| Вариант №3: Дочерние процессы выполняются через родительский (ограничение в рамках одной ноды). Точкой выполнения является родительский процесс, который внутри себя (параллельно или последовательно) выполняет дочерние процессы. За счет такого способа у нас также отсутствует конкуренция передачи сигнала в родительский процесс. Но мы ограничены выполнением дочерних процессов одной одной сервиса. Сложнее контролировать распределение нагрузки, если будет вложенный параллелизм. Также решает проблему, если дочерний процесс содержит ожидание (например асинхронный запрос-ответ), тут будет конкуренция сигнала от хендлера ответа к родительскому процессу. | Вариант №4: SimpleStreamTrigger + Timer. - Триггер проверяет условие завершения всех дочерних процессов (можно прикинуть количество незавершенных дочерних процессов).
- Если все обработано, то пробуждает процесс и деактивируется.
- Иначе:
- деактивируется (до поступления хотя бы одного сигнала),
- взводит признак стрима - процесс ожидает,
- взводит флаг новых сигналов на 0,
- выставляет задержку от оценки количества необработанных процессов (< N - малая задержка, иначе большая задержка).
- [Расширенный]: Дочерние процессы в блоке wakeup condition проверяют наличие незавершенных процессов.
- Если все процессы завершены или (незавершенных процессов мало и нет процессов с ошибкой), то можно опубликовать событие на SimpleStreamTrigger на сброс или установку минимальной задержки (в дополнение к сигналу).
- Читающей нагрузки будет немного больше чем в варианте 2 (чтение триггера на поступлении сигнала),
но пишущей нагрузки будет меньше чем в варианте 1 (запись - только на активации новым сигналом). - Если сигналов нет, то нет пустых срабатываний в отличие от варианта 2 (т.к. нет поступления сигнала от дочерних процессов).
| Вариант №5: SimpleStreamTrigger + Счетчик (ExternalCounter) в Redis. (на текущий момент самый лучший вариант). | InMemory счетчик, нагрузка на БД и конкуренция. Дочерний процесс уменьшает счетчик. И публикует событие только если счетчик достиг 0. Совмещает преимущества из варианта 1.1 (при этом не нагружает БД), в случае ошибки переключается на режим 1.4. | Проблема: изменение счетчика не привязано к основной транзакции БД. Возможно: - В начале транзакции (или шага) атомарно проверяем MemberSet, если есть запись то удаляем и увеличиваем счетчик на 1 (означает что процесс уже уменьшал счетчик, но потом было падение).
- До коммита транзакции атомарно добавляем значение в MemberSet и уменьшаем счетчик на 1. Если счетчик равен 0, то публикуем TriggerEvent
(добавляем в ручную компенсацию вызов из пункта 1 (при откате изоляции шага и компенсации транзакции)). - После коммита транзакции удаляем запись из MemberSet.
(Если мы падаем тут, то процесс уже перешел на другой шаг или даже завершился, поэтому наличие единичной остаточной записи в memberSet не будет критичным). - Trigger получает событие. (Необязательно) в хендлере может првоерить значения счетчиков и MemberSet.
MemberSet используется для уменьшения вероятности увеличить или уменьшить счетчик дважды одним экземпляром дочернего процесса. | В случае обнаружения повреждения обработка фактически переходит в режим 1.4: начинает публиковать событий каждый раз и используется задержка. Примеры проблемы: - Падение InMemory хранилища. Предполагается режим без снимков и удаление ключей.
Обнаружение (со стороны дочернего процесса) через отсутствие ключей (проверяется в транзакции). - Дублирование обновления счетчика.
Обнаружение (со стороны дочернего процесса) через значение счетчика < 0. Обнаружение (со стороны триггера) через активацию триггера (поступления сигнала от процесса), при этом обнаруживается что не все процессы завершены.
|
|
|
| 4 | Групповое действие | | | Действие, которое нужно применить к диапазону строк (сравнительно большому), независимо для каждой строки. Наличие у строк упорядоченного столбца (для выделения диапазонов). | | | | Родительские процесс определяет границы диапазона [min, max]. | select min(), max() where condition() | | Родительский процесс нарезает диапазон [min, max] на поддиапазоны. На каждый поддиапазон создается дочерний процесс. | | | Каждый дочерний процесс обрабатывает свой поддиапазон строк (параллельно). | Внутри поддиапазона может использоваться keyset пагинация. | | Родительский процесс ожидает завершения дочерних процессов (см. пример 1). | |
|
|
| 5 | Распределение заявок между исполнителями (Заготовка). | | Описание | Есть поток заявок на деталь (создание детали требует ресурсов, 1 станок, время). Есть N станков. Опционально: у станка есть уровень ресурсов и коэффициент скорости работы. Реализация системы распределения и обработки. | | Вариант 1 | | | Планирование без очереди к станку. | | | - У процесса планировщика есть
- StreamTrigger на поток заявок.
(Можно использовать расширение SignalCode, чтобы временно игнорировать откладывать этот сигнал пока все слоты станков заняты). - StreanTrigger на поток сигналов об освобождении слота станка.
| | | Процесс планировщик назначает заявку на свободный станок. Вопрос наиболее эффективной функции выбора (оценка наибольшее количество ресурсов, наилучшая скорость обработки и др.). Когда все станки заняты планировщик ожидает освобождения станков. |
| | Вариант 2 | | | Планировщик с очередью к станку. | | | - У процесса станка есть StreamTrigger, на который планировщик подает сигнал в добавления задачи в его очередь.
- У процесса планировщика есть StreamTrigger на поток заявок.
| | | При поступлении заявки процесс планировщик сразу назначает в очередь на какой либо станок. Хранение фактического уровня ресурсов и зарезервированного уровня ресурсов (с учетом очереди к станку). Вопрос наиболее эффективной функции выбора (оценка размера очереди, достаточности у станка ресурсов для ее обработки, общего количества ресурсов, скорости работы станка и др.). |
|
|