ЁЯПл The SchoolтА║ЁЯУИ ScalingтА║ЁЯУм рдзрдбрд╛ 12 тАФ Queues + рд╕рдВрдкреВрд░реНрдг рдЖрд░рд╛рдЦрдбрд╛: рдЧрд░реНрджреА gate рд╡рд░рдЪ рдЭреЗрд▓рд╛
ЁЯЦ╝я╕П See the drawing + lab ЁЯПа Course home ЁЯМ┐ Branch on GitHub тЬПя╕П View source
ЁЯЦ╝я╕П рдЖрдХреГрддреА рдЖрдгрд┐ labThe drawing + lab рдкреВрд░реНрдг рдкрд╛рдирд╛рд╡рд░ рдЙрдШрдбрд╛ тЖЧOpen full page тЖЧ

ЁЯУм рдзрдбрд╛ 12 тАФ Queues + рд╕рдВрдкреВрд░реНрдг рдЖрд░рд╛рдЦрдбрд╛: рдЧрд░реНрджреА gate рд╡рд░рдЪ рдЭреЗрд▓рд╛

ЁЯУН рддреБрдореНрд╣реА рдЗрдереЗ рдЖрд╣рд╛рдд: 13 рдкреИрдХреА рдзрдбрд╛ 12 ┬╖ рдорд╛рдЧреАрд▓: lesson-11-partitioning ┬╖ рдкреБрдвреАрд▓: lesson-13-surviving-failure


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

рдзрдбреЗ 01тАУ11, рдЖрдгрд┐ рд╢реЗрд╡рдЯрдЪреЗ рд╕рд╛рдзрди: request рдЖрдгрд┐ рд╣рд│реВ рдХрд╛рдо рдпрд╛рдВрдЪреНрдпрд╛рдордзреНрдпреЗ рдПрдХ queue (Amazon SQS), backlog рдиреБрд╕рд╛рд░ scale рд╣реЛрдгрд╛рд▒реНрдпрд╛ workers рд╕рд╣ тАФ рдЖрдгрд┐ рдордЧ scale рдХреЗрд▓реЗрд▓реНрдпрд╛ web app рдЪрд╛ рд╕рдВрдкреВрд░реНрдг рдЖрд░рд╛рдЦрдбрд╛. scale/demo.py рдордзрд▓реЗ queues() рдкреНрд░рдорд╛рдгрдкрддреНрд░рд╛рдВрдЪреНрдпрд╛ requests рдЪреА рдПрдХ рдЧрд░реНрджреА queue рдордзреВрди рдкрд╛рдард╡рддреЗ, рдЖрдгрд┐ рдЕрдкрдпрд╢реА рдХрд╛рдо рд╕реБрд░рдХреНрд╖рд┐рддрдкрдгреЗ рдкрд░рдд рдорд┐рд│рд╡рддрд╛ рдпреЗрдгреНрдпрд╛рд╕рд╛рдареА рдХрд╛рдп рд▓рд╛рдЧрддреЗ рддреЗ рд╕рд╛рдВрдЧрддреЗ: at-least-once delivery, visibility timeout, DLQ рдЖрдгрд┐ idempotent workers.

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

рдирд┐рдХрд╛рд▓рд╛рдирдВрддрд░ рдкреНрд░рддреНрдпреЗрдХ рдкрд╛рд▓рдХрд╛рд▓рд╛ рдЫрд╛рдкрд▓реЗрд▓реЗ рдкреНрд░рдорд╛рдгрдкрддреНрд░ рд╣рд╡реЗ рдЕрд╕рддреЗ. рдЫрдкрд╛рдИрд▓рд╛ рд╡реЗрд│ рд▓рд╛рдЧрддреЛ. рдЬрд░ рдкреНрд░рддреНрдпреЗрдХ рдкрд╛рд▓рдХ рдЫрдкрд╛рдИ рд╣реЛрдИрдкрд░реНрдпрдВрдд рдЦрд┐рдбрдХреАрд╡рд░ рдерд╛рдВрдмрд▓рд╛, рддрд░ рд░рд╛рдВрдЧ рдХрдзреАрдЪ рдкреБрдвреЗ рд╕рд░рдХрдд рдирд╛рд╣реА.

рдореНрд╣рдгреВрди рджреАрдкрд┐рдХрд╛ tokens ЁЯОЯя╕П рд╡рд╛рдкрд░рддреЗ. рдкрд╛рд▓рдХ рдкреНрд░рдорд╛рдгрдкрддреНрд░ рдорд╛рдЧрддреЛ, рддреНрдпрд╛рд▓рд╛ рд▓рдЧреЗрдЪ token рдорд┐рд│рддреЛ, рдЖрдгрд┐ рддреЛ рдЬрддреНрд░реЗрдЪреА рдордЬрд╛ рдШреНрдпрд╛рдпрд▓рд╛ рдЬрд╛рддреЛ. Token рдПрдХрд╛ рдкреЗрдЯреАрдд рдЬрд╛рддреЛ. рдорд╛рдЧрдЪреНрдпрд╛ рдЦреЛрд▓реАрдд рдЫрдкрд╛рдИ рдХрд░рдгрд╛рд░реЗ рдкреЗрдЯреАрддреВрди рдПрдХреЗрдХ token рдШреЗрддрд╛рдд рдЖрдгрд┐ рдЫрд╛рдкрддрд╛рдд.

рддреАрди рд╕реЗрдХрдВрдж рджрд░ рд╕реЗрдХрдВрджрд╛рд▓рд╛ 800 рдкрд╛рд▓рдХ рдорд╛рдЧрдгреА рдХрд░рддрд╛рдд. 3 рдЫрдкрд╛рдИ рдХрд░рдгрд╛рд▒реНрдпрд╛рдВрд╕рд╣ рдкреЗрдЯреАрдд 1,500 tokens рдЬрдорддрд╛рдд рдЖрдгрд┐ рджрд╣рд╛ рд╕реЗрдХрдВрджрд╛рдВрдирдВрддрд░рд╣реА рддреА рд░рд┐рдХрд╛рдореА рд╣реЛрдд рдирд╛рд╣реА. рдореНрд╣рдгреВрди рджреАрдкрд┐рдХрд╛ рдПрдХ рдирд┐рдпрдо рдЬреЛрдбрддреЗ: "рдкреЗрдЯреАрддрд▓реНрдпрд╛ рдкреНрд░рддреНрдпреЗрдХ 100 tokens рдорд╛рдЧреЗ рдПрдХ рдЫрдкрд╛рдИ рдХрд░рдгрд╛рд░рд╛, рдЬрд╛рд╕реНрддреАрдд рдЬрд╛рд╕реНрдд 10". рдкреЗрдЯреА рднрд░рд▓реА рдХреА рдЫрдкрд╛рдИ рдХрд░рдгрд╛рд░реЗ рд╡рд╛рдврд╡рд▓реЗ рдЬрд╛рддрд╛рдд; рд░рд┐рдХрд╛рдореА рдЭрд╛рд▓реА рдХреА рддреЗ рдШрд░реА рдЬрд╛рддрд╛рдд. рдкреЗрдЯреА рдЧрд░реНрджреА рдЭреЗрд▓рддреЗ: requests рдЦрд┐рдбрдХреАрд╡рд░ рдЧрд░реНрджреА рдХрд░рдгреНрдпрд╛рдРрд╡рдЬреА рдкреЗрдЯреАрдд рдерд╛рдВрдмрддрд╛рдд.

рдкреНрд░рдорд╛рдгрдкрддреНрд░ рдЫрд╛рдкрддрд╛рдирд╛ рдордзреНрдпреЗрдЪ рдЫрдкрд╛рдИ рдЕрдбрдХрд▓реА, рддрд░ рдХрд╛рд╣реА рд╡реЗрд│рд╛рдиреЗ token рдкреЗрдЯреАрдд рдкрд░рдд рдЬрд╛рддреЛ, рдЖрдгрд┐ рджреБрд╕рд░рд╛ рдЫрдкрд╛рдИ рдХрд░рдгрд╛рд░рд╛ рдкреБрдиреНрд╣рд╛ рдкреНрд░рдпрддреНрди рдХрд░рддреЛ. рдХрдзреА рдХрдзреА рддреЛрдЪ token рджреЛрдирджрд╛ рджрд┐рд▓рд╛ рдЬрд╛рддреЛ, рдореНрд╣рдгреВрди рдкреНрд░рддреНрдпреЗрдХ рдЫрдкрд╛рдИ рдХрд░рдгрд╛рд░рд╛ рдЖрдзреА рддрдкрд╛рд╕рддреЛ "рд╣реЗ рдкреНрд░рдорд╛рдгрдкрддреНрд░ рдЖрдзреАрдЪ рдЫрд╛рдкрд▓реЗ рдЖрд╣реЗ рдХрд╛?". рдкреБрдиреНрд╣рд╛ рдкреБрдиреНрд╣рд╛ рдЕрдкрдпрд╢реА рдард░рдгрд╛рд░рд╛ token рдПрдХрд╛ рдмрд╛рдЬреВрдЪреНрдпрд╛ рдкреЗрдЯреАрдд рдЬрд╛рддреЛ, рдореНрд╣рдгрдЬреЗ рдХреЛрдгреАрддрд░реА рддреЛ рдкрд╛рд╣реВ рд╢рдХреЗрд▓. Tokens рдкреЗрдЯреАрдд рдХрд╛рдпрдо рд░рд╛рд╣рдд рдирд╛рд╣реАрдд: рдХрд╛рд╣реА рджрд┐рд╡рд╕рд╛рдВрдирдВрддрд░ рдЬреБрдирд╛ token рдлреЗрдХреВрди рджрд┐рд▓рд╛ рдЬрд╛рддреЛ.

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

flowchart LR
    u["ЁЯСк parents"] --> api["ЁЯН│ API<br/>202 Accepted + token"]
    api -->|"SendMessage"| q["ЁЯУм SQS queue<br/>backlog"]
    q -->|"ReceiveMessage"| w["ЁЯЦия╕П workers<br/>1 per 100 waiting, max 10"]
    w -->|"DeleteMessage when done"| q
    q -.->|"failed too many times"| dlq["ЁЯзп dead-letter queue"]
    m["ЁЯУК ApproximateNumberOfMessagesVisible"] -.->|"scale workers"| w

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

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

ЁЯдФ рдХрд╛

рдХрд╛рд░рдг рдХрд╛рд╣реА рдХрд╛рдо рдПрдЦрд╛рджреНрдпрд╛ рдкрд╛рдирд╛рд▓рд╛ рд▓рд╛рдЧрд╛рд╡рд╛ рддреНрдпрд╛рдкреЗрдХреНрд╖рд╛ рд╣рд│реВ рдЕрд╕рддреЗ, рдЖрдгрд┐ рдЧрд░реНрджреА servers рдЬреЛрдбрд╛рдпрд▓рд╛ рд▓рд╛рдЧрдгрд╛рд▒реНрдпрд╛ рд╡реЗрд│реЗрдкреЗрдХреНрд╖рд╛ рдХрдореА рдХрд╛рд│ рдЯрд┐рдХрддреЗ. Queue "3 рд╕реЗрдХрдВрджрд╛рдВрдд 2,400 requests" рдЪреЗ "рдкреБрдврдЪреНрдпрд╛ рдХрд╛рд╣реА рд╕реЗрдХрдВрджрд╛рдВрдд рдкреВрд░реНрдг рдЭрд╛рд▓реЗрд▓реА 2,400 рдХрд╛рдореЗ" рдХрд░рддреЗ. рдкреБрдврдЪрд╛ рднрд╛рдЧ рдЬрд▓рдж рд░рд╛рд╣рддреЛ, рдЖрдгрд┐ рдорд╛рдЧрдЪрд╛ рднрд╛рдЧ рддреНрдпрд╛рд▓рд╛ рдЭреЗрдкреЗрд▓ рдЕрд╢рд╛ рд╡реЗрдЧрд╛рдиреЗ рдЪрд╛рд▓рддреЛ. рд╣реА рдЬрд╛рджреВ рдирд╛рд╣реА: message рддрд░реАрд╣реА expire рд╣реЛрдК рд╢рдХрддреЛ, рдХрд┐рдВрд╡рд╛ рдкреНрд░рддреНрдпреЗрдХ рд╡реЗрд│реА рдЕрдкрдпрд╢реА рдард░реВ рд╢рдХрддреЛ. Retries, DLQ рдЖрдгрд┐ idempotent workers рдЕрд╕рддреАрд▓ рддрд░ рддреЗ рдЕрдкрдпрд╢реА рдХрд╛рдо crash рдордзреНрдпреЗ рд╣рд░рд╡рдгреНрдпрд╛рдРрд╡рдЬреА рдЬрдкрд▓реЗ рдЬрд╛рддреЗ рдЖрдгрд┐ рд╕реБрд░рдХреНрд╖рд┐рддрдкрдгреЗ рдкрд░рдд рдорд┐рд│рд╡рддрд╛ рдпреЗрддреЗ.

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

scale/sim.py рдордзрд▓реЗ drain(arrivals, per_worker, workers_for) рдПрдХреЗрдХ рд╕реЗрдХрдВрдж рдЪрд╛рд▓рд╡рддреЗ: рддреЗ workers_for(backlog) рд▓рд╛ рдХрд┐рддреА workers рдЪрд╛рд▓рд╡рд╛рдпрдЪреЗ рддреЗ рд╡рд┐рдЪрд╛рд░рддреЗ тАФ рдпрд╛ рд╕реЗрдХрдВрджрд╛рдЪреЗ arrivals рдпреЗрдгреНрдпрд╛рдЖрдзреАрдЪреНрдпрд╛ backlog рд╡рд░реВрди рдард░рд╡рд▓реЗрд▓реЗ, рдЬрд╕реЗ рдЦрд▒реНрдпрд╛ autoscaler рд▓рд╛ metric рдереЛрдбрд╛ рдЙрд╢рд┐рд░рд╛ рджрд┐рд╕рддреЛ тАФ arrivals рдЬреЛрдбрддреЗ, рдЖрдгрд┐ рдкреНрд░рддреНрдпреЗрдХ worker рд▓рд╛ per_worker рдХрд╛рдореЗ рдкреВрд░реНрдг рдХрд░реВ рджреЗрддреЗ. рдкреНрд░рддреНрдпреЗрдХ row рдореНрд╣рдгрдЬреЗ (second, arrived, workers, done, backlog). queues() рдард░рд▓реЗрд▓реНрдпрд╛ 3 workers рдЪреА рддреБрд▓рдирд╛ "рдкреНрд░рддреНрдпреЗрдХ 100 рдерд╛рдВрдмрд▓реЗрд▓реНрдпрд╛рдВрдорд╛рдЧреЗ рдПрдХ worker, рдХрд┐рдорд╛рди 1, рдЬрд╛рд╕реНрддреАрдд рдЬрд╛рд╕реНрдд 10" рд╢реА рдХрд░рддреЗ.

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

python3 scale/demo.py queues
python3 - <<'EOF'
import sys; sys.path.insert(0, "scale"); from sim import drain
arrivals = [50, 800, 800, 800, 200, 50, 50, 50, 50, 50]
for name, rule in (("fixed 3", lambda b: 3), ("fixed 8", lambda b: 8), ("1 per 100 waiting, max 10", lambda b: min(10, max(1, -(-b // 100)))),
                   ("1 per 100 waiting, max 20", lambda b: min(20, max(1, -(-b // 100)))), ("1 per 50 waiting, max 20", lambda b: min(20, max(1, -(-b // 50))))):
    rows = drain(arrivals, 100, rule)
    print(f"{name:<26} peak backlog {max(r[4] for r in rows):>5} ┬╖ most workers {max(r[2] for r in rows):>2} ┬╖ left after 10 s {rows[-1][4]:>4}")
for s, a, w, d, b in drain(arrivals, 100, lambda b: min(10, max(1, -(-b // 100)))):
    print(f"second {s}: arrived {a:>3}, workers {w:>2}, done {d:>4}, waiting {b:>4}")
EOF
python3 scale/test_scale.py

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

queues рд╣реЗ рдЫрд╛рдкрддреЗ:

   fixed 3 workers                                      peak backlog  1500 ┬╖ left after 10 s  150
   scale on backlog (1 worker per 100 waiting, max 10)  peak backlog   800 ┬╖ left after 10 s    0
   the queue (SQS) absorbs the burst; workers process messages asynchronously, at their own pace
   SQS delivers at least once: a visibility timeout, retries, a DLQ and idempotent workers make failed work safe to recover

рддреБрдордЪрд╛ snippet рд╣реЗ рдЫрд╛рдкрддреЛ:

fixed 3                    peak backlog  1500 ┬╖ most workers  3 ┬╖ left after 10 s  150
fixed 8                    peak backlog     0 ┬╖ most workers  8 ┬╖ left after 10 s    0
1 per 100 waiting, max 10  peak backlog   800 ┬╖ most workers  8 ┬╖ left after 10 s    0
1 per 100 waiting, max 20  peak backlog   800 ┬╖ most workers  8 ┬╖ left after 10 s    0
1 per 50 waiting, max 20   peak backlog   700 ┬╖ most workers 14 ┬╖ left after 10 s    0
second 0: arrived  50, workers  1, done   50, waiting    0
second 1: arrived 800, workers  1, done  100, waiting  700
second 2: arrived 800, workers  7, done  700, waiting  800
second 3: arrived 800, workers  8, done  800, waiting  800
second 4: arrived 200, workers  8, done  800, waiting  200
second 5: arrived  50, workers  2, done  200, waiting   50
second 6: arrived  50, workers  1, done  100, waiting    0
second 7: arrived  50, workers  1, done   50, waiting    0
second 8: arrived  50, workers  1, done   50, waiting    0
second 9: arrived  50, workers  1, done   50, waiting    0

Tests тЬЕ L12 a queue absorbs the burst and drains to zero рдЖрдгрд┐ 13/13 passed рдиреЗ рд╕рдВрдкрддрд╛рдд.

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

Model рдордзреНрдпреЗ рдкреНрд░рддреНрдпреЗрдХ request queue рдордзреНрдпреЗ рдерд╛рдВрдмрд▓реА рдЖрдгрд┐ рдирдВрддрд░ рдкреВрд░реНрдг рдЭрд╛рд▓реА. Backlog рдиреБрд╕рд╛рд░ scaling рдиреЗ рдЧрд░реНрджреА рд╕реЗрдХрдВрдж 6 рдкрд░реНрдпрдВрдд рд╕рдВрдкрд╡рд▓реА рдЖрдгрд┐ рдкрд░рдд 1 worker рд╡рд░ рдЖрд▓реЗ; рдард░рд▓реЗрд▓реЗ 3 workers 10 рд╕реЗрдХрдВрджрд╛рдВрдирдВрддрд░рд╣реА 150 рдиреЗ рдорд╛рдЧреЗ рд╣реЛрддреЗ. рдард░рд▓реЗрд▓реЗ 8 workers backlog рддрдпрд╛рд░рдЪ рд╣реЛрдК рджреЗрдд рдирд╛рд╣реАрдд тАФ рдкрдг рджрд┐рд╡рд╕рднрд░ 8 workers рдЪреЗ рдкреИрд╕реЗ рдореЛрдЬрддрд╛рдд. рдХрдорд╛рд▓ 10 рд╡рд░реВрди 20 рдХреЗрд▓реНрдпрд╛рдиреЗ рдХрд╛рд╣реАрдЪ рдмрджрд▓рд▓реЗ рдирд╛рд╣реА, рдХрд╛рд░рдг рдирд┐рдпрдорд╛рдиреЗ рдХрдзреАрдЪ 8 рдкреЗрдХреНрд╖рд╛ рдЬрд╛рд╕реНрдд рдорд╛рдЧрд┐рддрд▓реЗ рдирд╛рд╣реАрдд; рдЬрд╛рд╕реНрдд рдЙрддреНрд╕реБрдХ рдирд┐рдпрдорд╛рдиреЗ (рдкреНрд░рддреНрдпреЗрдХ 50 рдорд╛рдЧреЗ 1) 14 рдкрд░реНрдпрдВрдд workers рд╡рд╛рдкрд░рд▓реЗ рдЖрдгрд┐ peak рдереЛрдбреЗ рдХрдореА рдареЗрд╡рд▓реЗ. рдЖрдгрд┐ рдЧрд░реНрджреАрдЪрд╛ рдкрд╣рд┐рд▓рд╛ рд╕реЗрдХрдВрдж рдиреЗрд╣рдореАрдЪ рд╣рд│реВ рдЕрд╕рддреЛ: autoscaler рд▓рд╛ backlog рддреЛ рддрдпрд╛рд░ рдЭрд╛рд▓реНрдпрд╛рдирдВрддрд░рдЪ рджрд┐рд╕рддреЛ.

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

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

рдЦрд▒реНрдпрд╛ account рд╡рд░ тАФ dead-letter queue рд╕рд╣ рдПрдХ queue (queue messages 4 рджрд┐рд╡рд╕ рдареЗрд╡рддреЗ, DLQ 14 рджрд┐рд╡рд╕; 5 receives рдирдВрддрд░ message DLQ рдордзреНрдпреЗ рдЬрд╛рддреЛ):

aws sqs create-queue --queue-name certificates-dlq --attributes '{"MessageRetentionPeriod": "1209600"}'
aws sqs create-queue --queue-name certificates --attributes '{
  "VisibilityTimeout": "120",
  "MessageRetentionPeriod": "345600",
  "RedrivePolicy": "{\"deadLetterTargetArn\":\"arn:aws:sqs:ap-south-1:111122223333:certificates-dlq\",\"maxReceiveCount\":\"5\"}"
}'
aws sqs get-queue-attributes --attribute-names ApproximateNumberOfMessagesVisible \
    --queue-url https://sqs.ap-south-1.amazonaws.com/111122223333/certificates

Kubernetes рд╡рд░ KEDA тАФ рдкреНрд░рддреНрдпреЗрдХ 100 рдерд╛рдВрдмрд▓реЗрд▓реНрдпрд╛ messages рдорд╛рдЧреЗ рд╕реБрдорд╛рд░реЗ рдПрдХ worker pod, 1 рддреЗ 10:

apiVersion: keda.sh/v1alpha1
kind: ScaledObject
metadata: { name: certificate-workers }
spec:
  scaleTargetRef: { name: certificate-worker }
  minReplicaCount: 1
  maxReplicaCount: 10
  triggers:
    - type: aws-sqs-queue
      authenticationRef: { name: keda-aws }
      metadata:
        queueURL: https://sqs.ap-south-1.amazonaws.com/111122223333/certificates
        queueLength: "100"
        awsRegion: ap-south-1

ЁЯПЧя╕П рд╕рдВрдкреВрд░реНрдг рдЖрд░рд╛рдЦрдбрд╛ тАФ scale рдХреЗрд▓реЗрд▓реА рдЬрддреНрд░рд╛

flowchart LR
    u["ЁЯСк parents"] --> cf["ЁЯМН CloudFront<br/>UI from S3 ┬╖ micro-cache ┬╖ shield"]
    cf -->|"/api"| lb["тЪЦя╕П ALB or API Gateway"]
    lb --> api["ЁЯН│ stateless API<br/>ASG ┬╖ pods + HPA/Karpenter ┬╖ Lambda"]
    api --> rc["ЁЯУМ Redis / Valkey"]
    api --> px["ЁЯЫОя╕П RDS Proxy"]
    px --> pri["ЁЯУТ primary"]
    px --> rep["ЁЯУЪ replicas"]
    api --> ddb["ЁЯЧВя╕П DynamoDB<br/>keys that spread"]
    api --> q["ЁЯУм SQS"] --> wk["ЁЯЦия╕П workers<br/>scale on backlog"]
рд╕реНрддрд░ рдХрд╢рд╛рдиреЗ scale рд╣реЛрддреЛ рдХрд╛рдп рдкрд╛рд╣рд╛рдпрдЪреЗ
UI CloudFront edges, versioned files, рдЧрд░реНрджреАрдЪреНрдпрд╛ рдкрд╛рдирд╛рдВрд╡рд░ s-maxage cache hit rate, origin requests
API load balancer рдорд╛рдЧреЗ stateless рдкреНрд░рддреА: ASG, HPA + Karpenter, рдХрд┐рдВрд╡рд╛ Lambda p99 latency, 5xx, concurrency, Pending pods
Reads Redis cache-aside, read replicas hit rate, ReplicaLag
Connections RDS Proxy / рдПрдХ pool DatabaseConnections рд╡рд┐рд░реБрджреНрдз max_connections
Writes рдкрд╕рд░рдгрд╛рд▒реНрдпрд╛ key рд╕рд╣ partitions throttled requests, hot keys
рд╣рд│реВ рдХрд╛рдо SQS + backlog рдиреБрд╕рд╛рд░ workers ApproximateNumberOfMessagesVisible, DLQ depth

ЁЯПн рдкреНрд░рддреНрдпрдХреНрд╖ рд╡рд╛рдкрд░рд╛рдд рд╣реЗ рдХрд╛ рдорд╣рддреНрддреНрд╡рд╛рдЪреЗ: рдкреНрд░рддреНрдпреЗрдХ рд╕реНрддрд░ рд╕реНрд╡рддрдГрдЪреНрдпрд╛ рд╕рдВрдХреЗрддрд╛рдиреБрд╕рд╛рд░ scale рд╣реЛрддреЛ. рдореЛрдареНрдпрд╛ рджрд┐рд╡рд╕рд╛рдЖрдзреА рд╕рдВрдкреВрд░реНрдг рдорд╛рд░реНрдЧрд╛рдЪрд╛ load-test рдХрд░рд╛ (рдзрдбрд╛ 02), рдорд╛рд╣реАрдд рдЕрд╕рд▓реЗрд▓рд╛ spike schedule рдХрд░рд╛ (рдзрдбреЗ 06, 08), рдЖрдгрд┐ API рд╕реНрддрд░ рдЖрддрд╛ рдЬреЗ рдкрд╛рдард╡реЗрд▓ рддреЗ database рд╕реНрддрд░ тАФ cache, replicas, proxy, partitions, queue тАФ рдЭреЗрд▓реВ рд╢рдХрддреЛ рдпрд╛рдЪреА рдЦрд╛рддреНрд░реА рдХрд░рд╛.

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

рдЬрддреНрд░рд╛ рдЖрддрд╛ рд╕реНрддрд░рд╛рдиреБрд╕рд╛рд░ рд╡рд╛рдвреВ рд╢рдХрддреЗ. рдкрдг рдореЛрдареНрдпрд╛ рдЬрддреНрд░реЗрдд рдЬрд╛рд╕реНрдд рднрд╛рдЧ рдЕрд╕рддрд╛рдд, рдЖрдгрд┐ рднрд╛рдЧ рдмрд┐рдШрдбрддрд╛рдд: рдПрдЦрд╛рджреА рд╣рд│реВ service, рдЕрдбрдХрд▓реЗрд▓рд╛ database, рд╣рд░рд╡рд▓реЗрд▓реА рдЗрдорд╛рд░рдд. рд╢реЗрд╡рдЯрдЪрд╛ рдзрдбрд╛: рдмрд┐рдШрд╛рдбрд╛рддреВрди рдЯрд┐рдХрдгреЗ тАФ timeouts, рд╕рднреНрдп retries, breaker switch рдЖрдгрд┐ рддреАрди рдЗрдорд╛рд░рддреА.

git checkout lesson-13-surviving-failure

ЁЯУм Lesson 12 тАФ Queues + the whole blueprint: take the burst at the gate

ЁЯУН You are here: Lesson 12 of 13 ┬╖ Previous: lesson-11-partitioning ┬╖ Next: lesson-13-surviving-failure


ЁЯУж What's in this branch

Lessons 01тАУ11, plus the last tool: a queue (Amazon SQS) between the request and the slow work, with workers that scale on the backlog тАФ and then the whole blueprint of a scaled web app. queues() in scale/demo.py sends a burst of certificate requests through a queue, and names what makes failed work safe to recover: at-least-once delivery, the visibility timeout, a DLQ and idempotent workers.

ЁЯзТ Explain like I'm 5

After the results, every parent wants a printed certificate. Printing takes time. If each parent waits at the counter while it prints, the queue never moves.

So Dipika uses tokens ЁЯОЯя╕П. A parent asks for a certificate, gets a token at once, and goes to enjoy the fair. The token goes into a box. In the back room, printers take tokens from the box, one by one, and print.

800 parents a second ask for three seconds. With 3 printers, the box fills up with 1,500 tokens and is still not empty ten seconds later. So Dipika adds a rule: "one printer for every 100 tokens in the box, up to 10". When the box fills, printers are added; when it empties, they go home. The box absorbs the rush: requests wait in it instead of crowding the counter.

If a printer jams in the middle of a certificate, the token goes back into the box after a while, and another printer tries again. Sometimes the same token is handed out twice, so each printer first checks "is this certificate already printed?". A token that fails again and again goes into a side box so someone can look at it. Tokens do not stay in the box forever: after some days, an old token is thrown away.

ЁЯЧ║я╕П Diagram

flowchart LR
    u["ЁЯСк parents"] --> api["ЁЯН│ API<br/>202 Accepted + token"]
    api -->|"SendMessage"| q["ЁЯУм SQS queue<br/>backlog"]
    q -->|"ReceiveMessage"| w["ЁЯЦия╕П workers<br/>1 per 100 waiting, max 10"]
    w -->|"DeleteMessage when done"| q
    q -.->|"failed too many times"| dlq["ЁЯзп dead-letter queue"]
    m["ЁЯУК ApproximateNumberOfMessagesVisible"] -.->|"scale workers"| w

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

тЭУ What

ЁЯдФ Why

Because some work is slower than a page should be, and bursts are shorter than the time it takes to add servers. A queue turns "2,400 requests in 3 seconds" into "2,400 jobs done over the next few seconds". The front stays fast, and the back runs at a pace it can survive. It is not magic: a message can still expire, or fail every time. With retries, a DLQ and idempotent workers, that failed work is kept and can be recovered safely, instead of being lost in a crash.

ЁЯФз How (in this repo)

drain(arrivals, per_worker, workers_for) in scale/sim.py plays one second at a time: it asks workers_for(backlog) how many workers to run тАФ decided from the backlog before this second's arrivals, as a real autoscaler sees a metric a little late тАФ adds the arrivals, and lets each worker finish per_worker jobs. Each row is (second, arrived, workers, done, backlog). queues() compares a fixed 3 workers with "one worker per 100 waiting, at least 1, at most 10".

ЁЯзк Try it

python3 scale/demo.py queues
python3 - <<'EOF'
import sys; sys.path.insert(0, "scale"); from sim import drain
arrivals = [50, 800, 800, 800, 200, 50, 50, 50, 50, 50]
for name, rule in (("fixed 3", lambda b: 3), ("fixed 8", lambda b: 8), ("1 per 100 waiting, max 10", lambda b: min(10, max(1, -(-b // 100)))),
                   ("1 per 100 waiting, max 20", lambda b: min(20, max(1, -(-b // 100)))), ("1 per 50 waiting, max 20", lambda b: min(20, max(1, -(-b // 50))))):
    rows = drain(arrivals, 100, rule)
    print(f"{name:<26} peak backlog {max(r[4] for r in rows):>5} ┬╖ most workers {max(r[2] for r in rows):>2} ┬╖ left after 10 s {rows[-1][4]:>4}")
for s, a, w, d, b in drain(arrivals, 100, lambda b: min(10, max(1, -(-b // 100)))):
    print(f"second {s}: arrived {a:>3}, workers {w:>2}, done {d:>4}, waiting {b:>4}")
EOF
python3 scale/test_scale.py

тЬЕ Verify тАФ what you should see

queues prints:

   fixed 3 workers                                      peak backlog  1500 ┬╖ left after 10 s  150
   scale on backlog (1 worker per 100 waiting, max 10)  peak backlog   800 ┬╖ left after 10 s    0
   the queue (SQS) absorbs the burst; workers process messages asynchronously, at their own pace
   SQS delivers at least once: a visibility timeout, retries, a DLQ and idempotent workers make failed work safe to recover

Your snippet prints:

fixed 3                    peak backlog  1500 ┬╖ most workers  3 ┬╖ left after 10 s  150
fixed 8                    peak backlog     0 ┬╖ most workers  8 ┬╖ left after 10 s    0
1 per 100 waiting, max 10  peak backlog   800 ┬╖ most workers  8 ┬╖ left after 10 s    0
1 per 100 waiting, max 20  peak backlog   800 ┬╖ most workers  8 ┬╖ left after 10 s    0
1 per 50 waiting, max 20   peak backlog   700 ┬╖ most workers 14 ┬╖ left after 10 s    0
second 0: arrived  50, workers  1, done   50, waiting    0
second 1: arrived 800, workers  1, done  100, waiting  700
second 2: arrived 800, workers  7, done  700, waiting  800
second 3: arrived 800, workers  8, done  800, waiting  800
second 4: arrived 200, workers  8, done  800, waiting  200
second 5: arrived  50, workers  2, done  200, waiting   50
second 6: arrived  50, workers  1, done  100, waiting    0
second 7: arrived  50, workers  1, done   50, waiting    0
second 8: arrived  50, workers  1, done   50, waiting    0
second 9: arrived  50, workers  1, done   50, waiting    0

The tests end with тЬЕ L12 a queue absorbs the burst and drains to zero and 13/13 passed.

ЁЯПБ What you just proved

In the model every request waited in the queue and was done later. Scaling on backlog cleared the burst by second 6 and went back to 1 worker; a fixed 3 was still 150 behind after 10 seconds. A fixed 8 never lets a backlog form тАФ but pays for 8 workers all day. Raising the maximum from 10 to 20 changed nothing, because the rule never asked for more than 8; a more eager rule (1 per 50) used up to 14 workers and kept the peak a little lower. And the first second of the burst is always slow: the autoscaler only sees the backlog after it forms.

тЪая╕П Common mistakes

ЁЯПн In production

On a real account тАФ a queue with a dead-letter queue (the queue keeps messages 4 days, the DLQ 14 days; after 5 receives a message moves to the DLQ):

aws sqs create-queue --queue-name certificates-dlq --attributes '{"MessageRetentionPeriod": "1209600"}'
aws sqs create-queue --queue-name certificates --attributes '{
  "VisibilityTimeout": "120",
  "MessageRetentionPeriod": "345600",
  "RedrivePolicy": "{\"deadLetterTargetArn\":\"arn:aws:sqs:ap-south-1:111122223333:certificates-dlq\",\"maxReceiveCount\":\"5\"}"
}'
aws sqs get-queue-attributes --attribute-names ApproximateNumberOfMessagesVisible \
    --queue-url https://sqs.ap-south-1.amazonaws.com/111122223333/certificates

KEDA on Kubernetes тАФ about one worker pod per 100 waiting messages, from 1 to 10:

apiVersion: keda.sh/v1alpha1
kind: ScaledObject
metadata: { name: certificate-workers }
spec:
  scaleTargetRef: { name: certificate-worker }
  minReplicaCount: 1
  maxReplicaCount: 10
  triggers:
    - type: aws-sqs-queue
      authenticationRef: { name: keda-aws }
      metadata:
        queueURL: https://sqs.ap-south-1.amazonaws.com/111122223333/certificates
        queueLength: "100"
        awsRegion: ap-south-1

ЁЯПЧя╕П The whole blueprint тАФ the scaled fair

flowchart LR
    u["ЁЯСк parents"] --> cf["ЁЯМН CloudFront<br/>UI from S3 ┬╖ micro-cache ┬╖ shield"]
    cf -->|"/api"| lb["тЪЦя╕П ALB or API Gateway"]
    lb --> api["ЁЯН│ stateless API<br/>ASG ┬╖ pods + HPA/Karpenter ┬╖ Lambda"]
    api --> rc["ЁЯУМ Redis / Valkey"]
    api --> px["ЁЯЫОя╕П RDS Proxy"]
    px --> pri["ЁЯУТ primary"]
    px --> rep["ЁЯУЪ replicas"]
    api --> ddb["ЁЯЧВя╕П DynamoDB<br/>keys that spread"]
    api --> q["ЁЯУм SQS"] --> wk["ЁЯЦия╕П workers<br/>scale on backlog"]
Tier Scales by Watch
UI CloudFront edges, versioned files, s-maxage on busy pages cache hit rate, origin requests
API stateless copies behind a load balancer: ASG, HPA + Karpenter, or Lambda p99 latency, 5xx, concurrency, Pending pods
Reads Redis cache-aside, read replicas hit rate, ReplicaLag
Connections RDS Proxy / a pool DatabaseConnections vs max_connections
Writes partitions with a key that spreads throttled requests, hot keys
Slow work SQS + workers on backlog ApproximateNumberOfMessagesVisible, DLQ depth

ЁЯПн Why this matters in production: every tier scales on its own signal. Before a big day, load-test the whole path (lesson 02), schedule the known spike (lessons 06, 08), and make sure the database tier тАФ cache, replicas, proxy, partitions, queue тАФ can take what the API tier will now send it.

тПня╕П Next

The fair can now grow tier by tier. But a bigger fair has more parts, and parts break: a slow service, a stuck database, a lost building. Last lesson: surviving failure тАФ timeouts, polite retries, a breaker switch and three buildings.

git checkout lesson-13-surviving-failure
тЖР PreviouspartitioningNext тЖТsurviving failure

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