ЁЯПл The SchoolтА║ЁЯУР System DesignтА║ЁЯУг рдзрдбрд╛ 08 тАФ Events, pub/sub рдЖрдгрд┐ outbox: рдПрдХ рддрдереНрдп, рдЕрдиреЗрдХ рдРрдХрдгрд╛рд░реЗ
ЁЯЦ╝я╕П See the drawing + lab ЁЯПа Course home ЁЯМ┐ Branch on GitHub тЬПя╕П View source
ЁЯЦ╝я╕П рдЖрдХреГрддреА рдЖрдгрд┐ labThe drawing + lab рдкреВрд░реНрдг рдкрд╛рдирд╛рд╡рд░ рдЙрдШрдбрд╛ тЖЧOpen full page тЖЧ

ЁЯУг рдзрдбрд╛ 08 тАФ Events, pub/sub рдЖрдгрд┐ outbox: рдПрдХ рддрдереНрдп, рдЕрдиреЗрдХ рдРрдХрдгрд╛рд░реЗ

ЁЯУН рддреБрдореНрд╣реА рдЗрдереЗ рдЖрд╣рд╛рдд: 18 рдкреИрдХреА рдзрдбрд╛ 08 ┬╖ рдорд╛рдЧреАрд▓: lesson-07-queues ┬╖ рдкреБрдвреАрд▓: lesson-09-services


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

рдзрдбреЗ 01тАУ07, рдЖрдгрд┐ events. рд╕реВрдЪрдирд╛ post рдЭрд╛рд▓реА рдХреА API рдПрдХ рддрдереНрдп рдЬрд╛рд╣реАрд░ рдХрд░рддреЛ тАФ NoticePosted тАФ рдЖрдгрд┐ SMS, email, search index рдЖрдгрд┐ analytics рдкреНрд░рддреНрдпреЗрдХ рд╕реНрд╡рддрдВрддреНрд░рдкрдгреЗ рдкреНрд░рддрд┐рд╕рд╛рдж рджреЗрддрд╛рдд (pub/sub). рдордЧ рдХрдареАрдг рднрд╛рдЧ: dual-write problem (рд╕реВрдЪрдирд╛ save рд╣реЛрддреЗ рдкрдг event рд╣рд░рд╡рддреЛ) рдЖрдгрд┐ рддреНрдпрд╛рд╡рд░рдЪрд╛ рдЙрдкрд╛рдп, relay рдХрд┐рдВрд╡рд╛ CDC рд╕рд╣ transactional outbox. рд╢реЗрд╡рдЯреА feed рдЪрд╛ рдкреНрд░рд╢реНрди: fan-out on write vs on read, рдЖрдгрд┐ celebrity hybrid. Outbox, dual_write() рдЖрдгрд┐ fan_out() design/blocks.py рдордзреНрдпреЗ, notify() design/designs.py рдордзреНрдпреЗ, рдЖрдгрд┐ events() design/demo.py рдордзреНрдпреЗ.

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

рдирд╡реА рд╕реВрдЪрдирд╛ рдЖрд▓реА рдХреА рд╢рд╛рд│реЗрддрд▓реНрдпрд╛ рдЕрдиреЗрдХрд╛рдВрдирд╛ рддрд┐рдЪреНрдпрд╛рд╢реА рдХрд╛рдо рдЕрд╕рддреЗ: SMS рдорджрддрдиреАрд╕, email рдорджрддрдиреАрд╕, рд╢рдмреНрджрд╕реВрдЪреА (word index) рдареЗрд╡рдгрд╛рд░реА рдЧреНрд░рдВрдердкрд╛рд▓, рдЖрдгрд┐ рд╡рд╛рд░реНрд╖рд┐рдХ рдЕрд╣рд╡рд╛рд▓рд╛рд╕рд╛рдареА рдореЛрдЬрдгреА рдХрд░рдгрд╛рд░реА рд╡реНрдпрдХреНрддреА. рдХрд╛рд░рдХреБрдирд╛рд▓рд╛ рдкреНрд░рддреНрдпреЗрдХрд╛рдХрдбреЗ рдЪрд╛рд▓рдд рдЬрд╛рд╡реЗ рд▓рд╛рдЧрд▓реЗ рдЕрд╕рддреЗ, рддрд░ рддреА рдПрдЦрд╛рджреНрдпрд╛рд▓рд╛ рд╡рд┐рд╕рд░рд▓реА рдЕрд╕рддреА, рдЖрдгрд┐ рдкреНрд░рддреНрдпреЗрдХ рдирд╡реНрдпрд╛ рдорджрддрдиреАрд╕рд╛рд╕рд╛рдареА рдПрдХ рдирд╡реА рдлреЗрд░реА рдЭрд╛рд▓реА рдЕрд╕рддреА.

рдореНрд╣рдгреВрди рджреАрдкрд┐рдХрд╛ рдПрдХ рдШрдВрдЯрд╛ ЁЯФФ рд╡рд╛рдЬрд╡рддреЗ рдЖрдгрд┐ рднреВрддрдХрд╛рд│рд╛рдд рдПрдХ рд╡рд╛рдХреНрдп рдореНрд╣рдгрддреЗ: "рд╕реВрдЪрдирд╛ 42 post рдЭрд╛рд▓реА." рдЬреНрдпрд╛рд▓рд╛ рдХрд╛рдо рдЖрд╣реЗ рддреЛ рдРрдХрдд рдЕрд╕рддреЛ рдЖрдгрд┐ рдЖрдкрд▓реЗ рдХрд╛рдо рдХрд░рддреЛ. рдирд╡рд╛ рдорджрддрдиреАрд╕ рдлрдХреНрдд рдРрдХрд╛рдпрд▓рд╛ рд╕реБрд░реБрд╡рд╛рдд рдХрд░рддреЛ. рдпрд╛рд▓рд╛рдЪ publish / subscribe рдореНрд╣рдгрддрд╛рдд.

рдкрдг рдЗрдереЗ рдПрдХ рд╕рд╛рдкрд│рд╛ рдЖрд╣реЗ. рдХрд╛рд░рдХреВрди рд╕реВрдЪрдирд╛ рдиреЛрдВрджрд╡рд╣реАрдд рд▓рд┐рд╣рд┐рддреЗ, рдЖрдгрд┐ рдордЧ рдШрдВрдЯреЗрдХрдбреЗ рдЪрд╛рд▓рдд рдЬрд╛рддреЗ. рд╡рд╛рдЯреЗрдд рддреА рдЕрдбрдЦрд│рд▓реА, рддрд░ рд╕реВрдЪрдирд╛ рдиреЛрдВрджрд╡рд╣реАрдд рдЕрд╕рддреЗ рдкрдг рдШрдВрдЯрд╛ рдХрдзреАрдЪ рд╡рд╛рдЬрдд рдирд╛рд╣реА. рдХреЛрдгреАрдЪ SMS рдкрд╛рдард╡рдд рдирд╛рд╣реА. рдХреЛрдгрд╛рд▓рд╛рдЪ рдХрд│рдд рдирд╛рд╣реА.

рдореНрд╣рдгреВрди рджреАрдкрд┐рдХрд╛ рдиреЛрдВрджрд╡рд╣реАрдЪреНрдпрд╛ рдЖрддрдЪ рдПрдХ "рдХрд░рд╛рдпрдЪреНрдпрд╛ рдШреЛрд╖рдгрд╛" рдкрд╛рди ЁЯУЭ рдЬреЛрдбрддреЗ. рдХрд╛рд░рдХреВрди рд╕реВрдЪрдирд╛ рдЖрдгрд┐ рдШреЛрд╖рдгрд╛ рдПрдХрд╛рдЪ рдкрд╛рдирд╛рд╡рд░, рдПрдХрд╛рдЪ рджрдорд╛рдд рд▓рд┐рд╣рд┐рддреЗ тАФ рджреЛрдиреНрд╣реА рдХрд┐рдВрд╡рд╛ рдПрдХрд╣реА рдирд╛рд╣реА. рдПрдХ рдзрд╛рд╡рдкрдЯреВ рджрд░ рдХрд╛рд╣реА рд╕реЗрдХрдВрджрд╛рдВрдиреА рддреЗ рдкрд╛рди рд╡рд╛рдЪрддреЛ рдЖрдгрд┐ рдкреНрд░рддреНрдпреЗрдХ рдУрд│реАрд╕рд╛рдареА рдШрдВрдЯрд╛ рд╡рд╛рдЬрд╡рддреЛ, рдордЧ рддреНрдпрд╛рд╡рд░ рдЦреВрдг рдХрд░рддреЛ. рдзрд╛рд╡рдкрдЯреВ рдЕрдбрдЦрд│рд▓рд╛ рддрд░реА рдУрд│ рддрд┐рдереЗрдЪ рдЕрд╕рддреЗ; рдкреБрдврдЪрд╛ рдзрд╛рд╡рдкрдЯреВ рддреА рд╡рд╛рдЬрд╡рддреЛ. рдХрд╛рд╣реАрдЪ рд╣рд░рд╡рдд рдирд╛рд╣реА. рдПрдХрд╛ рдУрд│реАрд╕рд╛рдареА рдШрдВрдЯрд╛ рджреЛрдирджрд╛ рд╡рд╛рдЬреВ рд╢рдХрддреЗ тАФ рдореНрд╣рдгреВрди рдкреНрд░рддреНрдпреЗрдХ рдРрдХрдгрд╛рд░рд╛ рддрдкрд╛рд╕рддреЛ "рд╕реВрдЪрдирд╛ 42 рдЪреЗ рдХрд╛рдо рдореА рдЖрдзреАрдЪ рдХреЗрд▓реЗ рдХрд╛?"

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

flowchart LR
    api["ЁЯФМ API: post notice"] -->|"ONE transaction"| db[("ЁЯРШ database<br/>notices + outbox rows")]
    db --> relay["ЁЯПГ relay / CDC<br/>reads outbox, publishes, marks sent"]
    relay --> bus{{"ЁЯУг topic: NoticePosted"}}
    bus --> sms["ЁЯУ▒ SMS queue"]
    bus --> mail["тЬЙя╕П email"]
    bus --> idx["ЁЯФО search index"]
    bus --> an["ЁЯУК analytics"]
    bad["тЭМ dual write:<br/>save тЖТ crash тЖТ publish<br/>= event lost"]

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

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

ЁЯдФ рдХрд╛

рдХрд╛рд░рдг рдХрд╛рдЧрдж (рдзрдбрд╛ 01) рд╕рд╛рдВрдЧрддреЛ "'posted' рджрд╛рдЦрд╡рд▓реНрдпрд╛рдирдВрддрд░ рд╕реВрдЪрдирд╛ рдХрдзреАрд╣реА рд╣рд░рд╡рддрд╛ рдХрд╛рдорд╛ рдирдпреЗ" тАФ рдЖрдгрд┐ рддреНрдпрд╛рдд рддрд┐рдЪрд╛ SMS рдкрдг рдпреЗрддреЛ. Dual write рдордзреНрдпреЗ, рдПрдХрд╛ рдЫреЛрдЯреНрдпрд╛ рдЦрд┐рдбрдХреАрддрд▓рд╛ рдХреЛрдгрддрд╛рд╣реА crash рд╣реЗ рд╡рдЪрди рдореЛрдбрддреЛ, рдЖрдгрд┐ рдХреБрдареЗрд╣реА error рджрд┐рд╕рдд рдирд╛рд╣реА. Outbox рдореБрд│реЗ event рд╕реВрдЪрдиреЗрдЗрддрдХрд╛рдЪ рдЯрд┐рдХрд╛рдК рд╣реЛрддреЛ. Pub/sub рдореБрд│реЗ notice API рд▓рд╣рд╛рди рд░рд╛рд╣рддреЛ: SMS, email, search рдЖрдгрд┐ analytics рдЕрд╕реНрддрд┐рддреНрд╡рд╛рдд рдЖрд╣реЗрдд рд╣реЗ рддреНрдпрд╛рд▓рд╛ рдорд╛рд╣реАрдд рдирд╕рддреЗ, рдореНрд╣рдгреВрди рдкреНрд░рддреНрдпреЗрдХ рдЬрдг рд╕реНрд╡рддрдВрддреНрд░рдкрдгреЗ рдмрджрд▓реВ, рдЕрдкрдпрд╢реА рд╣реЛрдК рдЖрдгрд┐ scale рд╣реЛрдК рд╢рдХрддреЛ.

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

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

python3 design/demo.py events
python3 - <<'EOF'
import sys; sys.path.insert(0, "design"); from blocks import Outbox, dual_write, fan_out
from designs import notify
print("dual write, crash тЖТ", dual_write(True))
o = Outbox()
for n in ("N-1", "N-2", "N-3"):
    o.place(n, crash_before_publish=True)
print("after 3 posts and a dead broker: waiting", o.outbox, "┬╖ published", o.published)
o.relay(); print("relay runs тЖТ published", [e[1] for e in o.published], "┬╖ waiting", o.outbox)
for followers in (40, 5_000, 2_000_000):
    print(f"{followers:>9,} followers, 3 posts/day: push {fan_out(followers, 3, 3, 400, 'push')} ┬╖ pull {fan_out(followers, 3, 3, 400, 'pull')}")
for followers in (40, 99_999, 100_000):
    print(f"{followers:>7,} followers тЖТ {notify(followers, 3)}")
EOF

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

events рд╣реЗ рдЫрд╛рдкрддреЛ:

тФАтФА 'NoticePosted' is one event; many services react (email, SMS, search index, analytics) тАФ pub/sub
   dual write, crash between DB and broker = False тЖТ ('notice saved', 'NoticePosted')
   dual write, crash between DB and broker = True  тЖТ ('notice saved', None)
   outbox, crash before publishing тЖТ saved (event waits in the outbox) ┬╖ relay runs later тЖТ saved + published [('NoticePosted', 'N-1')]
   fan-out push: a class with 40 parents, parents follow 3 classes, 5 posts/day, 400 reads/day тЖТ {'writes': 200, 'read_lookups': 400}
   fan-out pull: a class with 40 parents, parents follow 3 classes, 5 posts/day, 400 reads/day тЖТ {'writes': 5, 'read_lookups': 1200}
   events are facts in the past tense; consumers must be idempotent (the same event may arrive twice)

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

dual write, crash тЖТ ('notice saved', None)
after 3 posts and a dead broker: waiting [('NoticePosted', 'N-1'), ('NoticePosted', 'N-2'), ('NoticePosted', 'N-3')] ┬╖ published []
relay runs тЖТ published ['N-1', 'N-2', 'N-3'] ┬╖ waiting []
       40 followers, 3 posts/day: push {'writes': 120, 'read_lookups': 400} ┬╖ pull {'writes': 3, 'read_lookups': 1200}
    5,000 followers, 3 posts/day: push {'writes': 15000, 'read_lookups': 400} ┬╖ pull {'writes': 3, 'read_lookups': 1200}
2,000,000 followers, 3 posts/day: push {'writes': 6000000, 'read_lookups': 400} ┬╖ pull {'writes': 3, 'read_lookups': 1200}
     40 followers тЖТ {'mode': 'push on write', 'inbox_writes_per_day': 120}
 99,999 followers тЖТ {'mode': 'push on write', 'inbox_writes_per_day': 299997}
100,000 followers тЖТ {'mode': 'pull on read (celebrity)', 'inbox_writes_per_day': 0}

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

Dual write event рдЧреБрдкрдЪреВрдк рд╣рд░рд╡рддреЛ тАФ рд╕реВрдЪрдирд╛ save рд╣реЛрддреЗ рдЖрдгрд┐ SMS рдХрдзреАрдЪ рдЬрд╛рдгрд╛рд░ рдирд╛рд╣реА рдЕрд╕реЗ рдХрд╛рд╣реАрд╣реА рд╕рд╛рдВрдЧрдд рдирд╛рд╣реА. Outbox рдЕрд╕реЗрд▓ рддрд░ broker outage рдЪреНрдпрд╛ рд╡реЗрд│рдЪреНрдпрд╛ рддреАрди posts рддреАрди events рдерд╛рдВрдмрд╡реВрди рдареЗрд╡рддрд╛рдд, рдЖрдгрд┐ relay рдЪреА рдПрдХ рдлреЗрд░реА рддрд┐рдиреНрд╣реА рдХреНрд░рдорд╛рдиреЗ publish рдХрд░рддреЗ. Fan-out рдордзреНрдпреЗ, push рдЪрд╛ рдЦрд░реНрдЪ рд╢реНрд░реЛрддреГрд╡рд░реНрдЧрд╛рд╕реЛрдмрдд рд╡рд╛рдврддреЛ (2-million-follower рдЕрд╕рд▓реЗрд▓реНрдпрд╛ рдПрдХрд╛ рд▓реЗрдЦрдХрд╛рд╕рд╛рдареА рджрд┐рд╡рд╕рд╛рд▓рд╛ 6 million inbox writes) рддрд░ pull рдЪрд╛ рдЦрд░реНрдЪ рд╕рдкрд╛рдЯ рд░рд╛рд╣рддреЛ тАФ рдореНрд╣рдгреВрдирдЪ hybrid 40 рдЬрдгрд╛рдВрдЪреНрдпрд╛ рд╡рд░реНрдЧрд╛рд╕рд╛рдареА push рдХрд░рддреЛ рдЖрдгрд┐ celebrity рд╕рд╛рдареА pull рдХрд░рддреЛ.

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

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

PostgreSQL рдордзреАрд▓ outbox тАФ рд╕реВрдЪрдирд╛ рдЖрдгрд┐ рддрд┐рдЪрд╛ event рдПрдХрд╛рдЪ transaction рдордзреНрдпреЗ:

CREATE TABLE outbox (
  id           bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
  event_type   text        NOT NULL,           -- 'NoticePosted'
  aggregate_id text        NOT NULL,           -- the class id: keeps order per class
  payload      jsonb       NOT NULL,
  created_at   timestamptz NOT NULL DEFAULT now(),
  sent_at      timestamptz
);

BEGIN;
INSERT INTO notices (class_id, author_id, title, body, urgent) VALUES ('3A', 7, 'School closed', 'тАж', true) RETURNING id;
INSERT INTO outbox (event_type, aggregate_id, payload)
     VALUES ('NoticePosted', '3A', '{"notice_id": 42, "class_id": "3A", "urgent": true}');
COMMIT;

-- the relay, every second (SKIP LOCKED lets several relays share the work):
SELECT id, event_type, payload FROM outbox WHERE sent_at IS NULL ORDER BY id LIMIT 100 FOR UPDATE SKIP LOCKED;
-- publish each one, then:
UPDATE outbox SET sent_at = now() WHERE id = ANY($1);

рдЦрд▒реНрдпрд╛ account рд╡рд░ тАФ relay рдПрдХрд╛ SNS topic рд╡рд░ publish рдХрд░рддреЛ, рдЖрдгрд┐ рдкреНрд░рддреНрдпреЗрдХ consumer рдЪреА рд╕реНрд╡рддрдГрдЪреА SQS queue рддреНрдпрд╛рд▓рд╛ subscribe рдХреЗрд▓реЗрд▓реА рдЕрд╕рддреЗ (AWS рдЪреЗ рдкрд╛рд░рдВрдкрд░рд┐рдХ fan-out):

aws sns create-topic --name notice-posted
aws sns subscribe --topic-arn arn:aws:sns:ap-south-1:123456789012:notice-posted \
    --protocol sqs --notification-endpoint arn:aws:sqs:ap-south-1:123456789012:urgent-sms
aws sns publish --topic-arn arn:aws:sns:ap-south-1:123456789012:notice-posted \
    --message '{"type":"NoticePosted","notice_id":42,"class_id":"3A","urgent":true}'

DynamoDB рдордзреНрдпреЗ DynamoDB Streams рд╣реЗ рдЕрдВрдЧрднреВрдд CDC рдЖрд╣реЗ; PostgreSQL рдХрд┐рдВрд╡рд╛ MySQL рдордзреНрдпреЗ, Debezium log рд╡рд╛рдЪрддреЛ рдЖрдгрд┐ Kafka рдордзреНрдпреЗ рд▓рд┐рд╣рд┐рддреЛ. EventBridge рд╣реА рдЖрдгрдЦреА рдПрдХ managed bus рдЖрд╣реЗ, рдЬрд┐рдЪреЗ rules events рддреНрдпрд╛рдВрдЪреНрдпрд╛ content рдиреБрд╕рд╛рд░ рдкрд╛рдард╡рддрд╛рдд.

ЁЯПн рдкреНрд░рддреНрдпрдХреНрд╖ рд╡рд╛рдкрд░рд╛рдд рд╣реЗ рдХрд╛ рдорд╣рддреНрддреНрд╡рд╛рдЪреЗ: database рдХрдбреВрди broker рдХрдбреЗ рдЬрд╛рдгрд╛рд░рд╛ рдкреНрд░рддреНрдпреЗрдХ рдмрд╛рдг рдХрд╛рдврд╛ рдЖрдгрд┐ рд╡рд┐рдЪрд╛рд░рд╛ "рдиреЗрдордХреЗ рдЗрдереЗ crash рдЭрд╛рд▓реЗ рддрд░?" рдЙрддреНрддрд░ "рддреЗ рд╣рд░рд╡рддреЗ" рдЕрд╕реЗрд▓, рддрд░ launch рдЖрдзреА outbox (рдХрд┐рдВрд╡рд╛ CDC) рдЬреЛрдбрд╛ тАФ рд╣рд╛ bug tests рдордзреНрдпреЗ рджрд┐рд╕рдд рдирд╛рд╣реА.

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

SMS, email, search, analytics тАФ рдкреНрд░рддреНрдпреЗрдХрд╛рдЪреА рд╕реНрд╡рддрдГрдЪреА service рдЕрд╕рд╛рд╡реА рдХрд╛? рдХреА рдПрдХрдЪ program? рдирд┐рдпреЛрдЬрди рдХрд╛рд░реНрдпрд╛рд▓рдп рдкреНрд░рддреНрдпреЗрдХ рдЬрд╛рд╕реНрддреАрдЪреНрдпрд╛ hop рдЪреА рдХрд┐рдВрдордд рдореЛрдЬрддреЗ.

git checkout lesson-09-services

ЁЯУг Lesson 08 тАФ Events, pub/sub & the outbox: one fact, many listeners

ЁЯУН You are here: Lesson 08 of 18 ┬╖ Previous: lesson-07-queues ┬╖ Next: lesson-09-services


ЁЯУж What's in this branch

Lessons 01тАУ07, plus events. When a notice is posted, the API announces one fact тАФ NoticePosted тАФ and SMS, email, the search index and analytics each react on their own (pub/sub). Then the hard part: the dual-write problem (the notice is saved but the event is lost) and its fix, the transactional outbox with a relay or CDC. Last, the feed question: fan-out on write vs on read, and the celebrity hybrid. Outbox, dual_write() and fan_out() in design/blocks.py, notify() in design/designs.py, and events() in design/demo.py.

ЁЯзТ Explain like I'm 5

When a new notice arrives, many people in the school care: the SMS helpers, the email helper, the librarian who keeps the word index, and the person who counts things for the yearly report. If the clerk had to walk to each of them, she would forget one, and every new helper would mean a new walk.

So Dipika rings a bell ЁЯФФ and says one sentence, in the past tense: "Notice 42 was posted." Anyone who cares is listening and does their own job. A new helper just starts listening. That is publish / subscribe.

But there is a trap. The clerk writes the notice in the register, then walks to the bell. If she trips on the way, the notice is in the register but the bell never rings. Nobody sends the SMS. Nobody knows.

So Dipika adds an "announcements to make" page ЁЯУЭ inside the register. The clerk writes the notice and the announcement on the same page, in one go тАФ both or neither. A runner reads that page every few seconds and rings the bell for each line, then ticks it. If the runner trips, the line is still there; the next runner rings it. Nothing is lost. The bell may ring twice for one line тАФ so every listener checks "did I already do notice 42?"

ЁЯЧ║я╕П Diagram

flowchart LR
    api["ЁЯФМ API: post notice"] -->|"ONE transaction"| db[("ЁЯРШ database<br/>notices + outbox rows")]
    db --> relay["ЁЯПГ relay / CDC<br/>reads outbox, publishes, marks sent"]
    relay --> bus{{"ЁЯУг topic: NoticePosted"}}
    bus --> sms["ЁЯУ▒ SMS queue"]
    bus --> mail["тЬЙя╕П email"]
    bus --> idx["ЁЯФО search index"]
    bus --> an["ЁЯУК analytics"]
    bad["тЭМ dual write:<br/>save тЖТ crash тЖТ publish<br/>= event lost"]

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

тЭУ What

ЁЯдФ Why

Because the sheet (lesson 01) says "a notice must never be lost once 'posted' is shown" тАФ and that includes its SMS. With a dual write, any crash in a tiny window breaks that promise, and nothing shows an error. The outbox makes the event as durable as the notice itself. Pub/sub keeps the notice API small: it does not know that SMS, email, search and analytics exist, so each can change, fail and scale on its own.

ЁЯФз How (in this repo)

ЁЯзк Try it

python3 design/demo.py events
python3 - <<'EOF'
import sys; sys.path.insert(0, "design"); from blocks import Outbox, dual_write, fan_out
from designs import notify
print("dual write, crash тЖТ", dual_write(True))
o = Outbox()
for n in ("N-1", "N-2", "N-3"):
    o.place(n, crash_before_publish=True)
print("after 3 posts and a dead broker: waiting", o.outbox, "┬╖ published", o.published)
o.relay(); print("relay runs тЖТ published", [e[1] for e in o.published], "┬╖ waiting", o.outbox)
for followers in (40, 5_000, 2_000_000):
    print(f"{followers:>9,} followers, 3 posts/day: push {fan_out(followers, 3, 3, 400, 'push')} ┬╖ pull {fan_out(followers, 3, 3, 400, 'pull')}")
for followers in (40, 99_999, 100_000):
    print(f"{followers:>7,} followers тЖТ {notify(followers, 3)}")
EOF

тЬЕ Verify тАФ what you should see

events prints:

тФАтФА 'NoticePosted' is one event; many services react (email, SMS, search index, analytics) тАФ pub/sub
   dual write, crash between DB and broker = False тЖТ ('notice saved', 'NoticePosted')
   dual write, crash between DB and broker = True  тЖТ ('notice saved', None)
   outbox, crash before publishing тЖТ saved (event waits in the outbox) ┬╖ relay runs later тЖТ saved + published [('NoticePosted', 'N-1')]
   fan-out push: a class with 40 parents, parents follow 3 classes, 5 posts/day, 400 reads/day тЖТ {'writes': 200, 'read_lookups': 400}
   fan-out pull: a class with 40 parents, parents follow 3 classes, 5 posts/day, 400 reads/day тЖТ {'writes': 5, 'read_lookups': 1200}
   events are facts in the past tense; consumers must be idempotent (the same event may arrive twice)

Your snippet prints:

dual write, crash тЖТ ('notice saved', None)
after 3 posts and a dead broker: waiting [('NoticePosted', 'N-1'), ('NoticePosted', 'N-2'), ('NoticePosted', 'N-3')] ┬╖ published []
relay runs тЖТ published ['N-1', 'N-2', 'N-3'] ┬╖ waiting []
       40 followers, 3 posts/day: push {'writes': 120, 'read_lookups': 400} ┬╖ pull {'writes': 3, 'read_lookups': 1200}
    5,000 followers, 3 posts/day: push {'writes': 15000, 'read_lookups': 400} ┬╖ pull {'writes': 3, 'read_lookups': 1200}
2,000,000 followers, 3 posts/day: push {'writes': 6000000, 'read_lookups': 400} ┬╖ pull {'writes': 3, 'read_lookups': 1200}
     40 followers тЖТ {'mode': 'push on write', 'inbox_writes_per_day': 120}
 99,999 followers тЖТ {'mode': 'push on write', 'inbox_writes_per_day': 299997}
100,000 followers тЖТ {'mode': 'pull on read (celebrity)', 'inbox_writes_per_day': 0}

ЁЯПБ What you just proved

The dual write loses the event silently тАФ the notice is saved and nothing says the SMS will never go. With the outbox, three posts during a broker outage leave three events waiting, and one relay run publishes all three, in order. For fan-out, push cost grows with the audience (6 million inbox writes a day for one 2-million-follower author) while pull cost stays flat тАФ which is why the hybrid pushes for a class of 40 and pulls for a celebrity.

тЪая╕П Common mistakes

ЁЯПн In production

The outbox in PostgreSQL тАФ the notice and its event in one transaction:

CREATE TABLE outbox (
  id           bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
  event_type   text        NOT NULL,           -- 'NoticePosted'
  aggregate_id text        NOT NULL,           -- the class id: keeps order per class
  payload      jsonb       NOT NULL,
  created_at   timestamptz NOT NULL DEFAULT now(),
  sent_at      timestamptz
);

BEGIN;
INSERT INTO notices (class_id, author_id, title, body, urgent) VALUES ('3A', 7, 'School closed', 'тАж', true) RETURNING id;
INSERT INTO outbox (event_type, aggregate_id, payload)
     VALUES ('NoticePosted', '3A', '{"notice_id": 42, "class_id": "3A", "urgent": true}');
COMMIT;

-- the relay, every second (SKIP LOCKED lets several relays share the work):
SELECT id, event_type, payload FROM outbox WHERE sent_at IS NULL ORDER BY id LIMIT 100 FOR UPDATE SKIP LOCKED;
-- publish each one, then:
UPDATE outbox SET sent_at = now() WHERE id = ANY($1);

On a real account тАФ the relay publishes to one SNS topic, and each consumer has its own SQS queue subscribed to it (the classic AWS fan-out):

aws sns create-topic --name notice-posted
aws sns subscribe --topic-arn arn:aws:sns:ap-south-1:123456789012:notice-posted \
    --protocol sqs --notification-endpoint arn:aws:sqs:ap-south-1:123456789012:urgent-sms
aws sns publish --topic-arn arn:aws:sns:ap-south-1:123456789012:notice-posted \
    --message '{"type":"NoticePosted","notice_id":42,"class_id":"3A","urgent":true}'

With DynamoDB, DynamoDB Streams is built-in CDC; with PostgreSQL or MySQL, Debezium reads the log and writes to Kafka. EventBridge is another managed bus with rules that route events by their content.

ЁЯПн Why this matters in production: draw every arrow that crosses from the database to a broker and ask "what if we crash right here?" If the answer is "we lose it", add an outbox (or CDC) before launch тАФ this bug does not show up in tests.

тПня╕П Next

SMS, email, search, analytics тАФ should each one be its own service? Or one program? The planning office counts what every extra hop costs.

git checkout lesson-09-services
тЖР PreviousqueuesNext тЖТservices

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