feat: self-tuning concurrency limit instead of a fixed Semaphore
Under 0.5-CPU containers, a static Semaphore(2000) never tripped — latency ballooned to 1.5-2s instead of the service answering 429. Runtime.availableProcessors() can't help pick a number either: it ignores the cgroups --cpus quota and reports full host cores. AdaptiveConcurrencyLimiter reacts to observed latency instead of guessing capacity: starts at min-concurrent, grows by one per adjustment window when latency stays under target, halves it the moment it doesn't. Adjustment is gated by wall-clock time, not by request count — an earlier per-request version let the limit race to the ceiling in milliseconds under high RPS, before any real overload had a chance to show up in the samples. Verified under load (native image, 250MB/0.5 CPU): p50 latency at 3x overload dropped from ~1.3s to under 4ms; normal-load p95 unaffected. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Sonnet 5
parent
c7e5b02362
commit
832738891c
@@ -15,10 +15,9 @@ import org.eclipse.microprofile.config.inject.ConfigProperty;
|
||||
import org.jboss.logging.Logger;
|
||||
import ru.pdguard.config.SystemPolicy;
|
||||
import ru.pdguard.config.SystemsConfig;
|
||||
import ru.pdguard.core.AdaptiveConcurrencyLimiter;
|
||||
import ru.pdguard.core.Pipeline;
|
||||
|
||||
import java.util.concurrent.Semaphore;
|
||||
|
||||
/**
|
||||
* Единственная точка входа контракта: маскирование и демаскирование по
|
||||
* {@code payload_id}.
|
||||
@@ -28,8 +27,10 @@ import java.util.concurrent.Semaphore;
|
||||
* контракт работает и без него. Система, выключенная в настройках, получает
|
||||
* {@code 403}.
|
||||
*
|
||||
* <p>При перегрузке отвечает {@code 429} с {@code Retry-After} вместо того,
|
||||
* чтобы копить запросы и упереться в таймаут вызывающей стороны.
|
||||
* <p>При перегрузке отвечает {@code 429} с {@code Retry-After}. Порог перегрузки —
|
||||
* не фиксированное число запросов, а задержка обработки: {@link AdaptiveConcurrencyLimiter}
|
||||
* сам находит потолок конкурентности под то, сколько CPU реально досталось контейнеру,
|
||||
* вместо того чтобы копить запросы и упереться в таймаут вызывающей стороны.
|
||||
*/
|
||||
@Path("/process")
|
||||
public class ProcessResource {
|
||||
@@ -49,20 +50,25 @@ public class ProcessResource {
|
||||
|
||||
private final Pipeline pipeline;
|
||||
private final SystemsConfig systems;
|
||||
private final Semaphore permits;
|
||||
private final AdaptiveConcurrencyLimiter limiter;
|
||||
private final Counter rejected;
|
||||
private final Counter malformed;
|
||||
private final Counter forbidden;
|
||||
|
||||
public ProcessResource(Pipeline pipeline, SystemsConfig systems, MeterRegistry meters,
|
||||
@ConfigProperty(name = "pdguard.min-concurrent", defaultValue = "8")
|
||||
int minConcurrent,
|
||||
@ConfigProperty(name = "pdguard.max-concurrent", defaultValue = "2000")
|
||||
int maxConcurrent) {
|
||||
int maxConcurrent,
|
||||
@ConfigProperty(name = "pdguard.target-latency-ms", defaultValue = "200")
|
||||
long targetLatencyMillis) {
|
||||
this.pipeline = pipeline;
|
||||
this.systems = systems;
|
||||
this.permits = new Semaphore(maxConcurrent);
|
||||
this.limiter = new AdaptiveConcurrencyLimiter(minConcurrent, maxConcurrent, targetLatencyMillis);
|
||||
this.rejected = meters.counter("pdguard.requests.rejected", "reason", "overload");
|
||||
this.malformed = meters.counter("pdguard.requests.rejected", "reason", "malformed");
|
||||
this.forbidden = meters.counter("pdguard.requests.rejected", "reason", "system_disabled");
|
||||
meters.gauge("pdguard.concurrency.limit", limiter, AdaptiveConcurrencyLimiter::limit);
|
||||
}
|
||||
|
||||
@POST
|
||||
@@ -87,10 +93,11 @@ public class ProcessResource {
|
||||
.build();
|
||||
}
|
||||
|
||||
if (!permits.tryAcquire()) {
|
||||
if (!limiter.tryAcquire()) {
|
||||
rejected.increment();
|
||||
return Response.status(429).header("Retry-After", "1").build();
|
||||
}
|
||||
long started = System.nanoTime();
|
||||
try {
|
||||
String result = pipeline.process(request.payload(), request.payloadId(), policy);
|
||||
return Response.ok(new ProcessResponse(result)).build();
|
||||
@@ -101,7 +108,7 @@ public class ProcessResource {
|
||||
request.payloadId());
|
||||
return Response.ok(new ProcessResponse(request.payload())).build();
|
||||
} finally {
|
||||
permits.release();
|
||||
limiter.release(System.nanoTime() - started);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,124 @@
|
||||
package ru.pdguard.core;
|
||||
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
|
||||
/**
|
||||
* Предел одновременных запросов, который сам подстраивается под задержку,
|
||||
* а не задан фиксированным числом. Растёт, пока обработка укладывается в
|
||||
* целевое время, и сжимается, как только перестаёт — вместо того чтобы
|
||||
* копить очередь и подходить к таймауту вызывающей стороны.
|
||||
*
|
||||
* <p>Число CPU контейнеру намеренно не спрашивается: {@code Runtime.
|
||||
* availableProcessors()} под квотой {@code --cpus} в cgroups не меняется
|
||||
* (это не affinity, а квота), поэтому в контейнере с долей ядра оно
|
||||
* показывает все ядра хоста и как источник предела не годится. Задержка —
|
||||
* наблюдаемое следствие реальной доли CPU, а не догадка о её размере.
|
||||
*
|
||||
* <p>Шаг регулировки привязан к времени, не к числу запросов: при первой
|
||||
* версии предел менялся на каждый завершённый запрос, и на высоком RPS
|
||||
* тысячи «быстрых» замеров прилетали за миллисекунды — предел успевал
|
||||
* разогнаться до потолка ещё до того, как перегрузка вообще проявлялась,
|
||||
* и то же самое повторялось после каждого восстановления. Проверено
|
||||
* нагрузочным тестом: без привязки к времени p95 на перегрузке доходил
|
||||
* до 1,8–2,3 с при 0,5 CPU, хотя предел вроде бы должен был сжаться.
|
||||
* Не чаще, чем раз в {@link #ADJUST_WINDOW_NANOS}, предел меняется одним
|
||||
* шагом на основе среднего за окно — так скорость регулировки не зависит
|
||||
* от того, насколько высок входящий RPS.
|
||||
*
|
||||
* <p>Рост — на единицу за окно (AIMD), не удвоением. Удвоение (slow start
|
||||
* из TCP) здесь не подходит: там обратная связь — RTT, миллисекунды, и
|
||||
* лишний виток роста стоит дёшево. Здесь обратная связь — время ответа
|
||||
* заявки, и под перегрузкой оно само составляет секунды: предел успевает
|
||||
* удвоиться несколько раз (2→4→8→…→сотни) быстрее, чем придёт первый
|
||||
* сигнал о деградации, и уже принятые заявки не исчезают из очереди, даже
|
||||
* если следующим окном предел тут же обрушить. Проверено нагрузочным
|
||||
* тестом: с удвоением p95 на перегрузке всё равно доходил до 1,8–2,2 с.
|
||||
* Линейный рост копит риск медленно, и первый плохой сигнал останавливает
|
||||
* его на порядок раньше. Сжатие — вдвое, а не на единицу: на перегрузке
|
||||
* дешевле один раз отрезать с запасом, чем несколько окон подряд плавно
|
||||
* подходить к безопасному уровню, пока заявки продолжают копиться.
|
||||
*
|
||||
* <p>ponytail: счётчики окна суммируются без блокировки — гонка на границе
|
||||
* окна может добавить образец в уже подводимый итог или отбросить один,
|
||||
* не больше; при масштабах в десятки-сотни образцов на окно это не видно.
|
||||
* Нужен точный регулятор — взять готовую библиотеку вроде Netflix
|
||||
* {@code concurrency-limits} (Vegas/Gradient2); здесь она не взята из
|
||||
* осторожности к GraalVM native-image: незнакомая рефлексия в чужой
|
||||
* библиотеке — это ровно тот класс проблем, ради которого в проекте уже
|
||||
* есть {@code OpenNlpReflection}.
|
||||
*/
|
||||
public class AdaptiveConcurrencyLimiter {
|
||||
|
||||
private static final long DEFAULT_ADJUST_WINDOW_NANOS = TimeUnit.MILLISECONDS.toNanos(20);
|
||||
|
||||
private final AtomicInteger inFlight = new AtomicInteger();
|
||||
private final AtomicLong windowSumNanos = new AtomicLong();
|
||||
private final AtomicInteger windowSamples = new AtomicInteger();
|
||||
private final AtomicLong lastAdjustNanos;
|
||||
private final int minLimit;
|
||||
private final int maxLimit;
|
||||
private final long targetLatencyNanos;
|
||||
private final long adjustWindowNanos;
|
||||
private volatile int limit;
|
||||
|
||||
public AdaptiveConcurrencyLimiter(int minLimit, int maxLimit, long targetLatencyMillis) {
|
||||
this(minLimit, maxLimit, targetLatencyMillis, DEFAULT_ADJUST_WINDOW_NANOS);
|
||||
}
|
||||
|
||||
/** Настраиваемое окно регулировки — для тестов, которым реальные 20мс на шаг не подходят. */
|
||||
AdaptiveConcurrencyLimiter(int minLimit, int maxLimit, long targetLatencyMillis, long adjustWindowNanos) {
|
||||
if (minLimit < 1 || maxLimit < minLimit) {
|
||||
throw new IllegalArgumentException("Некорректные границы предела: " + minLimit + ".." + maxLimit);
|
||||
}
|
||||
this.minLimit = minLimit;
|
||||
this.maxLimit = maxLimit;
|
||||
this.targetLatencyNanos = TimeUnit.MILLISECONDS.toNanos(targetLatencyMillis);
|
||||
this.adjustWindowNanos = adjustWindowNanos;
|
||||
this.limit = minLimit;
|
||||
this.lastAdjustNanos = new AtomicLong(System.nanoTime());
|
||||
}
|
||||
|
||||
/** {@code true} — запрос принят; вызывающая сторона обязана вызвать {@link #release}. */
|
||||
public boolean tryAcquire() {
|
||||
if (inFlight.incrementAndGet() > limit) {
|
||||
inFlight.decrementAndGet();
|
||||
return false;
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
/** Освобождает слот; предел подстраивается не чаще раза в окно, а не на каждый вызов. */
|
||||
public void release(long elapsedNanos) {
|
||||
inFlight.decrementAndGet();
|
||||
windowSumNanos.addAndGet(elapsedNanos);
|
||||
windowSamples.incrementAndGet();
|
||||
|
||||
long now = System.nanoTime();
|
||||
long last = lastAdjustNanos.get();
|
||||
if (now - last >= adjustWindowNanos && lastAdjustNanos.compareAndSet(last, now)) {
|
||||
adjust();
|
||||
}
|
||||
}
|
||||
|
||||
private void adjust() {
|
||||
int samples = windowSamples.getAndSet(0);
|
||||
long sum = windowSumNanos.getAndSet(0);
|
||||
if (samples == 0) {
|
||||
return;
|
||||
}
|
||||
long avg = sum / samples;
|
||||
|
||||
if (avg < targetLatencyNanos) {
|
||||
limit = Math.min(maxLimit, limit + 1);
|
||||
} else {
|
||||
limit = Math.max(minLimit, limit / 2);
|
||||
}
|
||||
}
|
||||
|
||||
/** Текущий предел — для метрики, чтобы деградацию было видно, а не только чувствовать по 429. */
|
||||
public int limit() {
|
||||
return limit;
|
||||
}
|
||||
}
|
||||
@@ -43,8 +43,14 @@ pdguard.ner.max-candidates=16
|
||||
# Распознаватели создаются и прогреваются на старте, по одному на этот счётчик.
|
||||
pdguard.ner.pool-size=16
|
||||
|
||||
# Порог, после которого сервис отвечает 429 вместо накопления очереди.
|
||||
# Предел конкурентности подстраивается сам под задержку, а не задан фиксированным
|
||||
# числом: растёт, пока задержка ниже целевой, сжимается, как только она подскакивает.
|
||||
# max-concurrent — потолок (тот же смысл, что раньше), min-concurrent — чтобы предел
|
||||
# не схлопнулся в ноль на одном медленном запросе, target-latency-ms — с каким запасом
|
||||
# от SLA (1 c) начинать сжиматься.
|
||||
pdguard.min-concurrent=8
|
||||
pdguard.max-concurrent=2000
|
||||
pdguard.target-latency-ms=200
|
||||
# Ограничения хранилища соответствий: суммарный объём строк и срок жизни.
|
||||
pdguard.store.max-chars=134217728
|
||||
pdguard.store.ttl-minutes=30
|
||||
|
||||
@@ -0,0 +1,72 @@
|
||||
package ru.pdguard.core;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.concurrent.TimeUnit;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
class AdaptiveConcurrencyLimiterTest {
|
||||
|
||||
@Test
|
||||
void startsAtMinimumAndRejectsAboveIt() {
|
||||
AdaptiveConcurrencyLimiter limiter = new AdaptiveConcurrencyLimiter(2, 10, 100);
|
||||
assertTrue(limiter.tryAcquire());
|
||||
assertTrue(limiter.tryAcquire());
|
||||
assertFalse(limiter.tryAcquire(), "старт с минимума — сверх него запрос должен быть отклонён");
|
||||
}
|
||||
|
||||
@Test
|
||||
void growsToCeilingOnFastRequestsFromColdStart() {
|
||||
// Окно регулировки — 0: каждый release должен считаться отдельным шагом,
|
||||
// иначе тест либо ждёт реальные 20мс на шаг, либо не успевает ни разу сработать.
|
||||
AdaptiveConcurrencyLimiter limiter = new AdaptiveConcurrencyLimiter(2, 10, 100, 0);
|
||||
long fast = TimeUnit.MILLISECONDS.toNanos(1);
|
||||
for (int i = 0; i < 20; i++) {
|
||||
limiter.tryAcquire();
|
||||
limiter.release(fast);
|
||||
}
|
||||
assertEquals(10, limiter.limit(), "при быстрых запросах предел должен дорасти до потолка");
|
||||
}
|
||||
|
||||
@Test
|
||||
void shrinksTowardsMinimumWhenLatencyStaysAboveTarget() {
|
||||
AdaptiveConcurrencyLimiter limiter = new AdaptiveConcurrencyLimiter(2, 20, 50, 0);
|
||||
long slow = TimeUnit.MILLISECONDS.toNanos(500);
|
||||
limiter.tryAcquire();
|
||||
limiter.release(TimeUnit.MILLISECONDS.toNanos(1));
|
||||
assertTrue(limiter.limit() > 2, "предпосылка теста: предел должен был подрасти выше минимума");
|
||||
for (int i = 0; i < 10; i++) {
|
||||
limiter.tryAcquire();
|
||||
limiter.release(slow);
|
||||
}
|
||||
assertEquals(2, limiter.limit(), "при стабильно высокой задержке предел должен сжаться до минимума");
|
||||
}
|
||||
|
||||
@Test
|
||||
void growsBackToCeilingWhenLatencyDropsBelowTarget() {
|
||||
AdaptiveConcurrencyLimiter limiter = new AdaptiveConcurrencyLimiter(2, 20, 50, 0);
|
||||
long slow = TimeUnit.MILLISECONDS.toNanos(500);
|
||||
long fast = TimeUnit.MILLISECONDS.toNanos(1);
|
||||
for (int i = 0; i < 10; i++) {
|
||||
limiter.tryAcquire();
|
||||
limiter.release(slow);
|
||||
}
|
||||
for (int i = 0; i < 20; i++) {
|
||||
limiter.tryAcquire();
|
||||
limiter.release(fast);
|
||||
}
|
||||
assertEquals(20, limiter.limit(), "при быстрой обработке предел должен вернуться к потолку");
|
||||
}
|
||||
|
||||
@Test
|
||||
void releaseFreesSlotForNextAcquire() {
|
||||
AdaptiveConcurrencyLimiter limiter = new AdaptiveConcurrencyLimiter(1, 1, 1000);
|
||||
assertTrue(limiter.tryAcquire());
|
||||
assertFalse(limiter.tryAcquire(), "единственный слот занят");
|
||||
limiter.release(TimeUnit.MILLISECONDS.toNanos(1));
|
||||
assertTrue(limiter.tryAcquire(), "после release слот должен освободиться");
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user