Update README.md
This commit is contained in:
1 parent
fb017a9351
commit
bb6e8afbcf
1 file changed
+132
-96
@@ -1446,136 +1446,172 @@ futriix:~> cluster health
|
|||||||
|
|
||||||
## Backpressure
|
## Backpressure
|
||||||
|
|
||||||
Для равномерной загрузки каждого узла кластера, в субд futriix применяется механизм **"backpressure"**.
|
Назначение
|
||||||
**Backpressure (Обратное давление)** — это механизм управления потоками данных, предотвращающий переполнение узла кластера, когда он не успевает обрабатывать поступающие события. <br>
|
|
||||||
В нашем проекте данный механизм реализован для защиты системы от резких пиковых нагрузок: если буфер входящих сообщений достигает заданного порога, источник данных автоматически замедляется или приостанавливается до тех пор, пока потребитель не освободит ресурсы. Это гарантирует стабильность работы, отсутствие потерь данных и отказ от бесконечного накопления задач в очереди.
|
|
||||||
|
|
||||||
**Математическая формула вероятностного отклонения в Backpressure**
|
Для защищиты субд futriix от перегрузки при всплесках входящей нагрузки, в ней реализован механизм **Backpressure (клапан обратного давления)**, реализующий алгоритм **"Buffer Ring (кольцевой буфер обратного давления)."**
|
||||||
|
Он реализует адаптивное откладывание запросов вместо их немедленного отклонения: когда система близка к насыщению, запросы не отбрасываются, а ненадолго помещаются в кольцевой буфер, чтобы дать ядру время обработать уже принятые операции.
|
||||||
|
|
||||||
В системе Backpressure используется адаптивная вероятностная модель для отклонения запросов при перегрузке, которое рассчитывается по следующей формуле:
|
Такой подход отличается от классического «reject при перегрузке» тем, что сглаживает пики, а не режет их, и сохраняет больше полезной работы при кратковременных всплесках.
|
||||||
|
|
||||||
|
**Идея**
|
||||||
|
|
||||||
|
**Кольцевой буфер (Buffer Ring)** — это фиксированное кольцо слотов для ожидающих запросов. В отличие от очереди (FIFO), кольцо:
|
||||||
|
|
||||||
|
* Имеет постоянный размер N (не растёт под нагрузкой);
|
||||||
|
* Переиспользует слоты по кругу через индексы head и tail;
|
||||||
|
* При переполнении не расширяется, а применяет политику вытеснения;
|
||||||
|
* Поддерживает O(1) на вставку и извлечение — никаких аллокаций в горячем пути.
|
||||||
|
|
||||||
|
Каждый слот хранит «билет» запроса: метаданные (тип, deadline, приоритет) и функцию продолжения resume(), которая будет вызвана, когда наступит очередь.
|
||||||
|
|
||||||
```sh
|
```sh
|
||||||
P_reject = f(level) × g(load) × h(time)
|
tail (writer) head (reader)
|
||||||
|
↓ ↓
|
||||||
Где:
|
┌───┬───┬───┬───┬───┬───┬───┬───┐
|
||||||
|
│ 7 │ 8 │ 9 │ │ │ 4 │ 5 │ 6 │
|
||||||
* P_reject — итоговая вероятность отклонения запроса (0.0 - 1.0)
|
└───┴───┴───┴───┴───┴───┴───┴───┘
|
||||||
* f(level) — коэффициент на основе уровня перегрузки
|
↑ занято ↑ свободно ↑ занято
|
||||||
* g(load) — коэффициент на основе текущей нагрузки
|
|
||||||
* h(time) — коэффициент на основе времени (для защиты от "thundering herd")
|
|
||||||
```
|
```
|
||||||
|
|
||||||
**Компоненты формулы** </br>
|
**Алгоритм**
|
||||||
|
|
||||||
**Коэффициент уровня перегрузки f(level)**
|
1. Классификация нагрузки
|
||||||
|
|
||||||
|
На каждом цикле проверки (например, раз в check_interval_ms) менеджер снимает метрики: загрузку CPU, памяти, длину очереди, число соединений. По порогам из конфигурации определяется уровень:
|
||||||
|
|
||||||
|
| Уровень | Условие | Действие |
|
||||||
|
|---------|---------|----------|
|
||||||
|
| `Low` | CPU < `cpu_threshold` и очередь < `queue_size_threshold` | запрос проходит сразу |
|
||||||
|
| `Medium` | один из порогов превышен | запрос помещается в кольцо на короткую задержку `low_delay_ms` |
|
||||||
|
| `High` | CPU > `cpu_threshold` и очередь > `queue_size_threshold` | запрос помещается в кольцо; при переполнении применяется `high_reject_prob` |
|
||||||
|
| `Critical` | система не успевает дренировать кольцо | запросы отклоняются с `503 Service Unavailable` |
|
||||||
|
|
||||||
|
|
||||||
|
**2. Постановка в кольцо (enqueue)**
|
||||||
|
|
||||||
```sh
|
```sh
|
||||||
f(level) = {
|
func (r *Ring) Enqueue(req *Request) EnqueueResult:
|
||||||
0.00, если level = None
|
if r.size == r.capacity:
|
||||||
0.00, если level = Low (только задержка)
|
// Переполнение — выбираем политику вытеснения
|
||||||
0.30, если level = Medium
|
if req.Priority > r.slots[r.tail].Priority:
|
||||||
0.70, если level = High
|
evicted = r.slots[r.tail] // вытесняем менее приоритетный
|
||||||
0.90, если level = Critical
|
r.slots[r.tail] = req
|
||||||
}
|
r.tail = (r.tail + 1) mod r.capacity
|
||||||
```
|
return Enqueued(evicted)
|
||||||
</br>
|
else:
|
||||||
|
return Rejected(reason="ring_full")
|
||||||
|
|
||||||
**Коэффициент нагрузки g(load)**
|
r.slots[r.tail] = req
|
||||||
|
r.tail = (r.tail + 1) mod r.capacity
|
||||||
|
r.size++
|
||||||
|
return Enqueued(nil)
|
||||||
|
```
|
||||||
|
|
||||||
|
**3. Дренирование (dequeue)**
|
||||||
|
|
||||||
|
Отдельная горутина-дренажёр (drainer) периодически забирает запросы из головы кольца и передаёт их на выполнение:
|
||||||
|
|
||||||
```sh
|
```sh
|
||||||
g(load) = (cpu_usage + memory_usage + queue_factor + connection_factor) / 4
|
func (r *Ring) drainLoop():
|
||||||
|
ticker = NewTicker(r.drain_interval)
|
||||||
где:
|
for range ticker:
|
||||||
cpu_usage = current_cpu / cpu_threshold
|
for r.size > 0 and systemHasCapacity():
|
||||||
memory_usage = current_memory / memory_threshold
|
req = r.slots[r.head]
|
||||||
queue_factor = min(queue_size / queue_threshold, 1.0)
|
r.slots[r.head] = nil
|
||||||
connection_factor = min(connections / connection_threshold, 1.0)
|
r.head = (r.head + 1) mod r.capacity
|
||||||
|
r.size--
|
||||||
|
go req.resume() // выполняем вне критической секции
|
||||||
```
|
```
|
||||||
</br>
|
|
||||||
|
|
||||||
**Коэффициент времени h(time) (экспоненциальное сглаживание)**
|
|
||||||
|
Почему отдельная горутина: если бы запросы «дренировались» из того же потока, что их принимает, кольцо не сглаживало бы нагрузку, а просто перемещало её в вызывающий код. Дренажёр работает асинхронно и может приостанавливаться, если systemHasCapacity() возвращает `false`.
|
||||||
|
|
||||||
|
**4. Управление задержкой**
|
||||||
|
|
||||||
|
Каждый запрос в кольце имеет deadline — абсолютное время, к которому он должен быть либо выполнен, либо отклонён:
|
||||||
|
|
||||||
```sh
|
```sh
|
||||||
h(time) = 1 - e^(-λ × Δt)
|
req.deadline = now + config.LowDelayMs
|
||||||
|
|
||||||
где:
|
|
||||||
λ = 0.1 (константа скорости затухания)
|
|
||||||
Δt = время с последнего отклонения в секундах
|
|
||||||
```
|
```
|
||||||
</br>
|
Дренажёр перед `resume()` проверяет: если `now > req.deadline` — запрос отклоняется с `504 Gateway Timeout`, а не выполняется «просроченным». Это гарантирует, что клиент получит ответ в предсказуемое время, даже если система перегружена.
|
||||||
|
|
||||||
**Итоговая формула вероятности отклонения**
|
**5. Адаптивная настройка**
|
||||||
|
|
||||||
|
Пороги и задержки пересчитываются по наблюдаемым метрикам:
|
||||||
|
|
||||||
|
* queue_size_threshold — растёт, если кольцо регулярно переполняется, и падает, если оно почти всегда пусто.
|
||||||
|
* low_delay_ms — увеличивается при росте времени дренирования.
|
||||||
|
* high_reject_prob — рассчитывается как доля запросов, которые пришлось вытеснить на прошлом цикле.
|
||||||
|
|
||||||
|
Адаптация выполняется плавно(экспоненциальное сглаживание), чтобы избежать осцилляций.
|
||||||
|
|
||||||
|
**Псевдокод жизненного цикла запроса**
|
||||||
|
|
||||||
```sh
|
```sh
|
||||||
P_reject = f(level) × g(load) × (1 - e^(-0.1 × Δt))
|
Client → API Handler
|
||||||
|
│
|
||||||
|
├─► Backpressure.Check()
|
||||||
|
│ │
|
||||||
|
│ ├─ Low → выполнить немедленно
|
||||||
|
│ ├─ Medium → Enqueue(req) → resume() через low_delay_ms
|
||||||
|
│ ├─ High → Enqueue(req) → resume() при появлении capacity
|
||||||
|
│ └─ Critical → Reject(503)
|
||||||
|
│
|
||||||
|
└─► Ring.Enqueue(req)
|
||||||
|
│
|
||||||
|
├─ свободный слот → положить, вернуть Enqueued
|
||||||
|
├─ переполнение → вытеснить менее приоритетный
|
||||||
|
└─ невозможно → Reject(503)
|
||||||
|
|
||||||
|
Drainer Loop
|
||||||
|
│
|
||||||
|
├─ systemHasCapacity() == true → req = Ring.Dequeue(); go req.resume()
|
||||||
|
├─ systemHasCapacity() == false → sleep(drain_interval)
|
||||||
|
└─ req.deadline < now → Reject(504)
|
||||||
```
|
```
|
||||||
|
|
||||||
**Формула задержки (для уровня Low)**
|
**Параметры конфигурации**
|
||||||
|
|
||||||
```sh
|
| Параметр | Значение по умолчанию | Назначение |
|
||||||
D = D_base × (1 + α × load_factor)
|
|----------|----------------------|------------|
|
||||||
|
| `enabled` | `false` | Включить механизм |
|
||||||
|
| `cpu_threshold` | `0.8` | Порог CPU (0–1), выше которого запросы буферизуются |
|
||||||
|
| `memory_threshold` | `0.85` | Порог памяти (0–1) |
|
||||||
|
| `queue_size_threshold` | `10000` | Порог длины очереди |
|
||||||
|
| `connection_threshold` | `5000` | Порог числа соединений |
|
||||||
|
| `check_interval_ms` | `1000` | Интервал проверки нагрузки |
|
||||||
|
| `low_delay_ms` | `100` | Задержка для уровня `Medium` |
|
||||||
|
| `medium_reject_prob` | `0` | Вероятность отказа при `Medium` (%) |
|
||||||
|
| `high_reject_prob` | `50` | Вероятность отказа при `High` (%) |
|
||||||
|
| `ring_capacity` | `4 × queue_size_threshold` | Размер кольца |
|
||||||
|
|
||||||
где:
|
|
||||||
D_base = 100ms (базовая задержка)
|
|
||||||
α = 2.0 (коэффициент усиления)
|
|
||||||
load_factor = (cpu_usage + memory_usage) / 2
|
|
||||||
```
|
|
||||||
</br>
|
|
||||||
|
|
||||||
**Графическое представление**
|
**Гарантии**
|
||||||
|
|
||||||
```sh
|
* Ограниченная память. Кольцо никогда не превышает ring_capacity слотов.
|
||||||
Вероятность отклонения P_reject
|
* O(1) на операцию. Вставка и извлечение не зависят от размера кольца.
|
||||||
|
|
* Отсутствие блокировок на писателе. Enqueue не ждёт — он либо кладёт, либо вытесняет, либо отказывает.
|
||||||
1.0 | * * * Critical (90%)
|
* Predictable latency. Deadline гарантирует, что клиент получит ответ в предсказуемое время.
|
||||||
| *
|
* Прозрачность. Все отклонения логируются в аудит (BACKPRESSURE_REJECT) с указанием причины и текущего уровня нагрузки.
|
||||||
0.9 | * High (70%)
|
|
||||||
| *
|
|
||||||
0.7 | * Medium (30%)
|
|
||||||
| *
|
|
||||||
0.5 | *
|
|
||||||
| *
|
|
||||||
0.3 | * Low (0% - только задержка)
|
|
||||||
| *
|
|
||||||
0.1 |*
|
|
||||||
|_____________________________ Нагрузка
|
|
||||||
0 0.2 0.4 0.6 0.8 1.0
|
|
||||||
```
|
|
||||||
</br>
|
|
||||||
|
|
||||||
#### Пример рассчёта
|
|
||||||
|
|
||||||
**Исходные данные:**
|
**Ограничения**
|
||||||
|
|
||||||
* Уровень: High → f(level) = 0.70
|
* Не является очередью. Порядок FIFO не гарантируется при вытеснении по приоритету.
|
||||||
|
* Требует настройки под нагрузку. При слишком маленьком кольце и жёстких порогах возможны ложные срабатывания.
|
||||||
|
* Не заменяет rate limiting. Buffer Ring сглаживает всплески, но не ограничивает устойчивую скорость запросов — для этого нужен отдельный слой rate limiter перед API.
|
||||||
|
|
||||||
* CPU: 85% → cpu_usage = 0.85/0.80 = 1.0625
|
**Пример**
|
||||||
|
|
||||||
* Memory: 75% → memory_usage = 0.75/0.85 = 0.882
|
При 10 000 RPS и пороге `cpu_threshold` = 0.8:
|
||||||
|
|
||||||
* Queue: 8000/10000 = 0.8
|
1. На пике 15 000 RPS CPU достигает 0.9 → уровень `High`.
|
||||||
|
2. 5 000 «лишних» запросов попадают в кольцо (capacity 40 000).
|
||||||
|
3. Дренажёр выпускает их по мере освобождения CPU — примерно 200 запросов каждые 10 мс.
|
||||||
|
4. Если пик длится дольше `low_delay_ms`, «просроченные» запросы отклоняются с `504` — клиенты получают быструю деградацию вместо таймаутов.
|
||||||
|
5. Когда нагрузка падает до 8 000 RPS, кольцо пустеет, и запросы снова идут напрямую.
|
||||||
|
|
||||||
* Connections: 4000/5000 = 0.8
|
|
||||||
|
|
||||||
* Время с последнего отклонения: 2 секунды
|
|
||||||
</br>
|
|
||||||
**Рассчёт**
|
|
||||||
```sh
|
|
||||||
load_factor = (1.0625 + 0.882 + 0.8 + 0.8) / 4 = 0.886
|
|
||||||
|
|
||||||
time_factor = 1 - e^(-0.1 × 2) = 1 - e^(-0.2) = 1 - 0.819 = 0.181
|
|
||||||
|
|
||||||
P_reject = 0.70 × 0.886 × 0.181 = 0.112 = 11.2%
|
|
||||||
```
|
|
||||||
**Результат: ~11% запросов будут отклонены.**
|
|
||||||
</br>
|
|
||||||
|
|
||||||
**Преимущества формулы**
|
|
||||||
|
|
||||||
1. **Адаптивность** — реагирует на изменение нагрузки в реальном времени
|
|
||||||
2. **Сглаживание** — предотвращает резкие скачки отклонений
|
|
||||||
3. **Самовосстановление** — при снижении нагрузки вероятность автоматически уменьшается
|
|
||||||
4. **Предсказуемость** — поведение системы становится детерминированным и предсказуемым
|
|
||||||
|
|
||||||
<p align="right">(<a href="#readme-top">К началу</a>)</p>
|
<p align="right">(<a href="#readme-top">К началу</a>)</p>
|
||||||
|
|
||||||
|
|||||||
Reference in new issue
Block a user