36-წუთიანი სამუშაო სესია ასინქრონულ შეტყობინებებზე: რატომ ვსვამთ ბროკერს სერვისებს შორის, რა განსხვავებაა სამუშაო რიგსა და გამეორებად ლოგს შორის, Kafka-ს თემები და დანაწილებები, მიწოდების გარანტიები, რომლებსაც მართლა აქვს მნიშვნელობა, და პატერნები, რომლებიც ნაკადს პროდაქშენში წესრიგში ინახავს.
ჩექაუთი დასრულდა. ახლა ხუთი რამ უნდა მოხდეს — ბარათის ჩამოჭრა, ქვითრის გაგზავნა ელფოსტით, მარაგების განახლება, ანალიტიკის შეტყობინება, რეკომენდაციების გათბობა. თუ ჩექაუთი თითოეულს რიგრიგობით იძახებს და ელოდება, ის იმდენად სწრაფია, რამდენადაც ყველაზე ნელი, და იმდენად ხელმისაწვდომი, რამდენადაც ყველაზე არასტაბილური. შეტყობინებების ბროკერი ამ ჯაჭვს წყვეტს: ჩექაუთი აცხადებს „შეკვეთა განთავსდა“ და საქმეს აგრძელებს.
პირდაპირი გამოძახებები სერვისებს ერთ მყიფე ჯაჭვად კრავს. ბროკერი ჩექაუთს აძლევს საშუალებას, მოვლენა ერთხელ გამოაცხადოს, და თითოეულ კონსიუმერს — თავის დროზე იმოქმედოს.
async function placeOrder(o) { await save(o) await email.sendReceipt(o) // blocks… await inventory.reserve(o) // blocks… await analytics.track(o) // if this is down, checkout fails }
async function placeOrder(o) { await save(o) await broker.publish("order.placed", o) // one write, returns fast } // email, inventory, analytics each consume on their own // one of them being down can't fail the checkout
ჰგავს განსხვავებას იმას შორის, თითოეულ კოლეგას დაურეკოთ და ხაზზე ელოდოთ, თუ გუნდის დაფაზე ერთი ჩანაწერი გამოაკრათ, რომელსაც ყველა მაშინ წაიკითხავს, როცა თავისუფალია.
„შეტყობინებების ბროკერი“ ორ სრულიად განსხვავებულ დიზაინს აერთიანებს. არასწორს თუ აირჩევთ, სამუდამოდ ინსტრუმენტს შეებრძოლებით. გაყოფა მარტივია: შეტყობინება ქრება მას შემდეგ, რაც ვინმე დაამუშავებს, თუ რჩება, რომ ნებისმიერმა ხელახლა წაიკითხოს?
რიგი შლის m1-ს, როგორც კი მუშა მას ადასტურებს. ლოგი ყველა შეტყობინებას ინახავს; მკითხველები A და B სხვადასხვა ოფსეტზე დგანან და უკან დახევა შეუძლიათ.
order.placed ნაკადს დამოუკიდებლად კითხულობენ.Apache Kafka გამეორებადი ლოგის საეტალონო დიზაინია, და მისი ოთხი ცნება ხსნის, როგორ მასშტაბირდება ლოგი წამში მილიონობით შეტყობინებამდე და მაინც ინარჩუნებს თანმიმდევრობას იქ, სადაც ეს მნიშვნელოვანია. ისწავლეთ ეს ოთხი და ნაკადების სისტემების უმეტესობა ერთნაირად წაიკითხება.
order.placed). დანაწილება — ერთი თემა იყოფა N დალაგებულ ლოგად, და სწორედ დანაწილებები აძლევს თემას მანქანებზე მასშტაბირების საშუალებას. ოფსეტი — პოზიცია დანაწილების შიგნით. კონსიუმერთა ჯგუფი — კონსიუმერების ერთობლიობა, რომლებიც თემის სამუშაოს ინაწილებენ, სადაც თითოეულ დანაწილებას ჯგუფის ზუსტად ერთი წევრი კითხულობს.პროდიუსერი თითოეული შეტყობინების გასაღებს ჰეშავს, რომ დანაწილება აირჩიოს. billing ჯგუფი სამ დანაწილებას თავის სამ კონსიუმერზე ანაწილებს — თითოს თითო.
customer_id) → ერთი დანაწილება → ამ კლიენტისთვის თანმიმდევრობა დაცულია. გასაღების გარეშე → რიგრიგობით.// producer — the key pins related events to one partition await producer.send({ topic: "order.placed", messages: [{ key: order.customerId, value: JSON.stringify(order) }], }) // consumer — joins a group; Kafka assigns it partitions await consumer.subscribe({ topic: "order.placed" }) await consumer.run({ groupId: "billing", eachMessage: async ({ message }) => charge(message), })
ორი ჯგუფი ერთ თემაზე: billing და analytics თითოეული სრულ ნაკადს იღებს, დამოუკიდებლად.
ქსელი handshake-ის შუაში წყდება, ამიტომ ბროკერი ვერასოდეს იქნება დარწმუნებული, რომ კონსიუმერმა დაასრულა. ერთადერთი გულწრფელი კითხვაა, რომელი წარუმატებლობა გირჩევნიათ: შეტყობინების დაკარგვა თუ მისი ორჯერ მიწოდება. პასუხს განსაზღვრავს ის, როდის აკომიტებთ ოფსეტს სამუშაოს შესრულებასთან შედარებით.
დააკომიტეთ ოფსეტი ჯერ და კრახი შეტყობინებას დაკარგავს. დააკომიტეთ ბოლოს და კრახი მას ხელახლა მოიტანს. სატრანსპორტო შრეზე მესამე ვარიანტი არ არსებობს.
eachMessage(msg) { const o = JSON.parse(msg.value) ledger.addCharge(o.amount) // redelivery ⇒ charged twice } // at-least-once + non-idempotent = double billing
eachMessage(msg) { const o = JSON.parse(msg.value) ledger.upsertCharge({ id: o.orderId, // dedupe key amount: o.amount, }) // same id ⇒ same single row }
შეტყობინებების დინების ამუშავებას ერთი შუადღე სჭირდება. მათი თანმიმდევრობის შენარჩუნება, აფეთქებების ჩაწოვა, იმ შეტყობინების დამუშავება, რომელიც ყოველთვის ვარდება, და ბაგის შემდეგ გამეორება — აი, ეს არის ნამდვილი სამუშაო. ხუთი კიდე, რომელსაც ყველა გუნდი ხვდება.
თანმიმდევრობას მხოლოდ დანაწილების შიგნით იღებთ. თუ ერთი შეკვეთის ორი მოვლენა სხვადასხვა დანაწილებაში მოხვდება, კონსიუმერმა შეიძლება shipped დაინახოს paid-ზე ადრე. მიმართეთ სტაბილური გასაღებით, რომ ერთი ერთეულის ყველაფერი ერთ დანაწილებაზე დარჩეს.
producer.send({ topic: "order.events", messages: [{ value: evt }] }) // round-robin partition // paid & shipped may land in different partitions
producer.send({ topic: "order.events", messages: [{ key: evt.orderId, value: evt }] }) // same orderId ⇒ same partition ⇒ ordered
როცა პროდიუსერები უფრო სწრაფად წერენ, ვიდრე კონსიუმერები კითხულობენ, სხვაობა იზრდება. ამ სხვაობას სახელი აქვს — კონსიუმერის ჩამორჩენა: რამდენი შეტყობინებით ჩამორჩება ჯგუფი ლოგის თავს. რიგში გადატვირთული კონსიუმერი უკუწნევას გრძნობს; ლოგში ის უბრალოდ კიდევ უფრო ჩამორჩება.
შხამიანი შეტყობინება — დაზიანებული ან ისეთი, რომელიც ყოველთვის შეცდომას აგდებს — უსასრულოდ განმეორდება და დაბლოკავს ყველაფერს, რაც მის უკან დგას. შეზღუდეთ ხელახალი მცდელობები, შემდეგ კი გადაიტანეთ dead-letter რიგში (DLQ), რომ ნაკადის დანარჩენი ნაწილი დინებას აგრძელებდეს და უარყოფილი შეტყობინება ადამიანმა მოგვიანებით შეამოწმოს.
eachMessage(msg) { try { handle(msg) } catch (e) { if (msg.attempts >= 5) await dlq.send(msg) // quarantine else throw e // retry } }
ტრანსფორმაციის ბაგი გაუშვით პროდაქშენში? რადგან ლოგი ისტორიას ინახავს, შეგიძლიათ ჯგუფის ოფსეტი დააბრუნოთ და ხელახლა წაიკითხოთ. ახალი კონსიუმერი, რომელსაც სრული წარსული სჭირდება? დაიწყეთ ოფსეტ 0-დან. გამეორება მხოლოდ მაშინ მუშაობს, თუ თქვენი კონსიუმერები იდემპოტენტურია (ნაწილი 4) — თორემ დუბლიკატებსაც გაიმეორებთ.
პროდიუსერები და კონსიუმერები დამოუკიდებლად დეპლოიდებიან, ამიტომ გადარქმეული ველი ჩუმად ტეხს ქვემოთ მდგარ მკითხველებს. დაარეგისტრირეთ შეტყობინებების სქემები და დანერგეთ თავსებადი ევოლუცია (დაამატეთ არასავალდებულო ველები, არასოდეს გამოიყენოთ არსებული სხვა დანიშნულებით) — სქემების რეესტრი უარყოფს შეუთავსებელ ცვლილებას, სანამ ის პროდაქშენში მოხვდება; ეს პაიპლაინების დეკის სქემის შემოწმებების ნაკადური ბიძაშვილია.
რამდენიმე წარუმატებელი მცდელობის შემდეგ შხამიანი შეტყობინება dead-letter რიგში გადაინაცვლებს — ნაკადის დანარჩენი ნაწილი არასოდეს იბლოკება.
ბაზარი სუფთად იყოფა მე-2 ნაწილის რიგი-vs-ლოგი ხაზზე, რომელსაც ჯვარედინად ედება ის, თუ რამდენის ოპერირება გსურთ თავად. აი, წამყვანი სისტემები — თითოეული ერთსტრიქონიანი ძლიერი მხარითა და ხაფანგით.
დადებითი — მაღალი გამტარუნარიანობის მოვლენების ნაკადებისთვის დე-ფაქტო სტანდარტი; უზარმაზარი ეკოსისტემა, ნამდვილი გამეორება, exactly-once ნაკადების დამუშავებისთვის.
უარყოფითი — თვითჰოსტინგი და აწყობა ოპერაციულად მძიმეა; მარტივი ამოცანების რიგებისთვის ზედმეტია.
დადებითი — მოწიფული, მოქნილი შეტყობინებების ბროკერი მდიდარი მარშრუტიზაციით; იდეალურია ამოცანების რიგებისა და მოთხოვნა/პასუხისთვის.
უარყოფითი — გამეორებადი ლოგი არ არის; დადასტურებული შეტყობინება ქრება, დიდ მასშტაბზე კი გამტარუნარიანობით Kafka-ს ჩამორჩება.
დადებითი — სრულად მართული, თითქმის ნულოვანი ოპერაციებით: SQS მარტივი რიგებისთვის, Kinesis ლოგის/ნაკადის ფორმისთვის.
უარყოფითი — AWS-ზე მიბმა; Kafka-ზე ნაკლები ფუნქციონალი და ეკოსისტემა, ღირებულება კი მოხმარებასთან ერთად იზრდება.
დადებითი — ერთ სისტემაში აკეთებს რიგებსაც და ნაკადებსაც, ჩაშენებული გეო-რეპლიკაციითა და საფეხურებრივი საცავით.
უარყოფითი — მეტი მოძრავი ნაწილი (ეყრდნობა ცალკე შენახვის შრეს); Kafka-ზე პატარა საზოგადოება.
დადებითი — Kafka-ს API-ზე საუბრობს, მაგრამ ერთი ბინარია (არც JVM, არც ცალკე კოორდინატორი); უფრო მარტივი გასაშვები, ნაკლები შეყოვნებით.
უარყოფითი — ახალგაზრდა პროექტი და პატარა ეკოსისტემა; ბირთვი ღიაა, ზოგი ფუნქციონალი კომერციული.
უკვე AWS-ზე ხართ და ოპერაციები არ გინდათ? SQS / Kinesis. გჭირდებათ ნამდვილი მოვლენების ლოგი და თქვენივე ინფრასტრუქტურა? Kafka (ან Redpanda უფრო მსუბუქი გაშვებისთვის). კლასიკური ამოცანების რიგი მდიდარი მარშრუტიზაციით? RabbitMQ.
მივყვეთ ერთადერთ order.placed მოვლენას ყველაფერში, რაც განვიხილეთ: ერთხელ გამოქვეყნებული, კლიენტის მიხედვით დანაწილებული, დამოუკიდებელი ჯგუფების მიერ წაკითხული, თითოეული იდემპოტენტური, უარყოფილებისთვის კი DLQ.
ერთი გამოქვეყნება; ოთხი დამოუკიდებელი კონსიუმერთა ჯგუფი, თითოეული საკუთარ ოფსეტზე და იდემპოტენტური; წარუმატებლობები dead-letter რიგში გადადის. ხვალ მეხუთე კონსიუმერს დაამატებთ და ოფსეტ 0-დან გაიმეორებთ — ზემოთ არაფერი იცვლება.
„მოვლენა ერთხელ გამოაქვეყნეთ; დაე, ყველამ თავის საათზე წაიკითხოს.“
— მთელი მოხსენება, შეკუმშული
ხუთი სწრაფი კითხვა რიგებსა და ლოგებზე, Kafka-ზე, მიწოდების სემანტიკასა და პროდაქშენის კიდეებზე — მყისიერი პასუხი, ავტორიზაციის გარეშე.
ნავიგაცია ← → ღილაკებით ან სქროლით · უკან ბიბლიოთეკაში