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-теста ### Бенчмарк