Исходный код вики Примеры
Версия 8.45 от Alexandr Fokin на 2026/08/18 12:22
Последние авторы
| author | version | line-number | content |
|---|---|---|---|
| 1 | {{toc/}} | ||
| 2 | |||
| 3 | |||
| 4 | | |(% style="width:188px" %)**Пример задачи**|(% style="width:1268px" %)**Наборы решений** | ||
| 5 | |1|(% style="width:188px" %)1 родительский процесс и N дочерних процессов.|(% style="width:1268px" %)((( | ||
| 6 | |В данном примере имеется в виду, что дочерние процессы могут выполняться параллельно другу и независимо друг от друга, но в конце должны оповестить родительский процесс о необходимости продолжения обработки. | ||
| 7 | Если речь идет о каких-либо зависимостях порядка выполнения в дочерних процессах, то это может контролировать дочерний процесс (выделяя группу, которую сейчас можно запустить и ожидая окончания). | ||
| 8 | |((( | ||
| 9 | |((( | ||
| 10 | Вариант 1: CounterTrigger. | ||
| 11 | ))) | ||
| 12 | |((( | ||
| 13 | 1. Родительский процесс создает триггер со счетчиком N, создает и запускает дочерние процессы, засыпает. | ||
| 14 | 1. Дочерние процесс при завершении публикует TriggerEvent. | ||
| 15 | 1. TriggerConsumerRunner периодически считывает батч TriggerEvent, уменьшает считчик триггера и делает запись в БД. За счет агрегации событий завершения процессов мы уменьшаем нагрузку на БД. | ||
| 16 | 1. Когда все дочерние процессы отработали TriggerConsumerRunner получает значение счетчика 0 и взводит триггер. | ||
| 17 | 1. Триггер пробуждает родительский процесс для дальнейшего выполнения. | ||
| 18 | ))) | ||
| 19 | |TriggerEvent публикуются без использования TransactionOutbox напрямую в брокер после коммита транзакции (иначе мы бы нагружали БД). | ||
| 20 | |((( | ||
| 21 | Предполагаем, что основную часть времени система работает стабильно, но допускается ситуация, что транзакция закоммитилась, но TriggerEvent не смогли опубликоваться (остановка сервиса без graceful shutdown, проблемы соединения или работы с брокером сообщений). | ||
| 22 | |||
| 23 | Для таких случаев создается страхующий триггер (1 общий на тип процесса). Этот триггер запускается периодически и проходится по всем ожидающим процессам, проверяя условие (в реализации можно использовать keyset пагинацию) (в реализации можно использовать join для проверки условия). | ||
| 24 | Этот триггер выполняется периодически с более крупной временной задержкой. В случае обнаружения потери TriggerEvent, он поднимет заклинивший родительский процесс и он будет обработан (но позже). Можно установить этому триггеру низкий приоритет. | ||
| 25 | ))) | ||
| 26 | |[[image:Родительский дочерний процесс. Sequence.jpg]] | ||
| 27 | ))) | ||
| 28 | |((( | ||
| 29 | Вариант №2: | ||
| 30 | |||
| 31 | Мы просто ставим TimerTrigger на условно 1-5-10 минут (насколько важна задержка) и перепроверяем условие завершения. | ||
| 32 | В этом случае будет | ||
| 33 | |||
| 34 | * Из минус: что родительский процесс узнает о завершении дочерних процессов с задержкой. | ||
| 35 | Если дочерний процесс падает в ошибку, TimerTrigger все равно будет крутиться и создавать пустую нагрузку. | ||
| 36 | * Из плюсов: будет меньше пишущей нагрузки на БД чем в варианте 1 (но больше читающей - на проверку) т.к. у нас не будет CounterTrigger, но будет периодический запрос на проверку завершения всех дочерних процессов (аналогично страхующему триггер). | ||
| 37 | * [Расширенный]: Дочерние процессы в блоке wakeup condition проверяют наличие незавершенных процессов. | ||
| 38 | ** Если все процессы завершены или (незавершенных процессов мало и нет процессов с ошибкой), то можно опубликовать событие на TimerTrigger на сброс или установку минимальной задержки. | ||
| 39 | ))) | ||
| 40 | |((( | ||
| 41 | Вариант №3: | ||
| 42 | |||
| 43 | Дочерние процессы выполняются через родительский (ограничение в рамках одной ноды). | ||
| 44 | Точкой выполнения является родительский процесс, который внутри себя (параллельно или последовательно) выполняет дочерние процессы. | ||
| 45 | За счет такого способа у нас также отсутствует конкуренция передачи сигнала в родительский процесс. | ||
| 46 | Но мы ограничены выполнением дочерних процессов одной одной сервиса. | ||
| 47 | Сложнее контролировать распределение нагрузки, если будет вложенный параллелизм. | ||
| 48 | Также решает проблему, если дочерний процесс содержит ожидание (например асинхронный запрос-ответ), тут будет конкуренция сигнала от хендлера ответа к родительскому процессу. | ||
| 49 | ))) | ||
| 50 | |((( | ||
| 51 | Вариант №4: | ||
| 52 | |||
| 53 | SimpleStreamTrigger + Timer. | ||
| 54 | |||
| 55 | * Триггер проверяет условие завершения всех дочерних процессов (можно прикинуть количество незавершенных дочерних процессов). | ||
| 56 | ** Если все обработано, то пробуждает процесс и деактивируется. | ||
| 57 | ** Иначе: | ||
| 58 | *** деактивируется (до поступления хотя бы одного сигнала), | ||
| 59 | *** взводит признак стрима - процесс ожидает, | ||
| 60 | *** взводит флаг новых сигналов на 0, | ||
| 61 | *** выставляет задержку от оценки количества необработанных процессов (< N - малая задержка, иначе большая задержка). | ||
| 62 | * [Расширенный]: Дочерние процессы в блоке wakeup condition проверяют наличие незавершенных процессов. | ||
| 63 | ** Если все процессы завершены или (незавершенных процессов мало и нет процессов с ошибкой), то можно опубликовать событие на SimpleStreamTrigger на сброс или установку минимальной задержки (в дополнение к сигналу). | ||
| 64 | * Читающей нагрузки будет немного больше чем в варианте 2 (чтение триггера на поступлении сигнала), | ||
| 65 | но пишущей нагрузки будет меньше чем в варианте 1 (запись - только на активации новым сигналом). | ||
| 66 | * Если сигналов нет, то нет пустых срабатываний в отличие от варианта 2 (т.к. нет поступления сигнала от дочерних процессов). | ||
| 67 | ))) | ||
| 68 | |((( | ||
| 69 | |Вариант №5: | ||
| 70 | SimpleStreamTrigger + Счетчик в [[Redis>>doc:Разработка.Базы данных.NoSQL.Ключ-значение структура.Redis.WebHome]]. (на текущий момент самый лучший вариант). | ||
| 71 | |InMemory счетчик, нагрузка на БД и конкуренция. | ||
| 72 | Дочерний процесс уменьшает счетчик. И публикует событие только если счетчик достиг 0. | ||
| 73 | Совмещает преимущества из варианта 1.1 (при этом не нагружает БД), в случае ошибки переключается на режим 1.4. | ||
| 74 | |((( | ||
| 75 | Проблема: изменение счетчика не привязано к основной транзакции БД. | ||
| 76 | |||
| 77 | Возможно: | ||
| 78 | |||
| 79 | 1. В начале транзакции (или шага) атомарно проверяем MemberSet, если есть запись то удаляем и увеличиваем счетчик на 1 (означает что процесс уже уменьшал счетчик, но потом было падение). | ||
| 80 | 1. До коммита транзакции атомарно добавляем значение в MemberSet и уменьшаем счетчик на 1. Если счетчик равен 0, то публикуем TriggerEvent | ||
| 81 | (добавляем в ручную компенсацию вызов из пункта 1 (при откате изоляции шага и компенсации транзакции)). | ||
| 82 | 1. После коммита транзакции удаляем запись из MemberSet. | ||
| 83 | (Если мы падаем тут, то процесс уже перешел на другой шаг или даже завершился, поэтому наличие единичной остаточной записи в memberSet не будет критичным). | ||
| 84 | 1. Trigger получает событие. (Необязательно) в хендлере может првоерить значения счетчиков и MemberSet. | ||
| 85 | |||
| 86 | MemberSet используется для уменьшения вероятности увеличить или уменьшить счетчик дважды одним экземпляром дочернего процесса. | ||
| 87 | ))) | ||
| 88 | |((( | ||
| 89 | В случае обнаружения повреждения обработка фактически переходит в режим 1.4: начинает публиковать событий каждый раз и используется задержка. | ||
| 90 | |||
| 91 | Примеры проблемы: | ||
| 92 | |||
| 93 | * Падение InMemory хранилища. Предполагается режим без снимков и удаление ключей. | ||
| 94 | Обнаружение (со стороны дочернего процесса) через отсутствие ключей (проверяется в транзакции). | ||
| 95 | * Дублирование обновления счетчика. | ||
| 96 | Обнаружение (со стороны дочернего процесса) через значение счетчика < 0. | ||
| 97 | Обнаружение (со стороны триггера) через активацию триггера (поступления сигнала от процесса), при этом обнаруживается что не все процессы завершены. | ||
| 98 | ))) | ||
| 99 | ))) | ||
| 100 | ))) | ||
| 101 | |2|(% style="width:188px" %)Transaction outbox stream process.|(% style="width:1268px" %)[[image:TransactionOutbox. Sequence.jpg]] | ||
| 102 | |3|(% style="width:188px" %)Stream trigger|(% style="width:1268px" %)((( | ||
| 103 | | |((( | ||
| 104 | * Позволяет убрать лишние запросы пробуждения процесса (когда он и так запущен). | ||
| 105 | * __Позволяет полностью убрать задержку после остановки процесса__ (если есть новое сообщения, то он сразу же будет пробужден). | ||
| 106 | За счет того, что триггер точно знает, что есть новые сообщения и процесс только что уснул. | ||
| 107 | * Вводит 2 типа события, 1 сигнал о новом сообщении (содержит offset значение), 2 - процесс идет спать (содержит offset значение). | ||
| 108 | * Вводит дополнительное состояние в триггер: максимальный offset сообщения, максимальный offset обработанного процессом сообщения, флаг состояния сна процесса. | ||
| 109 | * В некоторых случаях позволяет не выполнять wakeup код в конце сессии обработки (если отключить wakeup, оставить только stream trigger) | ||
| 110 | (блокировка и обновление wakeup entity, проверка wakeup условия), __улучшает перформанс такта работы__. | ||
| 111 | ))) | ||
| 112 | |Алгоритм триггера.|((( | ||
| 113 | * При получении события о засыпании процесса: | ||
| 114 | Фиксирует смещение процесса обработки и сравнивает со смещением сообщения. | ||
| 115 | Если все сообщения обработаны, то не пробуждает процесс, иначе пробуждает процесс. | ||
| 116 | * При получении события о новом сообщении: | ||
| 117 | Фиксирует новое наибольшее смещение. | ||
| 118 | Если процесс не спит (по флагу в триггере), то ничего не делает. | ||
| 119 | Если процесс спит (по флагу), то пробуждает процесс. | ||
| 120 | |||
| 121 | Отслеживает смещение обработки процесса и последнего события. | ||
| 122 | Ожидает от процесса события о том, что он все обработал, его последнее смещение и он идет спать. | ||
| 123 | Если есть сообщения со смещением больше чем указал процесс, то делает гарантированное пробуждение процесса. | ||
| 124 | Когда поступает сигнал о новом сообщении (от отправителя сообщения), то обновляет данные о максимальном смещении и пробуждает процесс, если он спит | ||
| 125 | ))) | ||
| 126 | |Заготовка|[[https:~~/~~/github.com/cccc1808/cccc1808.ProcessEngine/tree/cccc1808/feature/trigger_stream_trigger>>https://github.com/cccc1808/cccc1808.ProcessEngine/tree/cccc1808/feature/trigger_stream_trigger]] | ||
| 127 | ))) | ||
| 128 | |4|(% style="width:188px" %)Групповое действие|(% style="width:1268px" %)((( | ||
| 129 | | |Действие, которое нужно применить к диапазону строк (сравнительно большому), независимо для каждой строки. | ||
| 130 | Наличие у строк упорядоченного столбца (для выделения диапазонов). | ||
| 131 | | |((( | ||
| 132 | |(% style="width:888px" %)Родительские процесс определяет границы диапазона [min, max].|(% style="width:266px" %){{code language="none"}}select min(), max() | ||
| 133 | where condition(){{/code}} | ||
| 134 | |(% style="width:888px" %)Родительский процесс нарезает диапазон [min, max] на поддиапазоны. На каждый поддиапазон создается дочерний процесс.|(% style="width:266px" %) | ||
| 135 | |(% style="width:888px" %)Каждый дочерний процесс обрабатывает свой поддиапазон строк (параллельно).|(% style="width:266px" %)Внутри поддиапазона может использоваться keyset пагинация. | ||
| 136 | |(% style="width:888px" %)Родительский процесс ожидает завершения дочерних процессов (см. пример 1).|(% style="width:266px" %) | ||
| 137 | ))) | ||
| 138 | ))) | ||
| 139 | |5|(% style="width:188px" %)Распределение заявок между исполнителями | ||
| 140 | (Заготовка).|(% style="width:1268px" %)((( | ||
| 141 | |(% style="width:94px" %)Описание|(% style="width:1156px" %)Есть поток заявок на деталь (создание детали требует ресурсов, 1 станок, время). | ||
| 142 | Есть N станков. Опционально: у станка есть уровень ресурсов и коэффициент скорости работы. | ||
| 143 | Реализация системы распределения и обработки. | ||
| 144 | |(% style="width:94px" %)Вариант 1|(% style="width:1156px" %)((( | ||
| 145 | | |Планирование без очереди к станку. | ||
| 146 | | |((( | ||
| 147 | * У процесса планировщика есть | ||
| 148 | ** StreamTrigger на поток заявок. | ||
| 149 | (Можно использовать расширение SignalCode, чтобы временно игнорировать откладывать этот сигнал пока все слоты станков заняты). | ||
| 150 | ** StreanTrigger на поток сигналов об освобождении слота станка. | ||
| 151 | ))) | ||
| 152 | | |Процесс планировщик назначает заявку на свободный станок. | ||
| 153 | Вопрос наиболее эффективной функции выбора (оценка наибольшее количество ресурсов, наилучшая скорость обработки и др.). | ||
| 154 | Когда все станки заняты планировщик ожидает освобождения станков. | ||
| 155 | ))) | ||
| 156 | |(% style="width:94px" %)Вариант 2|(% style="width:1156px" %)((( | ||
| 157 | | |Планировщик с очередью к станку. | ||
| 158 | | |((( | ||
| 159 | * У процесса станка есть StreamTrigger, на который планировщик подает сигнал в добавления задачи в его очередь. | ||
| 160 | * У процесса планировщика есть StreamTrigger на поток заявок. | ||
| 161 | ))) | ||
| 162 | | |При поступлении заявки процесс планировщик сразу назначает в очередь на какой либо станок. | ||
| 163 | Вопрос наиболее эффективной функции выбора (оценка размера очереди, достаточности у станка ресурсов для ее обработки, общего количества ресурсов, скорости работы станка и др.). | ||
| 164 | ))) | ||
| 165 | ))) | ||
| 166 | |||
| 167 | ---- | ||
| 168 | |||
| 169 | ==== Внутренние ссылки: ==== | ||
| 170 | |||
| 171 | ====== Дочерние страницы: ====== | ||
| 172 | |||
| 173 | {{children/}} | ||
| 174 | |||
| 175 | ====== Обратные ссылки: ====== | ||
| 176 | |||
| 177 | {{velocity}} | ||
| 178 | #set ($links = $doc.getBacklinks()) | ||
| 179 | #if ($links.size() > 0) | ||
| 180 | #foreach ($docname in $links) | ||
| 181 | #set ($rdoc = $xwiki.getDocument($docname).getTranslatedDocument()) | ||
| 182 | * [[$escapetool.xml($rdoc.fullName)]] | ||
| 183 | #end | ||
| 184 | #else | ||
| 185 | No back links for this page! | ||
| 186 | #end | ||
| 187 | {{/velocity}} | ||
| 188 | |||
| 189 | ---- |