Исходный код вики Примеры

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

Последние авторы
1 | |(% style="width:188px" %)**Пример задачи**|(% style="width:1268px" %)**Наборы решений**
2 |1|(% style="width:188px" %)(((
3 (% class="wikigeneratedid" id="H144043E43443844243543B44C44143A438439A043F44043E446435441441438N43443E44743544043D43844543F44043E44643544144143E432." %)
4 1 родительский процесс и N дочерних процессов.
5 )))|(% style="width:1268px" %)(((
6 |В данном примере имеется в виду, что дочерние процессы могут выполняться параллельно другу и независимо друг от друга, но в конце должны оповестить родительский процесс о необходимости продолжения обработки.
7 Если речь идет о каких-либо зависимостях порядка выполнения в дочерних процессах, то это может контролировать дочерний процесс (выделяя группу, которую сейчас можно запустить и ожидая окончания).
8 |(((
9 |(((
10 (% class="wikigeneratedid" id="H41243044043843043D4421:CounterTrigger." %)
11 Вариант 1: CounterTrigger.
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 )))
27 |[[image:Родительский дочерний процесс. Sequence.jpg]]
28 )))
29 |(((
30 Вариант №2:
31
32 Мы просто ставим TimerTrigger на условно 1-5-10 минут (насколько важна задержка) и перепроверяем условие завершения.
33 В этом случае будет
34
35 * Из минус: что родительский процесс узнает о завершении дочерних процессов с задержкой.
36 Если дочерний процесс падает в ошибку, TimerTrigger все равно будет крутиться и создавать пустую нагрузку.
37 * Из плюсов: будет меньше пишущей нагрузки на БД чем в варианте 1 (но больше читающей - на проверку) т.к. у нас не будет CounterTrigger, но будет периодический запрос на проверку завершения всех дочерних процессов (аналогично страхующему триггер).
38 * [Расширенный]: Дочерние процессы в блоке wakeup condition проверяют наличие незавершенных процессов.
39 ** Если все процессы завершены или (незавершенных процессов мало и нет процессов с ошибкой), то можно опубликовать событие на TimerTrigger на сброс или установку минимальной задержки.
40 )))
41 |(((
42 Вариант №3:
43
44 Дочерние процессы выполняются через родительский (ограничение в рамках одной ноды).
45 Точкой выполнения является родительский процесс, который внутри себя (параллельно или последовательно) выполняет дочерние процессы.
46 За счет такого способа у нас также отсутствует конкуренция передачи сигнала в родительский процесс.
47 Но мы ограничены выполнением дочерних процессов одной одной сервиса.
48 Сложнее контролировать распределение нагрузки, если будет вложенный параллелизм.
49 Также решает проблему, если дочерний процесс содержит ожидание (например асинхронный запрос-ответ), тут будет конкуренция сигнала от хендлера ответа к родительскому процессу.
50 )))
51 |(((
52 Вариант №4:
53
54 SimpleStreamTrigger + Timer.
55
56 * Триггер проверяет условие завершения всех дочерних процессов (можно прикинуть количество незавершенных дочерних процессов).
57 ** Если все обработано, то пробуждает процесс и деактивируется.
58 ** Иначе:
59 *** деактивируется (до поступления хотя бы одного сигнала),
60 *** взводит признак стрима - процесс ожидает,
61 *** взводит флаг новых сигналов на 0,
62 *** выставляет задержку от оценки количества необработанных процессов (< N - малая задержка, иначе большая задержка).
63 * [Расширенный]: Дочерние процессы в блоке wakeup condition проверяют наличие незавершенных процессов.
64 ** Если все процессы завершены или (незавершенных процессов мало и нет процессов с ошибкой), то можно опубликовать событие на SimpleStreamTrigger на сброс или установку минимальной задержки (в дополнение к сигналу).
65 * Читающей нагрузки будет немного больше чем в варианте 2 (чтение триггера на поступлении сигнала),
66 но пишущей нагрузки будет меньше чем в варианте 1 (запись - только на активации новым сигналом).
67 * Если сигналов нет, то нет пустых срабатываний в отличие от варианта 2 (т.к. нет поступления сигнала от дочерних процессов).
68 )))
69 |(((
70 |Вариант №5:
71 SimpleStreamTrigger + Счетчик (ExternalCounter) в [[Redis>>doc:Разработка.Базы данных.NoSQL.Ключ-значение структура.Redis.WebHome]]. (__на текущий момент самый лучший вариант__).
72 |InMemory счетчик, нагрузка на БД и конкуренция.
73 Дочерний процесс уменьшает счетчик. И публикует событие только если счетчик достиг 0.
74 Совмещает преимущества из варианта 1.1 (при этом не нагружает БД), в случае ошибки переключается на режим 1.4.
75 |(((
76 Проблема: изменение счетчика не привязано к основной транзакции БД.
77
78 Возможно:
79
80 1. В начале транзакции (или шага) атомарно проверяем MemberSet, если есть запись то удаляем и увеличиваем счетчик на 1 (означает что процесс уже уменьшал счетчик, но потом было падение).
81 1. До коммита транзакции атомарно добавляем значение в MemberSet и уменьшаем счетчик на 1. Если счетчик равен 0, то публикуем TriggerEvent
82 (добавляем в ручную компенсацию вызов из пункта 1 (при откате изоляции шага и компенсации транзакции)).
83 1. После коммита транзакции удаляем запись из MemberSet.
84 (Если мы падаем тут, то процесс уже перешел на другой шаг или даже завершился, поэтому наличие единичной остаточной записи в memberSet не будет критичным).
85 1. Trigger получает событие. (Необязательно) в хендлере может првоерить значения счетчиков и MemberSet.
86
87 MemberSet используется для уменьшения вероятности увеличить или уменьшить счетчик дважды одним экземпляром дочернего процесса.
88 )))
89 |(((
90 В случае обнаружения повреждения обработка фактически переходит в режим 1.4: начинает публиковать событий каждый раз и используется задержка.
91
92 Примеры проблемы:
93
94 * Падение InMemory хранилища. Предполагается режим без снимков и удаление ключей.
95 Обнаружение (со стороны дочернего процесса) через отсутствие ключей (проверяется в транзакции).
96 * Дублирование обновления счетчика.
97 Обнаружение (со стороны дочернего процесса) через значение счетчика < 0.
98 Обнаружение (со стороны триггера) через активацию триггера (поступления сигнала от процесса), при этом обнаруживается что не все процессы завершены.
99 )))
100 )))
101 )))
102 |2|(% style="width:188px" %)Transaction outbox stream process.|(% style="width:1268px" %)Смотри stream trigger.
103 |-|(% style="width:188px" %) |(% style="width:1268px" %)
104 |4|(% style="width:188px" %)Групповое действие|(% style="width:1268px" %)(((
105 | |Действие, которое нужно применить к диапазону строк (сравнительно большому), независимо для каждой строки.
106 Наличие у строк упорядоченного столбца (для выделения диапазонов).
107 | |(((
108 |(% style="width:888px" %)Родительские процесс определяет границы диапазона [min, max].|(% style="width:266px" %){{code language="none"}}select min(), max()
109 where condition(){{/code}}
110 |(% style="width:888px" %)Родительский процесс нарезает диапазон [min, max] на поддиапазоны. На каждый поддиапазон создается дочерний процесс.|(% style="width:266px" %)
111 |(% style="width:888px" %)Каждый дочерний процесс обрабатывает свой поддиапазон строк (параллельно).|(% style="width:266px" %)Внутри поддиапазона может использоваться keyset пагинация.
112 |(% style="width:888px" %)Родительский процесс ожидает завершения дочерних процессов (см. пример 1).|(% style="width:266px" %)
113 )))
114 )))
115 |5|(% style="width:188px" %)Распределение заявок между исполнителями
116 (Заготовка).|(% style="width:1268px" %)(((
117 |(% style="width:94px" %)Описание|(% style="width:1156px" %)Есть поток заявок на деталь (создание детали требует ресурсов, 1 станок, время).
118 Есть N станков. Опционально: у станка есть уровень ресурсов и коэффициент скорости работы.
119 Реализация системы распределения и обработки.
120 |(% style="width:94px" %)Вариант 1|(% style="width:1156px" %)(((
121 | |Планирование без очереди к станку.
122 | |(((
123 * У процесса планировщика есть
124 ** StreamTrigger на поток заявок.
125 (Можно использовать расширение SignalCode, чтобы временно игнорировать откладывать этот сигнал пока все слоты станков заняты).
126 ** StreanTrigger на поток сигналов об освобождении слота станка.
127 )))
128 | |Процесс планировщик назначает заявку на свободный станок.
129 Вопрос наиболее эффективной функции выбора (оценка наибольшее количество ресурсов, наилучшая скорость обработки и др.).
130 Когда все станки заняты планировщик ожидает освобождения станков.
131 )))
132 |(% style="width:94px" %)Вариант 2|(% style="width:1156px" %)(((
133 | |Планировщик с очередью к станку.
134 | |(((
135 * У процесса станка есть StreamTrigger, на который планировщик подает сигнал в добавления задачи в его очередь.
136 * У процесса планировщика есть StreamTrigger на поток заявок.
137 )))
138 | |При поступлении заявки процесс планировщик сразу назначает в очередь на какой либо станок.
139 Хранение фактического уровня ресурсов и зарезервированного уровня ресурсов (с учетом очереди к станку).
140 Вопрос наиболее эффективной функции выбора (оценка размера очереди, достаточности у станка ресурсов для ее обработки, общего количества ресурсов, скорости работы станка и др.).
141 )))
142 )))
143
144 ----
145
146 ==== Внутренние ссылки: ====
147
148 ====== Дочерние страницы: ======
149
150 {{children/}}
151
152 ====== Обратные ссылки: ======
153
154 {{velocity}}
155 #set ($links = $doc.getBacklinks())
156 #if ($links.size() > 0)
157 #foreach ($docname in $links)
158 #set ($rdoc = $xwiki.getDocument($docname).getTranslatedDocument())
159 * [[$escapetool.xml($rdoc.fullName)]]
160 #end
161 #else
162 No back links for this page!
163 #end
164 {{/velocity}}
165
166 ----