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