ЁЯУг рдзрдбрд╛ 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
тЭУ рдХрд╛рдп
- Event тАФ рдХрд╛рд╣реАрддрд░реА рдШрдбрд▓реЗ рдпрд╛рдЪреА рдиреЛрдВрдж, рднреВрддрдХрд╛рд│рд╛рдд рдирд╛рд╡ рджрд┐рд▓реЗрд▓реА:
NoticePosted,ParentJoinedClass. рддреЗ рдПрдХ рддрдереНрдп рдЖрд╣реЗ; рддреНрдпрд╛рд▓рд╛ рдХреЛрдгреА рдирд╛рд╣реА рдореНрд╣рдгреВ рд╢рдХрдд рдирд╛рд╣реА. (рдПрдХ command тАФSendSmsтАФ рдПрдХрд╛ receiver рд▓рд╛ рдХрд╛рд╣реАрддрд░реА рдХрд░рд╛рдпрд▓рд╛ рд╕рд╛рдВрдЧрддреЗ, рдЖрдгрд┐ рддреА рдирд╛рдХрд╛рд░рд▓реА рдЬрд╛рдК рд╢рдХрддреЗ.) - Pub/sub тАФ producer рдПрдХрд╛ topic рд╡рд░ publish рдХрд░рддреЛ; рдкреНрд░рддреНрдпреЗрдХ subscriber рд▓рд╛ рд╕реНрд╡рддрдГрдЪреА рдкреНрд░рдд рдорд┐рд│рддреЗ. рдХреЛрдг рдРрдХрдд рдЖрд╣реЗ рд╣реЗ producer рд▓рд╛ рдорд╛рд╣реАрдд рдирд╕рддреЗ. рдирд╡рд╛ рдРрдХрдгрд╛рд░рд╛ рдЬреЛрдбрд▓реНрдпрд╛рдиреЗ producer рдЪрд╛ code рдмрджрд▓рдд рдирд╛рд╣реА.
- Queue vs topic тАФ queue (рдзрдбрд╛ 07) рдкреНрд░рддреНрдпреЗрдХ job рдПрдХрд╛рдЪ worker рд▓рд╛ рджреЗрддреЗ; topic рдкреНрд░рддреНрдпреЗрдХ event рдкреНрд░рддреНрдпреЗрдХ subscriber рд▓рд╛ рджреЗрддреЛ (рдмрд╣реБрддреЗрдХ рд╡реЗрд│рд╛ рдкреНрд░рддреНрдпреЗрдХ subscriber рдорд╛рдЧреЗ рддреНрдпрд╛рдЪреА рд╕реНрд╡рддрдГрдЪреА queue рдЕрд╕рддреЗ).
- Dual-write problem тАФ app рджреЛрди systems рдордзреНрдпреЗ рд▓рд┐рд╣рд┐рддреЛ (рдЖрдзреА database, рдордЧ broker), рддреНрдпрд╛ рджреЛрдШрд╛рдВрд╡рд░ рдорд┐рд│реВрди рдХреЛрдгрддрд╛рд╣реА transaction рдирд╕рддреЛ. рджреЛрдШрд╛рдВрдЪреНрдпрд╛ рдордзреНрдпреЗ crash, timeout рдХрд┐рдВрд╡рд╛ broker outage рдЭрд╛рд▓реЗ рддрд░ рддреЗ рдПрдХрдореЗрдХрд╛рдВрд╢реА рдЬреБрд│рдд рдирд╛рд╣реАрдд: рд╕реВрдЪрдирд╛ save рдЭрд╛рд▓реА, event рд╣рд░рд╡рд▓рд╛ тАФ рдХрд┐рдВрд╡рд╛, рдЖрдзреА publish рдХреЗрд▓реЗ рддрд░, рдХрдзреАрдЪ save рди рдЭрд╛рд▓реЗрд▓реНрдпрд╛ рд╕реВрдЪрдиреЗрдЪрд╛ event.
- Transactional outbox тАФ business row рдЖрдгрд┐ рдПрдХ
outboxrow рдПрдХрд╛рдЪ local database transaction рдордзреНрдпреЗ рд▓рд┐рд╣рд╛. рдПрдХрддрд░ рджреЛрдиреНрд╣реА рдЕрд╕рддрд╛рдд рдХрд┐рдВрд╡рд╛ рдПрдХрд╣реА рдирд╕рддреЛ. рдордЧ рдПрдХ relay рди рдкрд╛рдард╡рд▓реЗрд▓реНрдпрд╛ outbox rows рд╡рд╛рдЪрддреЛ, рддреНрдпрд╛ publish рдХрд░рддреЛ, рдЖрдгрд┐ рддреНрдпрд╛рдВрдирд╛ sent рдореНрд╣рдгреВрди рдЦреВрдг рдХрд░рддреЛ. - CDC (change data capture) тАФ outbox рд▓рд╛ poll рдХрд░рдгреНрдпрд╛рдРрд╡рдЬреА, database рдЪрд╛рдЪ change log (PostgreSQL WAL, MySQL binlog, DynamoDB Streams) Debezium рд╕рд╛рд░рдЦреНрдпрд╛ tool рдиреЗ рд╡рд╛рдЪрд╛, рдЖрдгрд┐ рддрд┐рдереВрди publish рдХрд░рд╛.
- рдкреБрдиреНрд╣рд╛ at-least-once тАФ relay publish рдХрд░реВрди row рд▓рд╛ sent рдЕрд╢реА рдЦреВрдг рдХрд░рдгреНрдпрд╛рдЖрдзреАрдЪ crash рд╣реЛрдК рд╢рдХрддреЛ, рдореНрд╣рдгреВрди рддреЛ рдкреБрдиреНрд╣рд╛ publish рдХрд░рддреЛ. Consumers idempotent рд╣рд╡реЗрдд (process рдХреЗрд▓реЗрд▓реЗ event ids рдареЗрд╡рд╛).
- Fan-out on write (push) тАФ post рдЭрд╛рд▓реА рдХреА рддреА рдкреНрд░рддреНрдпреЗрдХ follower рдЪреНрдпрд╛ inbox рдордзреНрдпреЗ copy рдХрд░рд╛. рд╡рд╛рдЪрдгреЗ рдореНрд╣рдгрдЬреЗ рдПрдХ lookup. Post рдХрд░рд╛рдпрд▓рд╛ followers рдЗрддрдХреЗ writes рд▓рд╛рдЧрддрд╛рдд.
- Fan-out on read (pull) тАФ post рдПрдХрджрд╛рдЪ рд╕рд╛рдард╡рд╛; рдкреНрд░рддреНрдпреЗрдХ рд╡рд╛рдЪрдХ рддреЛ follow рдХрд░рдд рдЕрд╕рд▓реЗрд▓реНрдпрд╛ рд╕рдЧрд│реНрдпрд╛рдВрдХрдбреВрди рдЧреЛрд│рд╛ рдХрд░рддреЛ. Post рдХрд░рдгреЗ рдореНрд╣рдгрдЬреЗ рдПрдХ write. рд╡рд╛рдЪрд╛рдпрд▓рд╛ follows рдЗрддрдХреЗ lookups рд▓рд╛рдЧрддрд╛рдд.
- Celebrity hybrid тАФ рд╕рд╛рдорд╛рдиреНрдп рд▓реЗрдЦрдХрд╛рдВрд╕рд╛рдареА push, рдкреНрд░рдЪрдВрдб рд╢реНрд░реЛрддреГрд╡рд░реНрдЧ рдЕрд╕рд▓реЗрд▓реНрдпрд╛ рд▓реЗрдЦрдХрд╛рдВрд╕рд╛рдареА pull (рдЗрдереЗ тЙе 100,000 followers), рдЖрдгрд┐ рд╡рд╛рдЪрдХ feed рдЙрдШрдбрддреЛ рддреЗрд╡реНрд╣рд╛ рджреЛрдиреНрд╣реА рдПрдХрддреНрд░ рдХрд░рд╛.
ЁЯдФ рдХрд╛
рдХрд╛рд░рдг рдХрд╛рдЧрдж (рдзрдбрд╛ 01) рд╕рд╛рдВрдЧрддреЛ "'posted' рджрд╛рдЦрд╡рд▓реНрдпрд╛рдирдВрддрд░ рд╕реВрдЪрдирд╛ рдХрдзреАрд╣реА рд╣рд░рд╡рддрд╛ рдХрд╛рдорд╛ рдирдпреЗ" тАФ рдЖрдгрд┐ рддреНрдпрд╛рдд рддрд┐рдЪрд╛ SMS рдкрдг рдпреЗрддреЛ. Dual write рдордзреНрдпреЗ, рдПрдХрд╛ рдЫреЛрдЯреНрдпрд╛ рдЦрд┐рдбрдХреАрддрд▓рд╛ рдХреЛрдгрддрд╛рд╣реА crash рд╣реЗ рд╡рдЪрди рдореЛрдбрддреЛ, рдЖрдгрд┐ рдХреБрдареЗрд╣реА error рджрд┐рд╕рдд рдирд╛рд╣реА. Outbox рдореБрд│реЗ event рд╕реВрдЪрдиреЗрдЗрддрдХрд╛рдЪ рдЯрд┐рдХрд╛рдК рд╣реЛрддреЛ. Pub/sub рдореБрд│реЗ notice API рд▓рд╣рд╛рди рд░рд╛рд╣рддреЛ: SMS, email, search рдЖрдгрд┐ analytics рдЕрд╕реНрддрд┐рддреНрд╡рд╛рдд рдЖрд╣реЗрдд рд╣реЗ рддреНрдпрд╛рд▓рд╛ рдорд╛рд╣реАрдд рдирд╕рддреЗ, рдореНрд╣рдгреВрди рдкреНрд░рддреНрдпреЗрдХ рдЬрдг рд╕реНрд╡рддрдВрддреНрд░рдкрдгреЗ рдмрджрд▓реВ, рдЕрдкрдпрд╢реА рд╣реЛрдК рдЖрдгрд┐ scale рд╣реЛрдК рд╢рдХрддреЛ.
ЁЯФз рдХрд╕реЗ (рдпрд╛ repo рдордзреНрдпреЗ)
- design/blocks.py рдордзреАрд▓
dual_write(crash_between)рджреЛрди writes рдЪреНрдпрд╛ рдордзреНрдпреЗ crash рдЭрд╛рд▓рд╛ рддрд░("notice saved", None)рдкрд░рдд рдХрд░рддреЛ. Outbox.place(id, crash_before_publish)рд╕реВрдЪрдирд╛ рдЖрдгрд┐ рдПрдХ("NoticePosted", id)row рдПрдХрддреНрд░ рдЬреЛрдбрддреЛ (рддреЛ рдПрдХрдЪ transaction);relay()outbox rowspublishedрдордзреНрдпреЗ рд╣рд▓рд╡рддреЛ.fan_out(author_followers, reader_follows, posts_per_day, reads_per_day, mode)push рдЖрдгрд┐ pull рд╕рд╛рдареА рдПрдХрд╛ рджрд┐рд╡рд╕рд╛рдЪреЗ writes рдЖрдгрд┐ read lookups рдореЛрдЬрддреЛ.- design/designs.py рдордзреАрд▓
notify(followers, posts_per_day)100,000 followers рд▓рд╛ pull рдХрдбреЗ рд╡рд│рддреЛ.
ЁЯзк рдХрд░реВрди рдкрд╛рд╣рд╛
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 рдХрд░рддреЛ.
тЪая╕П рдиреЗрд╣рдореАрдЪреНрдпрд╛ рдЪреБрдХрд╛
- outbox рд╢рд┐рд╡рд╛рдп "рдЖрдзреА save, рдордЧ publish" (рдХрд┐рдВрд╡рд╛ "рдЖрдзреА publish, рдордЧ save") тАФ рдХреНрд╡рдЪрд┐рдд рд╣реЛрдгрд╛рд░рд╛, рдЧреБрдкрдЪреВрдк data loss
- events рд▓рд╛ commands рд╕рд╛рд░рдЦреА рдирд╛рд╡реЗ (
SendEmail) тАФ рдЖрддрд╛ producer рд▓рд╛ рдкреБрдиреНрд╣рд╛ рддреНрдпрд╛рдЪреЗ consumers рдорд╛рд╣реАрдд рдЕрд╕рддрд╛рдд - рдкреВрд░реНрдг database row рд╡рд╛рд╣реВрди рдиреЗрдгрд╛рд░реЗ рдкреНрд░рдЪрдВрдб events, рдХрд┐рдВрд╡рд╛ рдкреНрд░рддреНрдпреЗрдХ consumer рд▓рд╛ рдкрд░рдд call рдХрд░рд╛рдпрд▓рд╛ рд▓рд╛рд╡рдгрд╛рд░реЗ рдЕрдЧрджреА рдЫреЛрдЯреЗ events
- idempotent рдирд╕рд▓реЗрд▓реЗ consumers тАФ relay рдХрд╛рд╣реА events рдирдХреНрдХреАрдЪ рджреЛрдирджрд╛ publish рдХрд░реЗрд▓
- рдХреЛрдгреАрдЪ рд╕рд╛рдл рди рдХрд░рдгрд╛рд░реЗ outbox table тАФ рдкрд╛рдард╡рд▓реЗрд▓реНрдпрд╛ rows рдХрд╛рд╣реА рджрд┐рд╡рд╕ рдареЗрд╡рд╛, рдордЧ delete рдХрд░рд╛
- celebrity рдЪреА post рд▓рд╛рдЦреЛ inboxes рдордзреНрдпреЗ push рдХрд░рдгреЗ
- рд╕рдВрдкреВрд░реНрдг topic рд╡рд░ global ordering рдЕрдкреЗрдХреНрд╖рд┐рдд рдзрд░рдгреЗ тАФ ordering рдлрд╛рд░ рддрд░ рдкреНрд░рддреНрдпреЗрдХ key (рдкреНрд░рддреНрдпреЗрдХ рд╡рд░реНрдЧ) рдкреБрд░рддреЗ рдЕрд╕рддреЗ
ЁЯПн рдкреНрд░рддреНрдпрдХреНрд╖ рд╡рд╛рдкрд░рд╛рдд
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