From be67debf51216ed1451d8ba14c190547762c3371 Mon Sep 17 00:00:00 2001 From: coopops Date: Thu, 4 Jun 2026 05:43:48 +0000 Subject: [PATCH] =?UTF-8?q?docs(readme):=20=D1=80=D0=B0=D0=B7=D0=B4=D0=B5?= =?UTF-8?q?=D0=BB=20=D0=BF=D1=80=D0=BE=20backpressure-=D0=BE=D0=BA=D0=BD?= =?UTF-8?q?=D0=BE=20=D1=81=D0=BE=20=D1=81=D1=85=D0=B5=D0=BC=D0=BE=D0=B9=20?= =?UTF-8?q?+=20=D0=BA=D0=BE=D0=BD=D1=84=D0=B8=D0=B3=20=D0=B2=20=D0=BF?= =?UTF-8?q?=D1=80=D0=B8=D0=BC=D0=B5=D1=80=D0=B0=D1=85?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Описание окна поглощения writer↔consumer (highWater/lowWater гистерезис, пауза чтения SHiP по XLEN, антидребезг), ASCII-схема качелей backlog, рекомендация maxmemory+noeviction. Конфиг backpressure добавлен в оба YAML-примера, строка Backpr. в арх-диаграмму, актуализированы счётчики тестов (221/54). Co-Authored-By: Claude Opus 4.8 (1M context) --- README.md | 53 ++++++++++++++++++++++++++++++++++++++++++++++++++--- 1 file changed, 50 insertions(+), 3 deletions(-) diff --git a/README.md b/README.md index d54bcb6..0cec704 100644 --- a/README.md +++ b/README.md @@ -49,6 +49,11 @@ abiFallback: rpc-current xtrim: enabled: true intervalMs: 60000 +backpressure: + enabled: true # пауза чтения SHiP когда стрим переполнен (защита Redis RAM) + highWater: 100000 # порог паузы — число событий в стриме + lowWater: 50000 # порог возобновления (должен быть < highWater) + pollMs: 200 # период опроса XLEN во время паузы, мс logger: level: info pretty: false @@ -185,6 +190,11 @@ abiFallback: rpc-current xtrim: enabled: true intervalMs: 60000 +backpressure: + enabled: true # пауза чтения SHiP когда стрим переполнен (защита Redis RAM) + highWater: 100000 # порог паузы — число событий в стриме + lowWater: 50000 # порог возобновления (должен быть < highWater) + pollMs: 200 # период опроса XLEN во время паузы, мс reconnect: maxAttempts: 10 backoffSeconds: [1, 2, 5, 10, 30, 60, 120, 300, 600, 1800] @@ -224,7 +234,8 @@ docker run --rm --network host \ │ • Workers │ │ │ │ • ForkDet. │ │ Streams + │ │ • XtrimSup. │ │ Hashes + │ - └──────────────┘ │ ZSets │ + │ • Backpr. │ │ ZSets │ + └──────────────┘ │ │ │ │ ┌──────────────┐ │ │ │ ParserClient │◀─XREADGR──│ │ @@ -240,6 +251,42 @@ docker run --rm --network host \ - [Redis key taxonomy](docs/redis-key-taxonomy.md) - [Disaster recovery / fork scenarios](docs/disaster-recovery.md) +## Backpressure: защита памяти Redis на репарсинге + +Parser читает блоки из SHiP и пишет события в Redis-стрим. На длинном репарсинге (genesis → head, миллионы блоков) писатель легко обгоняет потребителя: события копятся в **памяти Redis** быстрее, чем `ParserClient` их вычитывает. RAM самого парсера ограничена окном SHiP (`max_messages_in_flight` + pull-based чтение с ack после обработки), а вот стрим в Redis может расти неограниченно — вплоть до OOM. `XtrimSupervisor` тут не спасает: он по дизайну не режет непрочитанные события. + +`BackpressureGate` вводит **окно поглощения с гистерезисом**. Перед обработкой каждого блока парсер смотрит backlog = `XLEN` стрима: + +- backlog ≥ **highWater** → чтение SHiP встаёт на **паузу** (блок не подтверждается, SHiP сам перестаёт слать после своего in-flight окна); +- во время паузы парсер опрашивает `XLEN` каждые `pollMs` и ждёт, пока потребитель сольёт backlog ≤ **lowWater**; +- слилось → чтение **возобновляется**. + +Ничего не теряется — события остаются в стриме, парсер лишь перестаёт доливать. Если потребитель вообще не подключён, backlog дорастает до `highWater` и писатель замирает, удерживая Redis в пределах ~`highWater` событий. + +**Почему не дёргается.** Разрыв high/low (по умолчанию 100k / 50k) — это антидребезг (гистерезис). Один цикл «пауза ↔ работа» перемалывает 50k событий, поэтому переключения происходят редко, а не на каждом блоке. В steady-state (потребитель у головы цепи) backlog почти нулевой → пауза не включается вообще, оверхед = один `XLEN` на блок (O(1)). + +``` +backlog +100k | /| /| /| <- highWater -> ПАУЗА + | / | / | / | + 50k | / |______/ |______/ |___ <- lowWater -> ПУСК + | / (пауза) (пауза) + 0 +--+----------------------------> время + пишем ждём пишем ждём пишем +``` + +Конфиг (включён по умолчанию): + +```yaml +backpressure: + enabled: true # выключить можно enabled: false + highWater: 100000 # порог паузы — число событий в стриме + lowWater: 50000 # порог возобновления (< highWater) + pollMs: 200 # период опроса XLEN во время паузы, мс +``` + +Рекомендация по инфраструктуре: на Redis выставить `maxmemory` + `maxmemory-policy noeviction` — тогда при достижении лимита `XADD` вернёт ошибку (парсер её обработает), а не сработает OOM-killer. + ## Известные ограничения ### Схлопывание дельт create+remove внутри одного блока @@ -298,8 +345,8 @@ pnpm --filter @coopenomics/parser2 test:integration # нужен Docker для ``` **Текущее покрытие:** -- `@coopenomics/parser2`: 81% statements / 74% functions (205 unit-тестов) -- `@coopenomics/coopos-ship-reader`: 81% statements / 86% functions (51 unit-тест) +- `@coopenomics/parser2`: 221 unit-тестов + integration (xtrim/backpressure на реальном Redis) +- `@coopenomics/coopos-ship-reader`: 54 unit-теста ### Бенчмарк -- 2.52.0