Изменения документа Примеры

Редактировал(а) Alexandr Fokin 2026/08/26 19:05

От версии 8.4
отредактировано Alexandr Fokin
на 2026/04/29 11:32
Изменить комментарий: К данной версии нет комментариев
К версии 8.40
отредактировано Alexandr Fokin
на 2026/07/27 21:00
Изменить комментарий: К данной версии нет комментариев

Сводка

Подробности

Свойства страницы
Родительский документ
... ... @@ -1,1 +1,1 @@
1 -Проекты и репозитории.Библиотеки.Движок cccc1808\. ProcessEngine.WebHome
1 +Проекты и репозитории.Библиотеки.Движок cccc1808\. ProcessEngine.Тема\. Сигналы и триггеры.WebHome
Содержимое
... ... @@ -1,4 +1,4 @@
1 -|1|Родительский процесс, N дочерних процессов.|(((
1 +|1|(% style="width:188px" %)1 родительский процесс и N дочерних процессов.|(% style="width:1268px" %)(((
2 2  |В данном примере имеется в виду, что дочерние процессы могут выполняться параллельно другу и независимо друг от друга, но в конце должны оповестить родительский процесс о необходимости продолжения обработки.
3 3  Если речь идет о каких-либо зависимостях порядка выполнения в дочерних процессах, то это может контролировать дочерний процесс (выделяя группу, которую сейчас можно запустить и ожидая окончания).
4 4  |(((
... ... @@ -22,23 +22,87 @@
22 22  |[[image:Родительский дочерний процесс. Sequence.jpg]]
23 23  )))
24 24  |(((
25 озможен вариант №2:
25 +Вариант №2:
26 26  
27 -Мы просто ставит timerTrigger на условно 1-5-10 минут (насколько важна задержка) и перепроверяем условие завершения.
27 +Мы просто ставим TimerTrigger на условно 1-5-10 минут (насколько важна задержка) и перепроверяем условие завершения.
28 28  В этом случае будет
29 29  
30 -* Из минус: что родительский процесс узнает о завершении дочерних процессов с задержкой (хотя в задержке можно использовать функцию от количества необработанных дочерних процессов, но тогда нужно считать количество или хотя бы что оно не больше N).
31 -* Из плюсов: будет меньше пишущей нагрузки на БД (но больше читающей - на проверку) т.к. у нас не будет CounterTrigger, но будет периодический запрос на проверку завершения всех дочерних процессов (аналогично страхующему триггер). \
30 +* Из минус: что родительский процесс узнает о завершении дочерних процессов с задержкой.
31 +Если дочерний процесс падает в ошибку, TimerTrigger все равно будет крутиться и создавать пустую нагрузку.
32 +* Из плюсов: будет меньше пишущей нагрузки на БД чем в варианте 1 (но больше читающей - на проверку) т.к. у нас не будет CounterTrigger, но будет периодический запрос на проверку завершения всех дочерних процессов (аналогично страхующему триггер).
33 +* [Расширенный]: Дочерние процессы в блоке wakeup condition проверяют наличие незавершенных процессов.
34 +** Если все процессы завершены или (незавершенных процессов мало и нет процессов с ошибкой), то можно опубликовать событие на TimerTrigger на сброс или установку минимальной задержки.
32 32  )))
36 +|(((
37 +Вариант №3:
38 +
39 +Дочерние процессы выполняются через родительский (ограничение в рамках одной ноды).
40 +Точкой выполнения является родительский процесс, который внутри себя (параллельно или последовательно) выполняет дочерние процессы.
41 +За счет такого способа у нас также отсутствует конкуренция передачи сигнала в родительский процесс.
42 +Но мы ограничены выполнением дочерних процессов одной одной сервиса.
43 +Сложнее контролировать распределение нагрузки, если будет вложенный параллелизм.
44 +Также решает проблему, если дочерний процесс содержит ожидание (например асинхронный запрос-ответ), тут будет конкуренция сигнала от хендлера ответа к родительскому процессу.
33 33  )))
34 -|2|Transaction outbox stream process.|[[image:TransactionOutbox. Sequence.jpg]]
35 -|3|Stream trigger|(((
46 +|(((
47 +Вариант №4:
48 +
49 +SimpleStreamTrigger + Timer.
50 +
51 +* Триггер проверяет условие завершения всех дочерних процессов (можно прикинуть количество незавершенных дочерних процессов).
52 +** Если все обработано, то пробуждает процесс и деактивируется.
53 +** Иначе:
54 +*** деактивируется (до поступления хотя бы одного сигнала),
55 +*** взводит признак стрима - процесс ожидает,
56 +*** взводит флаг новых сигналов на 0,
57 +*** выставляет задержку от оценки количества необработанных процессов (< N - малая задержка, иначе большая задержка).
58 +* [Расширенный]: Дочерние процессы в блоке wakeup condition проверяют наличие незавершенных процессов.
59 +** Если все процессы завершены или (незавершенных процессов мало и нет процессов с ошибкой), то можно опубликовать событие на SimpleStreamTrigger на сброс или установку минимальной задержки (в дополнение к сигналу).
60 +* Читающей нагрузки будет немного больше чем в варианте 2 (чтение триггера на поступлении сигнала),
61 +но пишущей нагрузки будет меньше чем в варианте 1 (запись - только на активации новым сигналом).
62 +* Если сигналов нет, то нет пустых срабатываний в отличие от варианта 2 (т.к. нет поступления сигнала от дочерних процессов).
63 +)))
64 +|(((
65 +|Вариант №5: SimpleStreamTrigger + Счетчик в [[Redis>>doc:Разработка.Базы данных.NoSQL.Ключ-значение структура.Redis.WebHome]].
66 +|InMemory счетчик, нагрузка на БД и конкуренция.
67 +Дочерний процесс уменьшает счетчик. И публикует событие только если счетчик достиг 0.
68 +Совмещает преимущества из варианта 1.1 (при этом не нагружает БД), в случае ошибки переключается на режим 1.4.
69 +|(((
70 +Проблема: изменение счетчика не привязано к основной транзакции БД.
71 +
72 +Возможно:
73 +
74 +1. В начале транзакции (или шага) атомарно проверяем MemberSet, если есть запись то удаляем и увеличиваем счетчик на 1 (означает что процесс уже уменьшал счетчик, но потом было падение).
75 +1. До коммита транзакции атомарно добавляем значение в MemberSet и уменьшаем счетчик на 1. Если счетчик равен 0, то публикуем TriggerEvent
76 +(добавляем в ручную компенсацию вызов из пункта 1 (при откате изоляции шага и компенсации транзакции)).
77 +1. После коммита транзакции удаляем запись из MemberSet.
78 +(Если мы падаем тут, то процесс уже перешел на другой шаг или даже завершился, поэтому наличие единичной остаточной записи в memberSet не будет критичным).
79 +1. Trigger получает событие. (Необязательно) в хендлере может првоерить значения счетчиков и MemberSet.
80 +
81 +MemberSet используется для уменьшения вероятности увеличить или уменьшить счетчик дважды одним экземпляром дочернего процесса.
82 +)))
83 +|(((
84 +В случае обнаружения повреждения обработка фактически переходит в режим 1.4: начинает публиковать событий каждый раз и используется задержка.
85 +
86 +Примеры проблемы:
87 +
88 +* Падение InMemory хранилища. Предполагается режим без снимков и удаление ключей.
89 +Обнаружение (со стороны дочернего процесса) через отсутствие ключей (проверяется в транзакции).
90 +* Дублирование обновления счетчика.
91 +Обнаружение (со стороны дочернего процесса) через значение счетчика < 0.
92 +Обнаружение (со стороны триггера) через активацию триггера (поступления сигнала от процесса), при этом обнаруживается что не все процессы завершены.
93 +)))
94 +)))
95 +)))
96 +|2|(% style="width:188px" %)Transaction outbox stream process.|(% style="width:1268px" %)[[image:TransactionOutbox. Sequence.jpg]]
97 +|3|(% style="width:188px" %)Stream trigger|(% style="width:1268px" %)(((
36 36  | |(((
37 37  * Позволяет убрать лишние запросы пробуждения процесса (когда он и так запущен).
38 -* Позволяет полностью убрать задержку после остановки процесса (если есть новое сообщения, то он сразу же будет пробужден).
100 +* __Позволяет полностью убрать задержку после остановки процесса__ (если есть новое сообщения, то он сразу же будет пробужден).
39 39  За счет того, что триггер точно знает, что есть новые сообщения и процесс только что уснул.
40 40  * Вводит 2 типа события, 1 сигнал о новом сообщении (содержит offset значение), 2 - процесс идет спать (содержит offset значение).
41 41  * Вводит дополнительное состояние в триггер: максимальный offset сообщения, максимальный offset обработанного процессом сообщения, флаг состояния сна процесса.
104 +* В некоторых случаях позволяет не выполнять wakeup код в конце сессии обработки (если отключить wakeup, оставить только stream trigger)
105 +(блокировка и обновление wakeup entity, проверка wakeup условия), __улучшает перформанс такта работы__.
42 42  )))
43 43  |Алгоритм триггера.|(((
44 44  * При получении события о засыпании процесса:
... ... @@ -54,17 +54,40 @@
54 54  Если есть сообщения со смещением больше чем указал процесс, то делает гарантированное пробуждение процесса.
55 55  Когда поступает сигнал о новом сообщении (от отправителя сообщения), то обновляет данные о максимальном смещении и пробуждает процесс, если он спит
56 56  )))
57 -| |TODO:
121 +|Заготовка|[[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 -|4|Групповое действие|(((
60 -| |Действие, которое нужно применить к диапазону строк, независимо для каждой строки.
61 -Наличие у строк упорядоченного столбца.
123 +|4|(% style="width:188px" %)Групповое действие|(% style="width:1268px" %)(((
124 +| |Действие, которое нужно применить к диапазону строк (сравнительно большому), независимо для каждой строки.
125 +Наличие у строк упорядоченного столбца (для выделения диапазонов).
62 62  | |(((
63 63  |(% style="width:888px" %)Родительские процесс определяет границы диапазона [min, max].|(% style="width:266px" %){{code language="none"}}select min(), max()
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" %)
131 +|(% style="width:888px" %)Родительский процесс ожидает завершения дочерних процессов (см. пример 1).|(% style="width:266px" %)
68 68  )))
69 -| |
70 70  )))
134 +
135 +----
136 +
137 +==== Внутренние ссылки: ====
138 +
139 +====== Дочерние страницы: ======
140 +
141 +{{children/}}
142 +
143 +====== Обратные ссылки: ======
144 +
145 +{{velocity}}
146 +#set ($links = $doc.getBacklinks())
147 +#if ($links.size() > 0)
148 + #foreach ($docname in $links)
149 + #set ($rdoc = $xwiki.getDocument($docname).getTranslatedDocument())
150 + * [[$escapetool.xml($rdoc.fullName)]]
151 + #end
152 +#else
153 + No back links for this page!
154 +#end
155 +{{/velocity}}
156 +
157 +----
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