Merge pull request 'docs(readme): backpressure-окно со схемой + конфиг' (#11) from dev into main
Release / Release (push) Successful in 9m21s
Release / Release (push) Successful in 9m21s
This commit was merged in pull request #11.
This commit is contained in:
@@ -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-теста
|
||||
|
||||
### Бенчмарк
|
||||
|
||||
|
||||
Reference in New Issue
Block a user