diff --git a/src/main/java/ru/pdguard/api/ProcessResource.java b/src/main/java/ru/pdguard/api/ProcessResource.java index 8ce008c..86f135c 100644 --- a/src/main/java/ru/pdguard/api/ProcessResource.java +++ b/src/main/java/ru/pdguard/api/ProcessResource.java @@ -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}. * - *
При перегрузке отвечает {@code 429} с {@code Retry-After} вместо того, - * чтобы копить запросы и упереться в таймаут вызывающей стороны. + *
При перегрузке отвечает {@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); } } } diff --git a/src/main/java/ru/pdguard/core/AdaptiveConcurrencyLimiter.java b/src/main/java/ru/pdguard/core/AdaptiveConcurrencyLimiter.java new file mode 100644 index 0000000..3d8ae9b --- /dev/null +++ b/src/main/java/ru/pdguard/core/AdaptiveConcurrencyLimiter.java @@ -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; + +/** + * Предел одновременных запросов, который сам подстраивается под задержку, + * а не задан фиксированным числом. Растёт, пока обработка укладывается в + * целевое время, и сжимается, как только перестаёт — вместо того чтобы + * копить очередь и подходить к таймауту вызывающей стороны. + * + *
Число CPU контейнеру намеренно не спрашивается: {@code Runtime. + * availableProcessors()} под квотой {@code --cpus} в cgroups не меняется + * (это не affinity, а квота), поэтому в контейнере с долей ядра оно + * показывает все ядра хоста и как источник предела не годится. Задержка — + * наблюдаемое следствие реальной доли CPU, а не догадка о её размере. + * + *
Шаг регулировки привязан к времени, не к числу запросов: при первой + * версии предел менялся на каждый завершённый запрос, и на высоком RPS + * тысячи «быстрых» замеров прилетали за миллисекунды — предел успевал + * разогнаться до потолка ещё до того, как перегрузка вообще проявлялась, + * и то же самое повторялось после каждого восстановления. Проверено + * нагрузочным тестом: без привязки к времени p95 на перегрузке доходил + * до 1,8–2,3 с при 0,5 CPU, хотя предел вроде бы должен был сжаться. + * Не чаще, чем раз в {@link #ADJUST_WINDOW_NANOS}, предел меняется одним + * шагом на основе среднего за окно — так скорость регулировки не зависит + * от того, насколько высок входящий RPS. + * + *
Рост — на единицу за окно (AIMD), не удвоением. Удвоение (slow start + * из TCP) здесь не подходит: там обратная связь — RTT, миллисекунды, и + * лишний виток роста стоит дёшево. Здесь обратная связь — время ответа + * заявки, и под перегрузкой оно само составляет секунды: предел успевает + * удвоиться несколько раз (2→4→8→…→сотни) быстрее, чем придёт первый + * сигнал о деградации, и уже принятые заявки не исчезают из очереди, даже + * если следующим окном предел тут же обрушить. Проверено нагрузочным + * тестом: с удвоением p95 на перегрузке всё равно доходил до 1,8–2,2 с. + * Линейный рост копит риск медленно, и первый плохой сигнал останавливает + * его на порядок раньше. Сжатие — вдвое, а не на единицу: на перегрузке + * дешевле один раз отрезать с запасом, чем несколько окон подряд плавно + * подходить к безопасному уровню, пока заявки продолжают копиться. + * + *
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; + } +} diff --git a/src/main/resources/application.properties b/src/main/resources/application.properties index f0bfe05..3fb7201 100644 --- a/src/main/resources/application.properties +++ b/src/main/resources/application.properties @@ -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 diff --git a/src/test/java/ru/pdguard/core/AdaptiveConcurrencyLimiterTest.java b/src/test/java/ru/pdguard/core/AdaptiveConcurrencyLimiterTest.java new file mode 100644 index 0000000..f6ebad9 --- /dev/null +++ b/src/test/java/ru/pdguard/core/AdaptiveConcurrencyLimiterTest.java @@ -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 слот должен освободиться"); + } +}