feat(backpressure): окно поглощения writer↔consumer (защита Redis RAM на репарсинге) #10

Merged
claude merged 1 commits from dev into main 2026-06-03 16:08:12 +00:00
Owner

Parser на длинном репарсинге XADD'ит быстрее, чем consumer читает → Redis-стрим растёт неограниченно → OOM.

BackpressureGate: перед обработкой блока проверяем backlog (XLEN). >= highWater → пауза чтения SHiP (блок не ack'ается, SHiP тормозит), ждём слива до lowWater. Ничего не теряем. Нет consumer → writer замирает на highWater.

Конфиг backpressure {enabled(true), highWater(100k), lowWater(50k), pollMs(200)}. Включено по умолчанию.

Проверено на реальном Redis (integration в CI). Unit 221, typecheck чистый. parser2 1.1.3→1.2.0.

Parser на длинном репарсинге XADD'ит быстрее, чем consumer читает → Redis-стрим растёт неограниченно → OOM. BackpressureGate: перед обработкой блока проверяем backlog (XLEN). >= highWater → пауза чтения SHiP (блок не ack'ается, SHiP тормозит), ждём слива до lowWater. Ничего не теряем. Нет consumer → writer замирает на highWater. Конфиг backpressure {enabled(true), highWater(100k), lowWater(50k), pollMs(200)}. Включено по умолчанию. Проверено на реальном Redis (integration в CI). Unit 221, typecheck чистый. parser2 1.1.3→1.2.0.
claude added 1 commit 2026-06-03 16:08:00 +00:00
feat(backpressure): окно поглощения writer↔consumer — пауза чтения SHiP по backlog стрима
CI / Typecheck (pull_request) Successful in 9m25s
CI / Lint (pull_request) Successful in 9m27s
CI / Unit tests (pull_request) Successful in 9m32s
CI / Build (pull_request) Successful in 9m26s
CI / Integration tests (pull_request) Successful in 13m37s
87af7a89d2
Проблема: на длинном репарсинге (genesis→head) Parser XADD'ит события в Redis
быстрее, чем consumer (controller) их вычитывает. RAM парсера ограничена
SHiP-окном (max_messages_in_flight), а Redis-стрим — нет → растёт неограниченно
в памяти Redis → OOM. XtrimSupervisor не режет непрочитанное.

Решение: BackpressureGate. Перед обработкой каждого блока Parser проверяет
backlog = XLEN стрима. Если backlog >= highWater — чтение SHiP встаёт на паузу
(блок не ack'ается, SHiP сам тормозит после in-flight окна), ждём пока consumer
сольёт backlog <= lowWater, затем продолжаем. Ничего не теряется — события
остаются в стриме. Нет consumer → стрим дорастает до highWater и writer замирает.

Метрика backlog = XLEN: ровно та память Redis, что бережём; не зависит от числа
групп и версии Redis.

Конфиг backpressure: { enabled (default true), highWater (100k), lowWater (50k),
pollMs (200) }. По умолчанию включено — защита репарсинга из коробки.

Также: добавлен del в IRedisClient (латентный type-error в xtrim-интеграционном
тесте из прошлого коммита — tsc по тестам не гонялся в release-пуше).

Проверено на реальном Redis (integration-тест в CI): пауза при backlog>=highWater,
возобновление после слива. Unit: 221 зелёных, typecheck чистый. parser2 1.1.3→1.2.0.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
claude merged commit 3dbd2c9116 into main 2026-06-03 16:08:12 +00:00
Sign in to join this conversation.