ЁЯПл The SchoolтА║ЁЯМР Distributed SystemsтА║ЁЯЪз рдзрдбрд╛ 12 тАФ Backpressure & рд╕рдВрдкреВрд░реНрдг рдЪрд┐рддреНрд░: "рдирдВрддрд░ рдпрд╛" рд▓рд╡рдХрд░ рд╕рд╛рдВрдЧрд╛
ЁЯЦ╝я╕П See the drawing + lab ЁЯПа Course home ЁЯМ┐ Branch on GitHub тЬПя╕П View source
ЁЯЦ╝я╕П рдЖрдХреГрддреА рдЖрдгрд┐ labThe drawing + lab рдкреВрд░реНрдг рдкрд╛рдирд╛рд╡рд░ рдЙрдШрдбрд╛ тЖЧOpen full page тЖЧ

ЁЯЪз рдзрдбрд╛ 12 тАФ Backpressure & рд╕рдВрдкреВрд░реНрдг рдЪрд┐рддреНрд░: "рдирдВрддрд░ рдпрд╛" рд▓рд╡рдХрд░ рд╕рд╛рдВрдЧрд╛

ЁЯУН рддреБрдореНрд╣реА рдЗрдереЗ рдЖрд╣рд╛рдд: 12 рдкреИрдХреА рдзрдбрд╛ 12 ┬╖ рдорд╛рдЧреЗ: lesson-11-exactly-once ┬╖ рд╢реЗрд╡рдЯрдЪрд╛ рдзрдбрд╛ ЁЯОУ


ЁЯУж рдпрд╛ рдмреНрд░рдБрдЪрдордзреНрдпреЗ рдХрд╛рдп рдЖрд╣реЗ

рдзрдбреЗ 01тАУ11, рдЖрдгрд┐ backpressure: node рдХрд░реВ рд╢рдХрддреЗ рддреНрдпрд╛рдкреЗрдХреНрд╖рд╛ рдЬрд╛рд╕реНрдд рдХрд╛рдо рдЖрд▓реЗ рдХреА unbounded рд░рд╛рдВрдЧ рддреНрдпрд╛ overload рдЪреЗ рд░реВрдкрд╛рдВрддрд░ рд╕рддрдд рд╡рд╛рдврдд рдЬрд╛рдгрд╛рд▒реНрдпрд╛ рдерд╛рдВрдмрдгреНрдпрд╛рдд рдХрд░рддреЗ; bounded рд░рд╛рдВрдЧ "рдирдВрддрд░ рдпрд╛" (429 / 503 + Retry-After) рдЕрд╕реЗ рд╕рд╛рдВрдЧреВрди рд▓рд╡рдХрд░ рдирдХрд╛рд░ рджреЗрддреЗ рдЖрдгрд┐ рдерд╛рдВрдмрдгреЗ рдХрдореА рдареЗрд╡рддреЗ. рд╢рд┐рд╡рд╛рдп load shedding, backoff рдЖрдгрд┐ jitter рд╕рд╣ retries, rate limiting рдЖрдгрд┐ circuit breaking тАФ рдЖрдгрд┐ рдХреЛрд░реНрд╕рдЪрд╛ рд╕рдВрдкреВрд░реНрдг рдирдХрд╛рд╢рд╛. dist/demo.py рдордзреАрд▓ backpressure() рдЖрдгрд┐ dist/sim.py рдордзреАрд▓ BoundedQueue.

ЁЯзТ 5 рд╡рд░реНрд╖рд╛рдВрдЪреНрдпрд╛ рдореБрд▓рд╛рд▓рд╛ рд╕рдордЬрд╛рд╡рд▓реНрдпрд╛рд╕рд╛рд░рдЦреЗ

рдкреНрд░рд╡реЗрд╢рд╛рдЪреНрдпрд╛ рджрд┐рд╡рд╢реА рдкреБрдгреЗ office рдордзреНрдпреЗ рджрд░ second рд▓рд╛ 120 рдкрд╛рд▓рдХ рдпреЗрддрд╛рдд. рдХрд╛рд░рдХреВрди рджрд░ second рд▓рд╛ 100 рдЬрдгрд╛рдВрдирд╛ рдорджрдд рдХрд░реВ рд╢рдХрддрд╛рдд. рдкреНрд░рддреНрдпреЗрдХ second рд▓рд╛ рдЖрдгрдЦреА 20 рдкрд╛рд▓рдХ рд░рд╛рдВрдЧреЗрдд рдЬреЛрдбрд▓реЗ рдЬрд╛рддрд╛рдд.

рдХрд╛рд╣реА рд▓реЛрдХрд╛рдВрдирд╛ рдкрд░рдд рдкрд╛рдард╡рдгреЗ рдирд┐рд╖реНрдареБрд░ рд╡рд╛рдЯрддреЗ. рдкрдг рджреБрд╕рд░рд╛ рдкрд░реНрдпрд╛рдп рдореНрд╣рдгрдЬреЗ рд╕рдЧрд│реНрдпрд╛рдВрдирд╛рдЪ рдЦреВрдк рд╡реЗрд│ рдерд╛рдВрдмрд╛рд╡реЗ рд▓рд╛рдЧрддреЗ, рдЖрдгрд┐ office рд╕рдЧрд│реНрдпрд╛рдВрд╕рд╛рдареАрдЪ рдХрд╛рдо рдХрд░рдгреЗ рдерд╛рдВрдмрд╡рддреЗ.

рдЖрдгрд┐ рдкрд░рдд рдкрд╛рдард╡рд▓реЗрд▓реНрдпрд╛ рдкрд╛рд▓рдХрд╛рдВрд╕рд╛рдареА рдЖрдгрдЦреА рдПрдХ рдирд┐рдпрдо: рд╕рдЧрд│реЗ рдПрдХрд╛рдЪ рдХреНрд╖рдгреА рдкрд░рдд рдпреЗрдК рдирдХрд╛. рдереЛрдбреЗ рдерд╛рдВрдмрд╛, рдордЧ рдЖрдгрдЦреА рдереЛрдбреЗ рдЬрд╛рд╕реНрдд, рдЖрдгрд┐ рдереЛрдбрд╛рд╕рд╛ random рд╡реЗрд│ рдЬреЛрдбрд╛, рдореНрд╣рдгрдЬреЗ рдкреБрдврдЪреА рдЧрд░реНрджреА рд╡рд┐рдЦреБрд░рд▓реА рдЬрд╛рдИрд▓.

ЁЯЧ║я╕П рдЖрдХреГрддреА

flowchart LR
    in["ЁЯСк 120 requests/s arrive"] --> gate{"ЁЯЪз queue full?<br/>limit 200"}
    gate -->|"yes"| shed["тЫФ 429 / 503 + Retry-After<br/>shed 100 in 10 s"]
    gate -->|"no"| q["ЁЯУЛ bounded queue"]
    q --> srv["ЁЯзСтАНЁЯТ╝ serve 100/s<br/>wait stays about 1.0 s"]
    un["unbounded queue: wait 1.0 s after 5 s,<br/>2.0 s after 10 s, and still growing"]

ЁЯЧ║я╕П рд░реЗрдЦрд╛рдЯрд▓реЗрд▓реА рдЖрд╡реГрддреНрддреА + рдПрдХ lab: https://school-edh.pages.dev/distributed-systems/lesson-diagrams.html#l12

тЭУ рдХрд╛рдп

ЁЯдФ рдХрд╛

рдХрд╛рд░рдг рдпрд╛рдЖрдзреАрдЪрд╛ рдкреНрд░рддреНрдпреЗрдХ рдзрдбрд╛ рддрд╛рдгрд╛рдЦрд╛рд▓реА рдЬрд╛рд╕реНрддреАрдЪреЗ рдХрд╛рдо рдирд┐рд░реНрдорд╛рдг рдХрд░рддреЛ: timeouts рдирдВрддрд░рдЪреЗ retries (01, 11), failovers (05), repairs (06), elections (09). рдорд░реНрдпрд╛рджрд╛ рдирд╕рд▓реЗрд▓реА system рдЫреЛрдЯреНрдпрд╛ overload рдЪреЗ рдкреВрд░реНрдг overload рдордзреНрдпреЗ рд░реВрдкрд╛рдВрддрд░ рдХрд░рддреЗ тАФ рд░рд╛рдВрдЧ рд╡рд╛рдврддреЗ, рдкреНрд░рддреНрдпреЗрдХ request time out рд╣реЛрддреЗ, рдкреНрд░рддреНрдпреЗрдХ client retry рдХрд░рддреЛ, рдЖрдгрд┐ load рджреБрдкреНрдкрдЯ рд╣реЛрддреЛ. рд▓рд╡рдХрд░ "рдирдВрддрд░ рдпрд╛" рд╕рд╛рдВрдЧрдгреЗ рд╣реЗрдЪ рдмрд╛рдХреАрдЪреА system рдЪрд╛рд▓реВ рдареЗрд╡рддреЗ.

ЁЯФз рдХрд╕реЗ (рдпрд╛ repo рдордзреНрдпреЗ)

dist/sim.py рдордзреАрд▓ BoundedQueue(limit) рджрд░ second рд▓рд╛ рдПрдХ tick рдЪрд╛рд▓рд╡рддреЗ: run(arrivals, served_per_tick) рд░рд┐рдХрд╛рдореНрдпрд╛ рдЬрд╛рдЧреЗрдкрд░реНрдпрдВрдд (limit - q; limit=None рдореНрд╣рдгрдЬреЗ рдорд░реНрдпрд╛рджрд╛ рдирд╛рд╣реА) рдирд╡реНрдпрд╛ arrivals рдШреЗрддреЗ, рдЙрд░рд▓реЗрд▓реНрдпрд╛ shed рдордзреНрдпреЗ рдореЛрдЬрддреЗ, served_per_tick рдкрд░реНрдпрдВрдд serve рдХрд░рддреЗ, рдЖрдгрд┐ рдерд╛рдВрдмрдгреЗ queue ├╖ served_per_tick seconds рдореНрд╣рдгреВрди рдиреЛрдВрджрд╡рддреЗ. dist/demo.py рдордзреАрд▓ backpressure() 100/s service рд╡рд┐рд░реБрджреНрдз 120 requests/s рдЪреЗ 10 seconds рдЪрд╛рд▓рд╡рддреЗ, рдорд░реНрдпрд╛рджреЗрд╢рд┐рд╡рд╛рдп рдЖрдгрд┐ 200 рдЪреНрдпрд╛ рдорд░реНрдпрд╛рджреЗрд╕рд╣.

ЁЯзк рдХрд░реВрди рдкрд╛рд╣рд╛

python3 dist/demo.py backpressure
python3 - <<'EOF'
import sys; sys.path.insert(0, "dist"); from sim import BoundedQueue
for rate in (90, 120, 200):
    for limit in (None, 200, 400):
        q = BoundedQueue(limit); w = q.run([rate] * 10, served_per_tick=100)
        print(f"{rate:>3}/s arrive ┬╖ limit {limit or 'none':>4} тЖТ wait after 10 s {w[-1]:>4.1f} s ┬╖ shed {q.shed:>4}")
EOF
python3 dist/demo.py
python3 dist/test_dist.py

тЬЕ рддрдкрд╛рд╕рд╛ тАФ рддреБрдореНрд╣рд╛рд▓рд╛ рдХрд╛рдп рджрд┐рд╕рд╛рдпрд▓рд╛ рд╣рд╡реЗ

backpressure рд╣реЗ print рдХрд░рддреЗ:

тФАтФА 120 requests/s arrive, 100/s can be served ┬╖ queue limit none тЖТ wait after 5 s 1.0 s, after 10 s 2.0 s ┬╖ shed 0
тФАтФА 120 requests/s arrive, 100/s can be served ┬╖ queue limit  200 тЖТ wait after 5 s 1.0 s, after 10 s 1.0 s ┬╖ shed 100

рддреБрдордЪрд╛ snippet рд╣реЗ print рдХрд░рддреЛ:

 90/s arrive ┬╖ limit none тЖТ wait after 10 s  0.0 s ┬╖ shed    0
 90/s arrive ┬╖ limit  200 тЖТ wait after 10 s  0.0 s ┬╖ shed    0
 90/s arrive ┬╖ limit  400 тЖТ wait after 10 s  0.0 s ┬╖ shed    0
120/s arrive ┬╖ limit none тЖТ wait after 10 s  2.0 s ┬╖ shed    0
120/s arrive ┬╖ limit  200 тЖТ wait after 10 s  1.0 s ┬╖ shed  100
120/s arrive ┬╖ limit  400 тЖТ wait after 10 s  2.0 s ┬╖ shed    0
200/s arrive ┬╖ limit none тЖТ wait after 10 s 10.0 s ┬╖ shed    0
200/s arrive ┬╖ limit  200 тЖТ wait after 10 s  1.0 s ┬╖ shed  900
200/s arrive ┬╖ limit  400 тЖТ wait after 10 s  3.0 s ┬╖ shed  700

рдкреВрд░реНрдг demo тЬЕ done тАФ the branches agree рдиреЗ рд╕рдВрдкрддреЛ, рдЖрдгрд┐ tests 12/12 passed рдиреЗ.

ЁЯПБ рддреБрдореНрд╣реА рдЖрддреНрддрд╛рдЪ рдХрд╛рдп рд╕рд┐рджреНрдз рдХреЗрд▓реЗ

рдХреНрд╖рдорддреЗрдЪреНрдпрд╛ рдЦрд╛рд▓реА (90/s) рдорд░реНрдпрд╛рджреЗрд▓рд╛ рдХрд╛рд╣реАрдЪ рдорд╣рддреНрддреНрд╡ рдирд╕рддреЗ. 200/s рд╡рд░ рдорд░реНрдпрд╛рджрд╛ рдирд╕рддрд╛рдирд╛ рдерд╛рдВрдмрдгреЗ 10 seconds рдордзреНрдпреЗ 10 seconds рдкрд░реНрдпрдВрдд рдкреЛрд╣реЛрдЪрд▓реЗ рдЖрдгрд┐ рдЕрдЬреВрдирд╣реА рд╡рд╛рдврдд рд╣реЛрддреЗ; 200 рдЪреНрдпрд╛ рдорд░реНрдпрд╛рджреЗрд╕рд╣ рддреЗ 1 second рд╡рд░рдЪ рд░рд╛рд╣рд┐рд▓реЗ, 900 requests рдЯрд╛рдХреВрди рджреЗрдгреНрдпрд╛рдЪреНрдпрд╛ рдХрд┐рдВрдорддреАрд╡рд░, рдЬреНрдпрд╛рдВрдирд╛ рдЭрдЯрдкрдЯ "рдирдВрддрд░ рдпрд╛" рдорд┐рд│рд╛рд▓реЗ. рдорд░реНрдпрд╛рджрд╛ рдореНрд╣рдгрдЬреЗ рдХрдорд╛рд▓ рдерд╛рдВрдмрдгреНрдпрд╛рдЪреА рдирд┐рд╡рдб: рдЗрдереЗ 400 рд╕реБрдорд╛рд░реЗ 3 seconds рдкрд░реНрдпрдВрдд рдерд╛рдВрдмреВ рджреЗрддреЗ. рддреА user (рдХрд┐рдВрд╡рд╛ caller рдЪрд╛ timeout) рдЦрд░реЛрдЦрд░ рдХрд┐рддреА рд╡реЗрд│ рдерд╛рдВрдмреЗрд▓ рддреНрдпрд╛рд╡рд░реВрди рдирд┐рд╡рдбрд╛.

тЪая╕П рдиреЗрд╣рдореАрдЪреНрдпрд╛ рдЪреБрдХрд╛

ЁЯПн рдкреНрд░рддреНрдпрдХреНрд╖ рд╡рд╛рдкрд░рд╛рдд

Edge рд╡рд░ rate limiting тАФ рдкреНрд░рддреНрдпреЗрдХ client address рд╕рд╛рдареА token-bucket рдкрджреНрдзрддреАрдЪреА рдорд░реНрдпрд╛рджрд╛ рдЕрд╕рд▓реЗрд▓реЗ NGINX, рдорд░реНрдпрд╛рджрд╛ рдУрд▓рд╛рдВрдбрд▓реНрдпрд╛рд╡рд░ 429 рдЙрддреНрддрд░ рджреЗрддреЗ. On a real account:

limit_req_zone $binary_remote_addr zone=per_ip:10m rate=10r/s;
server {
    location /api/ {
        limit_req zone=per_ip burst=20 nodelay;
        limit_req_status 429;
        proxy_pass http://timetable;
    }
}

Load shed рдХрд░рдгрд╛рд░реА service рдХреЗрд╡реНрд╣рд╛ рдкрд░рдд рдпрд╛рдпрдЪреЗ рддреЗ рд╕рд╛рдВрдЧрддреЗ:

HTTP/1.1 503 Service Unavailable
Retry-After: 5

Envoy circuit breaking тАФ рдкреНрд░рддреНрдпреЗрдХ upstream cluster рд╕рд╛рдареА connections, рдерд╛рдВрдмрд▓реЗрд▓реНрдпрд╛ requests, active requests рдЖрдгрд┐ active retries рд╡рд░ рдорд░реНрдпрд╛рджрд╛; рдорд░реНрдпрд╛рджреЗрдкрд▓реАрдХрдбреЗ Envoy request рд░рд╛рдВрдЧреЗрдд рдареЗрд╡рдгреНрдпрд╛рдРрд╡рдЬреА рд▓рдЧреЗрдЪ fail рдХрд░рддреЗ (defaults 1024, 1024, 1024 рдЖрдгрд┐ 3 рдЖрд╣реЗрдд):

clusters:
  - name: timetable
    circuit_breakers:
      thresholds:
        - priority: DEFAULT
          max_connections: 200
          max_pending_requests: 100      # the bounded queue
          max_requests: 400
          max_retries: 3

Code рдордзреНрдпреЗ backoff рдЖрдгрд┐ full jitter рд╕рд╣ retries:

import random, time
def call_with_retries(call, attempts=4, base=0.1, cap=5.0):
    for n in range(attempts):
        try:
            return call()
        except TemporaryError:
            if n == attempts - 1: raise
            time.sleep(random.uniform(0, min(cap, base * 2 ** n)))   # full jitter

AWS SDKs рд╣реЗ рддреБрдордЪреНрдпрд╛рд╕рд╛рдареА рдХрд░рддрд╛рдд (AWS_RETRY_MODE=standard рдХрд┐рдВрд╡рд╛ adaptive); Kubernetes рдЪрд╛ API server API Priority and Fairness рдиреЗ load shed рдХрд░рддреЛ рдЖрдгрд┐ рдПрдЦрд╛рджреНрдпрд╛ priority level рдЪреНрдпрд╛ queues рднрд░рд▓реНрдпрд╛ рдХреА 429 рдЙрддреНрддрд░ рджреЗрддреЛ.

ЁЯПн рдкреНрд░рддреНрдпрдХреНрд╖ рд╡рд╛рдкрд░рд╛рдд рд╣реЗ рдХрд╛ рдорд╣рддреНрддреНрд╡рд╛рдЪреЗ: рдкреНрд░рддреНрдпреЗрдХ service рд╕рд╛рдареА рддрд┐рдЪреНрдпрд╛ рд░рд╛рдВрдЧреЗрдЪреНрдпрд╛ рдорд░реНрдпрд╛рджрд╛, рднрд░рд▓реНрдпрд╛рд╡рд░ рддреА рдХрд╛рдп рдкрд░рдд рдХрд░рддреЗ, рдХреЛрдгрддреЗ рдХрд╛рдо рддреА рдЖрдзреА рдЯрд╛рдХреВрди рджреЗрддреЗ, рдЖрдгрд┐ рдкреНрд░рддреНрдпреЗрдХ caller рдЪреА retry policy рд▓рд┐рд╣реВрди рдареЗрд╡рд╛. рдордЧ рдХреНрд╖рдорддреЗрдкреЗрдХреНрд╖рд╛ рдЬрд╛рд╕реНрдд load-test рдХрд░рд╛ рдЖрдгрд┐ 429/503 рдЪреА рд╕рдВрдЦреНрдпрд╛ рд╡рд╛рдврдд рдЕрд╕рддрд╛рдирд╛ p99 рдерд╛рдВрдмрдгреЗ рд╕рдкрд╛рдЯ рд░рд╛рд╣рддреЗ рдХрд╛ рддреЗ рддрдкрд╛рд╕рд╛.

ЁЯОУ рд╢рд╛рдЦрд╛рдВрдЪреЗ рдПрдХрдордд рдЭрд╛рд▓реЗ

рдирд┐рд░реЛрдкреЗ рд╣рд░рд╡рддрд╛рдд, рдЖрдгрд┐ рд╢рд╛рдВрддрддреЗрддреВрди рдЬрд╡рд│рдЬрд╡рд│ рдХрд╛рд╣реАрдЪ рдХрд│рдд рдирд╛рд╣реА тЖТ рдкреНрд░рддреНрдпреЗрдХ рд╢рд╛рдЦреЗрдЪреЗ рдШрдбреНрдпрд╛рд│ рдереЛрдбреЗрд╕реЗ рдЪреБрдХреАрдЪреЗ рдЕрд╕рддреЗ, рдореНрд╣рдгреВрди counters рдХрд╛рд░рдгрд╛рд▓рд╛ рдкрд░рд┐рдгрд╛рдорд╛рдЪреНрдпрд╛ рдЖрдзреА рдХреНрд░рдо рджреЗрддрд╛рдд тЖТ vector clocks рдЦрд░рд╛ conflict рдЖрдгрд┐ рдЙрд╢реАрд░ рдпрд╛рдВрддрд▓рд╛ рдлрд░рдХ рдУрд│рдЦрддрд╛рдд тЖТ heartbeats рдлрдХреНрдд рд╕рдВрд╢рдп рдШреЗрдК рд╢рдХрддрд╛рдд тЖТ рдиреЛрдВрджрд╡рд╣реАрдЪреНрдпрд╛ рдкреНрд░рддреАрдВрдирд╛ рдЬреЗ рдорд┐рд│рд╛рд▓реЗрдЪ рдирд╡реНрд╣рддреЗ рддреЗ рд╣рд░рд╡рддреЗ тЖТ quorums рд╕рд░реНрд╡рд╛рдд рдирд╡реАрди рд╢реЛрдзрдгреНрдпрд╛рд╕рд╛рдареА рдПрдХрдореЗрдХрд╛рдВрд╡рд░ рдпреЗрддрд╛рдд тЖТ рдкреНрд░рддреНрдпреЗрдХ рд╡рд╛рдЪрдгрд╛рд▒реНрдпрд╛рд▓рд╛ рдирд╛рд╡ рдЕрд╕рд▓реЗрд▓реЗ рд╡рдЪрди рдорд┐рд│рддреЗ тЖТ рддреБрдЯрд▓реЗрд▓рд╛ рд░рд╕реНрддрд╛ "рдирдХрд╛рд░ рджреНрдпрд╛ рдХрд┐рдВрд╡рд╛ рдЬреБрдиреЗ рдЙрддреНрддрд░ рджреНрдпрд╛" рд╣реА рдирд┐рд╡рдб рд▓рд╛рджрддреЛ тЖТ рдмрд╣реБрдордд рдкреНрд░рддреНрдпреЗрдХ term рдордзреНрдпреЗ рдПрдХ leader рдирд┐рд╡рдбрддреЗ тЖТ lease рд▓рд╛ fencing token рд▓рд╛рдЧрддреЛ тЖТ at-least-once рдЕрдзрд┐рдХ key рдореНрд╣рдгрдЬреЗ effectively once тЖТ рдЖрдгрд┐ bounded рд░рд╛рдВрдЧ рд▓рд╡рдХрд░ "рдирдВрддрд░ рдпрд╛" рд╕рд╛рдВрдЧрддреЗ. рддреБрдореНрд╣реА рдлрдХреНрдд distributed systems рд╢рд┐рдХрд▓рд╛ рдирд╛рд╣реАрдд тАФ рддреБрдореНрд╣реА рджреЛрди boxes рдордзрд▓рд╛ рдХреЛрдгрддрд╛рд╣реА рдмрд╛рдг рдкрд╛рд╣реВрди рдпреЛрдЧреНрдп рдкреНрд░рд╢реНрди рд╡рд┐рдЪрд╛рд░реВ рд╢рдХрддрд╛: рд╣реА рдЪрд┐рдареНрдареА рд╣рд░рд╡рд▓реА, рдЙрд╢рд┐рд░рд╛ рдЖрд▓реА, рдкреБрдиреНрд╣рд╛ рдЖрд▓реА, рдХрд┐рдВрд╡рд╛ рд░рд╕реНрддрд╛ рддреБрдЯрд▓рд╛ рддрд░ рдХрд╛рдп? ЁЯМРЁЯССЁЯОУ

тПня╕П рдкреБрдвреЗ

рдЗрддрд░ рд╢рд╛рд│рд╛. System Design рд╢рд╛рд│рд╛ рд╣реЗ рд╕рдЧрд│реЗ рдЦрд▒реНрдпрд╛ designs рдордзреНрдпреЗ рд╡рд╛рдкрд░рддреЗ; Scaling рд╢рд╛рд│рд╛ traffic рд╣рд╛рддрд╛рд│рддреЗ; Database рд╢рд╛рд│рд╛ рдиреЛрдВрджрд╡рд╣реАрдд рдЖрдгрдЦреА рдЦреЛрд▓рд╡рд░ рдЬрд╛рддреЗ; Kubernetes рд╢рд╛рд│рд╛ etcd рдЖрдгрд┐ leases рд╡рд░ рдЪрд╛рд▓рддреЗ; School portal рдордзреНрдпреЗ рдмрд╛рдХреА рд╕рдЧрд│реЗ рдЖрд╣реЗ.

git checkout main
python3 dist/demo.py     # one last run, for fun

ЁЯЪз Lesson 12 тАФ Backpressure & the whole picture: say "try later" early

ЁЯУН You are here: Lesson 12 of 12 ┬╖ Previous: lesson-11-exactly-once ┬╖ The last lesson ЁЯОУ


ЁЯУж What's in this branch

Lessons 01тАУ11, plus backpressure: when more work arrives than a node can do, an unbounded queue turns the overload into ever-growing waiting; a bounded queue refuses early with "try later" (429 / 503 + Retry-After) and keeps the wait short. Also load shedding, retries with backoff and jitter, rate limiting and circuit breaking тАФ and the whole map of the course. backpressure() in dist/demo.py and BoundedQueue in dist/sim.py.

ЁЯзТ Explain like I'm 5

On admission day, 120 parents a second come to the Pune office. The clerks can help 100 a second. Every second, 20 more parents join the line.

Turning some people away sounds unkind. But the other choice is that everyone waits too long, and the office stops working for all of them.

And one more rule for the parents sent away: do not all come back at the same moment. Wait a bit, then a bit longer, and add a small random wait, so the next crowd is spread out.

ЁЯЧ║я╕П Diagram

flowchart LR
    in["ЁЯСк 120 requests/s arrive"] --> gate{"ЁЯЪз queue full?<br/>limit 200"}
    gate -->|"yes"| shed["тЫФ 429 / 503 + Retry-After<br/>shed 100 in 10 s"]
    gate -->|"no"| q["ЁЯУЛ bounded queue"]
    q --> srv["ЁЯзСтАНЁЯТ╝ serve 100/s<br/>wait stays about 1.0 s"]
    un["unbounded queue: wait 1.0 s after 5 s,<br/>2.0 s after 10 s, and still growing"]

ЁЯЧ║я╕П Drawn version + a lab: https://school-edh.pages.dev/distributed-systems/lesson-diagrams.html#l12

тЭУ What

ЁЯдФ Why

Because every lesson before this one creates extra work under stress: retries after timeouts (01, 11), failovers (05), repairs (06), elections (09). A system with no limits turns a small overload into a total one тАФ the queue grows, every request times out, every client retries, and the load doubles. Saying "try later" early is what keeps the rest of the system working.

ЁЯФз How (in this repo)

BoundedQueue(limit) in dist/sim.py runs one tick per second: run(arrivals, served_per_tick) takes the new arrivals up to the free room (limit - q; limit=None means no limit), counts the rest in shed, serves up to served_per_tick, and records the wait as queue ├╖ served_per_tick seconds. backpressure() in dist/demo.py runs 10 seconds of 120 requests/s against 100/s of service, with no limit and with a limit of 200.

ЁЯзк Try it

python3 dist/demo.py backpressure
python3 - <<'EOF'
import sys; sys.path.insert(0, "dist"); from sim import BoundedQueue
for rate in (90, 120, 200):
    for limit in (None, 200, 400):
        q = BoundedQueue(limit); w = q.run([rate] * 10, served_per_tick=100)
        print(f"{rate:>3}/s arrive ┬╖ limit {limit or 'none':>4} тЖТ wait after 10 s {w[-1]:>4.1f} s ┬╖ shed {q.shed:>4}")
EOF
python3 dist/demo.py
python3 dist/test_dist.py

тЬЕ Verify тАФ what you should see

backpressure prints:

тФАтФА 120 requests/s arrive, 100/s can be served ┬╖ queue limit none тЖТ wait after 5 s 1.0 s, after 10 s 2.0 s ┬╖ shed 0
тФАтФА 120 requests/s arrive, 100/s can be served ┬╖ queue limit  200 тЖТ wait after 5 s 1.0 s, after 10 s 1.0 s ┬╖ shed 100

Your snippet prints:

 90/s arrive ┬╖ limit none тЖТ wait after 10 s  0.0 s ┬╖ shed    0
 90/s arrive ┬╖ limit  200 тЖТ wait after 10 s  0.0 s ┬╖ shed    0
 90/s arrive ┬╖ limit  400 тЖТ wait after 10 s  0.0 s ┬╖ shed    0
120/s arrive ┬╖ limit none тЖТ wait after 10 s  2.0 s ┬╖ shed    0
120/s arrive ┬╖ limit  200 тЖТ wait after 10 s  1.0 s ┬╖ shed  100
120/s arrive ┬╖ limit  400 тЖТ wait after 10 s  2.0 s ┬╖ shed    0
200/s arrive ┬╖ limit none тЖТ wait after 10 s 10.0 s ┬╖ shed    0
200/s arrive ┬╖ limit  200 тЖТ wait after 10 s  1.0 s ┬╖ shed  900
200/s arrive ┬╖ limit  400 тЖТ wait after 10 s  3.0 s ┬╖ shed  700

The full demo ends with тЬЕ done тАФ the branches agree, and the tests with 12/12 passed.

ЁЯПБ What you just proved

Below capacity (90/s), the limit never matters. At 200/s with no limit, the wait reached 10 seconds in 10 seconds and was still growing; with a limit of 200 it stayed at 1 second, at the price of shedding 900 requests that got a fast "try later". The limit is a choice of maximum wait: 400 allows up to about 3 seconds here. Pick it from how long a user (or the caller's timeout) will really wait.

тЪая╕П Common mistakes

ЁЯПн In production

Rate limiting at the edge тАФ NGINX with a token-bucket style limit per client address, answering 429 when it is exceeded. On a real account:

limit_req_zone $binary_remote_addr zone=per_ip:10m rate=10r/s;
server {
    location /api/ {
        limit_req zone=per_ip burst=20 nodelay;
        limit_req_status 429;
        proxy_pass http://timetable;
    }
}

A service that sheds load says when to come back:

HTTP/1.1 503 Service Unavailable
Retry-After: 5

Envoy circuit breaking тАФ per upstream cluster, caps on connections, waiting requests, active requests and active retries; beyond a cap, Envoy fails the request at once instead of queueing it (the defaults are 1024, 1024, 1024 and 3):

clusters:
  - name: timetable
    circuit_breakers:
      thresholds:
        - priority: DEFAULT
          max_connections: 200
          max_pending_requests: 100      # the bounded queue
          max_requests: 400
          max_retries: 3

Retries with backoff and full jitter in code:

import random, time
def call_with_retries(call, attempts=4, base=0.1, cap=5.0):
    for n in range(attempts):
        try:
            return call()
        except TemporaryError:
            if n == attempts - 1: raise
            time.sleep(random.uniform(0, min(cap, base * 2 ** n)))   # full jitter

The AWS SDKs do this for you (AWS_RETRY_MODE=standard or adaptive); Kubernetes' API server sheds load with API Priority and Fairness and answers 429 when a priority level's queues are full.

ЁЯПн Why this matters in production: for each service, write down its queue limits, what it returns when full, which work it sheds first, and the retry policy of every caller. Then load-test above capacity and check that the p99 wait stays flat while the 429/503 count rises.

ЁЯОУ The branches agree

Messengers get lost, and silence tells you almost nothing тЖТ every branch's clock is a little wrong, so counters order cause before effect тЖТ vector clocks tell a real conflict from a delay тЖТ heartbeats can only suspect тЖТ copies of the register lose what they had not received тЖТ quorums overlap to find the newest тЖТ each reader gets a named promise тЖТ a cut road forces "refuse or answer stale" тЖТ a majority elects one leader per term тЖТ a lease needs a fencing token тЖТ at-least-once plus a key is effectively once тЖТ and a bounded queue says "try later" early. You didn't just learn distributed systems тАФ you can look at any arrow between two boxes and ask the right questions: what if this note is lost, late, repeated, or the road is cut? ЁЯМРЁЯССЁЯОУ

тПня╕П Next

The other schools. The System Design school uses all of this in real designs; the Scaling school handles the traffic; the Database school goes deeper on the register; the Kubernetes school runs on etcd and leases; the School portal has the rest.

git checkout main
python3 dist/demo.py     # one last run, for fun
тЖР Previousexactly onceFinished! Take the quiz тЖТcheck what stuck

This page is the lesson's README from the lesson-12-backpressure branch, shown here so the whole School stays on one site. Code files open on GitHub at the same branch.