Изменения документа Примеры
Редактировал(а) Alexandr Fokin 2026/08/26 19:05
От версии 8.6
отредактировано Alexandr Fokin
на 2026/04/29 11:34
на 2026/04/29 11:34
Изменить комментарий:
К данной версии нет комментариев
К версии 8.45
отредактировано Alexandr Fokin
на 2026/08/18 12:22
на 2026/08/18 12:22
Изменить комментарий:
К данной версии нет комментариев
Сводка
-
Свойства страницы (2 изменено, 0 добавлено, 0 удалено)
-
Объекты (0 изменено, 1 добавлено, 0 удалено)
Подробности
- Свойства страницы
-
- Родительский документ
-
... ... @@ -1,1 +1,1 @@ 1 -Проекты и репозитории.Библиотеки.Движок cccc1808\. ProcessEngine.WebHome 1 +Проекты и репозитории.Библиотеки.Движок cccc1808\. ProcessEngine.Тема\. Сигналы и триггеры.WebHome - Содержимое
-
... ... @@ -1,3 +1,7 @@ 1 +{{toc/}} 2 + 3 + 4 +| |(% style="width:188px" %)**Пример задачи**|(% style="width:1268px" %)**Наборы решений** 1 1 |1|(% style="width:188px" %)1 родительский процесс и N дочерних процессов.|(% style="width:1268px" %)((( 2 2 |В данном примере имеется в виду, что дочерние процессы могут выполняться параллельно другу и независимо друг от друга, но в конце должны оповестить родительский процесс о необходимости продолжения обработки. 3 3 Если речь идет о каких-либо зависимостях порядка выполнения в дочерних процессах, то это может контролировать дочерний процесс (выделяя группу, которую сейчас можно запустить и ожидая окончания). ... ... @@ -22,23 +22,88 @@ 22 22 |[[image:Родительский дочерний процесс. Sequence.jpg]] 23 23 ))) 24 24 |((( 25 -В озможен вариант №2:29 +Вариант №2: 26 26 27 -Мы просто стави тtimerTrigger на условно 1-5-10 минут (насколько важна задержка) и перепроверяем условие завершения.31 +Мы просто ставим TimerTrigger на условно 1-5-10 минут (насколько важна задержка) и перепроверяем условие завершения. 28 28 В этом случае будет 29 29 30 -* Из минус: что родительский процесс узнает о завершении дочерних процессов с задержкой (хотя в задержке можно использовать функцию от количества необработанных дочерних процессов, но тогда нужно считать количество или хотя бы что оно не больше N). 31 -* Из плюсов: будет меньше пишущей нагрузки на БД (но больше читающей - на проверку) т.к. у нас не будет CounterTrigger, но будет периодический запрос на проверку завершения всех дочерних процессов (аналогично страхующему триггер). \ 34 +* Из минус: что родительский процесс узнает о завершении дочерних процессов с задержкой. 35 +Если дочерний процесс падает в ошибку, TimerTrigger все равно будет крутиться и создавать пустую нагрузку. 36 +* Из плюсов: будет меньше пишущей нагрузки на БД чем в варианте 1 (но больше читающей - на проверку) т.к. у нас не будет CounterTrigger, но будет периодический запрос на проверку завершения всех дочерних процессов (аналогично страхующему триггер). 37 +* [Расширенный]: Дочерние процессы в блоке wakeup condition проверяют наличие незавершенных процессов. 38 +** Если все процессы завершены или (незавершенных процессов мало и нет процессов с ошибкой), то можно опубликовать событие на TimerTrigger на сброс или установку минимальной задержки. 32 32 ))) 40 +|((( 41 +Вариант №3: 42 + 43 +Дочерние процессы выполняются через родительский (ограничение в рамках одной ноды). 44 +Точкой выполнения является родительский процесс, который внутри себя (параллельно или последовательно) выполняет дочерние процессы. 45 +За счет такого способа у нас также отсутствует конкуренция передачи сигнала в родительский процесс. 46 +Но мы ограничены выполнением дочерних процессов одной одной сервиса. 47 +Сложнее контролировать распределение нагрузки, если будет вложенный параллелизм. 48 +Также решает проблему, если дочерний процесс содержит ожидание (например асинхронный запрос-ответ), тут будет конкуренция сигнала от хендлера ответа к родительскому процессу. 33 33 ))) 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 +))) 34 34 |2|(% style="width:188px" %)Transaction outbox stream process.|(% style="width:1268px" %)[[image:TransactionOutbox. Sequence.jpg]] 35 35 |3|(% style="width:188px" %)Stream trigger|(% style="width:1268px" %)((( 36 36 | |((( 37 37 * Позволяет убрать лишние запросы пробуждения процесса (когда он и так запущен). 38 -* Позволяет полностью убрать задержку после остановки процесса (если есть новое сообщения, то он сразу же будет пробужден). 105 +* __Позволяет полностью убрать задержку после остановки процесса__ (если есть новое сообщения, то он сразу же будет пробужден). 39 39 За счет того, что триггер точно знает, что есть новые сообщения и процесс только что уснул. 40 40 * Вводит 2 типа события, 1 сигнал о новом сообщении (содержит offset значение), 2 - процесс идет спать (содержит offset значение). 41 41 * Вводит дополнительное состояние в триггер: максимальный offset сообщения, максимальный offset обработанного процессом сообщения, флаг состояния сна процесса. 109 +* В некоторых случаях позволяет не выполнять wakeup код в конце сессии обработки (если отключить wakeup, оставить только stream trigger) 110 +(блокировка и обновление wakeup entity, проверка wakeup условия), __улучшает перформанс такта работы__. 42 42 ))) 43 43 |Алгоритм триггера.|((( 44 44 * При получении события о засыпании процесса: ... ... @@ -54,10 +54,10 @@ 54 54 Если есть сообщения со смещением больше чем указал процесс, то делает гарантированное пробуждение процесса. 55 55 Когда поступает сигнал о новом сообщении (от отправителя сообщения), то обновляет данные о максимальном смещении и пробуждает процесс, если он спит 56 56 ))) 57 -| |TODO: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]] 58 58 ))) 59 59 |4|(% style="width:188px" %)Групповое действие|(% style="width:1268px" %)((( 60 -| |Действие, которое нужно применить к диапазону строк, независимо для каждой строки. 129 +| |Действие, которое нужно применить к диапазону строк (сравнительно большому), независимо для каждой строки. 61 61 Наличие у строк упорядоченного столбца (для выделения диапазонов). 62 62 | |((( 63 63 |(% style="width:888px" %)Родительские процесс определяет границы диапазона [min, max].|(% style="width:266px" %){{code language="none"}}select min(), max() ... ... @@ -64,6 +64,57 @@ 64 64 where condition(){{/code}} 65 65 |(% style="width:888px" %)Родительский процесс нарезает диапазон [min, max] на поддиапазоны. На каждый поддиапазон создается дочерний процесс.|(% style="width:266px" %) 66 66 |(% style="width:888px" %)Каждый дочерний процесс обрабатывает свой поддиапазон строк (параллельно).|(% style="width:266px" %)Внутри поддиапазона может использоваться keyset пагинация. 67 -|(% style="width:888px" %)Родительский процесс ожидает завершения дочерних процессов.|(% style="width:266px" %) 136 +|(% style="width:888px" %)Родительский процесс ожидает завершения дочерних процессов (см. пример 1).|(% style="width:266px" %) 68 68 ))) 69 69 ))) 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 +----
- XWiki.XWikiComments[0]
-
- Автор
-
... ... @@ -1,0 +1,1 @@ 1 +XWiki.cccc1808 - Комментарий
-
... ... @@ -1,0 +1,4 @@ 1 +Замечание: конфигурация задержки trigger consumer вычитывания и накопления батча trigger events. 2 + 3 +* Для примера 1 предпочтительная более большая задержка т.к. это уменьшит нагрузку на БД (агрегирует больше сигналов от дочерних процессов в одну операцию обновления). Throughput. 4 +* Для примера 3 в контексте inbox stream trigger, может быть предпочтительная более низкая задержка, чтобы не раздувать задержку от поступления сообщения до его обработки. Latency. - Дата
-
... ... @@ -1,0 +1,1 @@ 1 +2026-05-01 15:36:23.922