34-წუთიანი სამუშაო სესია მონაცემთა პაიპლაინებზე: ამოღება, გარდაქმნა და ტვირთვა; რატომ შეცვალა ELT-მ თამაშის წესები; პარტია vs სტრიმინგი; და ის იდემპოტენტურობა, ორკესტრირება და ხარისხის შემოწმებები, რომლებიც პაიპლაინს ისეთს ხდის, რომ მშვიდად დაიძინოთ.
შეკვეთები აპლიკაციის ბაზაშია, კლიენტები — CRM-ში, ჩამოჭრები — გადახდების პროცესორში, სარეკლამო ხარჯი — ათეულ SaaS დაშბორდში. ყოველი სასარგებლო შეკითხვა — "რა იყო გასული კვირის მარჟა რეგიონების მიხედვით?" — მოითხოვს, რომ ეს ყველაფერი შეერთდეს, გაიწმინდოს და სანდო იყოს. პაიპლაინი სწორედ ისაა, რაც ამას განზრახ აკეთებს და არა ხელით.
წყარო-სისტემა, თითოეული თავისი სქემით, ფორმატითა და უცნაურობებით.
ადგილი, სადაც ანალიტიკოსები და დაშბორდები ნამდვილად აგზავნიან შეკითხვებს — საწყობი.
ახალი სტრიქონები მუდმივად მოდის; გუშინდელი ექსპორტი უკვე მოძველებულია.
სქემები იცვლება, API-ები ახლდება. დააპროექტეთ მტვრევადობაზე და არა იდეალურ სტოპკადრზე.
ბევრი არეული წყარო ერთ პაიპლაინში იყრის თავს და ერთ საწყობში ეშვება, რომელსაც მთელი კომპანია ენდობა.
// Monday morning: finance needs last week's revenue const rows = await prod.query("SELECT * FROM orders") // on the live DB sendEmail("finance@co", toCSV(rows)) // no history · no schema check · breaks on a rename · run by a human
extract("orders", { from: "replica" }) // not prod .transform(cleanRevenue) // typed + tested .load("warehouse.revenue") // versioned, full history // scheduled · monitored · re-runnable · alerts on failure
ყოველი პაიპლაინი — ინსტრუმენტი რაც უნდა იყოს — სწორედ ამ სამ რამეს აკეთებს. სიტყვები რომ ერთხელ დაალაგოთ, დანარჩენი საუბარი მხოლოდ იმაზეა, სად და როდის უშვებთ მათ.
ამოიღეთ ჩანაწერები მონაცემთა ბაზებიდან, API-ებიდან, ფაილებიდან ან მოვლენების ნაკადებიდან. ორი რთული ადგილია: არ გადატვირთოთ წყარო და აიღოთ მხოლოდ ის, რაც ახალია — აწარმოეთ watermark (ბოლო id ან დროის ნიშნული), რომ დელტა წაიკითხოთ და არა მთელი ცხრილი.
მიუთითეთ ტიპები, გაასწორეთ კოდირებები, წაშალეთ დუბლიკატები, მიაერთეთ საცნობარო მონაცემები, გამოიყენეთ ბიზნესწესები. სწორედ აქ გადადის მონაცემი ნედლიდან სანდომდე — და აქვე ცხოვრობს ბაგების უმეტესობა, ამიტომ ეს ნაწილი ტიპიზებული და დატესტილი უნდა იყოს.
დასვით შედეგი საწყობში. ბრმა ჩასმის ნაცვლად ამჯობინეთ upsert, წერეთ პარტიებად და გახადეთ ჩაწერა იდემპოტენტური, რომ ხელახალმა მცდელობამ ორჯერ ვერ დათვალოს (მე-5 ნაწილი).
// transform tangled into the read, no types rows.forEach(r => { r.amount = r.amount + "" // stringify money → lose precision r.date = r.created.slice(0,10) // assumes a format that'll change }) // bad rows pass through silently and poison the warehouse
function clean(o: RawOrder): Order { return { amount: Money.parse(o.amount), // typed money placedAt: Date.parse(o.created_at), // throws if invalid region: regionOf(o.country), // validated lookup } }
ჰგავს საწყობის მიმღებ დოკს: ამოღება არის ჩამოსული სატვირთო, გარდაქმნა — საქონლის შემოწმება და მონიშვნა, ტვირთვა კი — სწორ თაროზე დალაგება.
ETL და ELT ერთსა და იმავე სამ ეტაპს იყენებს — უბრალოდ ვერ თანხმდებიან იმაზე, როდის უნდა გარდაქმნა. ეს ერთი გადანაცვლება თანამედროვე data engineering-ის ყველაზე დიდი ცვლილებაა და მთლიანად იმით არის განპირობებული, რაც საწყობს დღეს შეუძლია.
ETL ცალკე მანქანაზე გარდაქმნის და მხოლოდ მზა ფორმას ტვირთავს. ELT ნედლ მონაცემებს დებს, შემდეგ კი გარდაქმნას საწყობს ატარებინებს SQL-ით.
თანამედროვე ნაგულისხმევი არჩევანი ELT-ია: ამოიღეთ, ჩატვირთეთ ნედლად, შემდეგ საწყობში გარდაქმენით ვერსიების კონტროლში მყოფი SQL-ით. ETL-ს მაშინ მიმართეთ, როცა პრივატულობა, მოცულობა ან სამიზნე გარდაქმნას უფრო ადრე აიძულებს.
ETL vs ELT იმაზე იყო, სად გარდაქმნით. ეს ცალკე შეკითხვაა: რამდენად ხშირად მოძრაობს მონაცემი. ელოდებით და დიდ გროვას განრიგით გადაიტანთ (ვთქვათ, ყოველ ღამე), თუ ყოველ ახალ ჩანაწერს იმავე წამში ამუშავებთ? კომპანიების უმეტესობას ბოლოს ორივე სჭირდება — სხვადასხვა შეკითხვა სხვადასხვა დაგვიანებას ითმენს.
პარტია ელოდება და შემდეგ მთელ ფანჯარას ერთბაშად ამუშავებს. სტრიმინგი ყოველ მოვლენაზე მისვლისთანავე მოქმედებს.
ქსელი წყდება, ნოუდები კვდებიან, დეპლოი სამუშაოს 90%-ზე კლავს. საკითხი ისაა, არა გავიმეოროთ თუ არა — არამედ ის, გააფუჭებს თუ არა გამეორება თქვენს მონაცემებს. ამ თვისებას სახელი აქვს: იდემპოტენტურობა.
გამეორებული INSERT 42-ე შეკვეთას ორჯერ თვლის. order_id-ზე გასაღებული MERGE ყოველთვის ერთსა და იმავე ერთ სტრიქონს დებს.
INSERT-ის ნაცვლად upsert/MERGE ბიზნეს-გასაღებზე.now() ან შემთხვევითი id-ები არ უნდა იყოს.async function load(batch: Row[]) { for (const row of batch) await db.insert("orders", row) // blind append } // crash after row 800 of 1000 → retry re-inserts 1–800
async function load(batch: Row[]) { await db.upsert("orders", batch, { key: "order_id" // same id ⇒ same row }) } // retry from row 1 → rows 1–800 just overwrite themselves
დროებითი ჩავარდნისას ორკესტრატორი უბრალოდ ხელახლა უშვებს ამოცანას. რადგან ტვირთვა იდემპოტენტურია, ავტომატური გამეორება უსაფრთხოა — ღამის 3 საათზე ადამიანს არავინ იძახებს, რომ ჯერ "დუბლები დაასუფთაოს".
ბექფილი პაიპლაინს წარსულ თარიღებზე უშვებს ხელახლა — ბაგის გასასწორებლად, ახალი სვეტის დასამატებლად ან ისტორიის ჩასატვირთად. იდემპოტენტური, დროით დანაწილებული ნაბიჯები საშუალებას გაძლევთ გასული მარტი თავიდან გაატაროთ ისე, რომ სხვა არაფერს შეეხოთ.
პაიპლაინი ბევრი ურთიერთდამოკიდებული ამოცანაა, რომელიც განრიგით ეშვება და რომელსაც გამეორება, მონიტორინგი და ნდობა სჭირდება. ამას ორი დისციპლინა აწესრიგებს: ორკესტრირება და მონაცემთა ხარისხი.
ყოველი ნოუდი ამოცანაა; წიბოები — დამოკიდებულებები. ხარისხის check ბლოკავს publish-ს — ცუდი მონაცემი მომხმარებლამდე არასოდეს აღწევს.
კოდის ტესტები თქვენს გარდაქმნის ლოგიკას ამტკიცებს. ხარისხის შემოწმებები კი ამტკიცებს, რომ ყოველ გაშვებაზე მონაცემი საღია — და პაიპლაინს კარიბჭეს უყენებს, რომ ცუდი მონაცემი ხმამაღლა ჩავარდეს და არა ჩუმად გაჟონოს დაშბორდებში.
სვეტები არსებობს, ტიპები ემთხვევა, მოულოდნელი გადარქმევები არაა. თუ წყარომ ველი დაამატა ან წაშალა, გაშვება უნდა ჩავარდეს და არა ჩუმად შეავსოს რეპორტი null-ებით.
უახლესი სტრიქონი SLA-ზე ახალგაზრდაა. მოძველებული მონაცემი ჩუმი ჩავარდნაა — დაშბორდი კვლავ იხატება, უბრალოდ ტყუის.
დღევანდელი პარტია ნორმის დასაშვებ ფარგლებშია. გაშვება, რომელიც 0 სტრიქონს ტვირთავს (ან 50×-ს), ჩვეულებრივ გაფუჭებულ წყაროს ნიშნავს და არა რეკორდულ გაყიდვების დღეს.
პირველადი გასაღები მართლაც უნიკალურია (იჭერს მე-5 ნაწილის დუბლიკატების ბაგს) და not-null სვეტები არასოდეს არის null. ყველაზე იაფი შემოწმებები, ყველაზე მაღალი დაჭერის მაჩვენებლით.
შეკვეთებში ყოველი customer_id არსებობს customers ცხრილში. ობოლი გასაღებები ჩუმად ქრება inner join-ში — და შემოსავალიც მათთან ერთად ქრება.
Lineage არის რუკა იმისა, თუ როგორ გამოიყვანეს ყოველი ცხრილი და სვეტი თავისი წყაროებიდან. როცა დაშბორდი არასწორად გამოიყურება, lineage საშუალებას გაძლევთ ის ყოველი გარდაქმნის გავლით დამნაშავე წყარომდე მიადევნოთ — და დაინახოთ, კიდევ რა იმტვრევა, თუ სვეტს შეცვლით.
ამ ყველაფერს თითქმის არასოდეს აწყობთ ხელით. თანამედროვე სტეკი, როგორც წესი, სამი საქმეა, სპეციალიზებულ ინსტრუმენტებზე გადანაწილებული: კონექტორი ამოღებისა და ტვირთვისთვის, dbt გარდაქმნისთვის და ერთი ორკესტრატორი, რომ ეს ყველაფერი განრიგით გაეშვას. აი, სად დგანან წამყვანი ინსტრუმენტები და რაში არის თითოეული კარგი (და ცუდი).
ორკესტრატორი (ქვემოთ) დირიჟორობს; კონექტორები E+L-ს ართმევენ თავს, dbt საწყობის შიგნით T-ს იბარებს, ხარისხის შემოწმებები კი წყვეტენ, რა მოხვდება დაშბორდებამდე.
დადებითი — ინდუსტრიის ნაგულისხმევი არჩევანი: უზარმაზარი საზოგადოება, თითქმის ყველგან ეშვება, ინტეგრაცია ყველაფრისთვის.
უარყოფითი — მძიმე Python-ბოილერპლეიტი და ის ამოცანებს გეგმავს ისე, რომ თქვენს მონაცემებს სინამდვილეში ვერ იგებს.
დადებითი — მონაცემთა აქტივებით (ცხრილები, რომლებსაც აწარმოებთ) აზროვნებს, ტიპიზებული შემავალი/გამომავალი მონაცემებით და შესანიშნავი ლოკალური ტესტირებით.
უარყოფითი — Airflow-ზე ახალგაზრდა და პატარაა და აზროვნების ახალი წესის სწავლას მოითხოვს.
დადებითი — მსუბუქი და წმინდა Python; არსებული სკრიპტის მართულ ნაკადად ქცევა სწრაფია.
უარყოფითი — ნაკლები რამაა ყუთშივე ჩადებული; გარშემო სტეკის მეტ ნაწილს თავად აწყობთ.
დადებითი — ასობით მზა კონექტორი (წინასწარ აგებული წამკითხველები წყაროებისთვის), ღია კოდით და თვითჰოსტირებადი.
უარყოფითი — თავად გაშვება შეიძლება მძიმე იყოს და კონექტორების ხარისხი წყაროდან წყაროზე იცვლება.
დადებითი — სრულად მართული კონექტორები: ერთხელ მოაწყობთ და შემდეგ პრაქტიკულად აღარაფერს სჭირდება მოვლა.
უარყოფითი — ფასი გადატანილი მონაცემების მოცულობაზეა მიბმული, ამიტომ დიდ მოცულობაზე ანგარიში სწრაფად იზრდება.
დადებითი — გარდაქმნები ვერსიების კონტროლში მყოფი SQL-ის სახით, ჩაშენებული ტესტებითა და დოკუმენტაციით, ასე რომ მათი ფლობა ანალიტიკოსებსაც შეუძლიათ.
უარყოფითი — მხოლოდ გარდაქმნა: თავად არც ამოიღებს, არც ჩატვირთავს და არც განრიგში ჩასვამს — ზემოთ ჩამოთვლილი ინსტრუმენტები სჭირდება.
ერთი და იგივე საქმე — შეკვეთების საწყობში გადატანა — ჯერ ისე დაწერილი, როგორც ჩვეულებრივ იწყება, შემდეგ კი ისე, როგორც უნდა დასრულდეს.
// runs from a laptop, no schedule, no checks const rows = await prod.query("SELECT * FROM orders") // full table, on prod rows.forEach(r => warehouse.insert(r)) // blind append // crash midway ⇒ dups · rename ⇒ silent nulls // no history · no alert · no idea if it's right
extract("orders", { from: "replica", since: watermark }) .transform(clean) // typed · tested .check([schema, freshness, unique]) // gate .load({ key: "order_id", mode: "merge" }) // idempotent // orchestrated · incremental · re-runnable · alerts on fail
"გადაიტანეთ მონაცემები ერთხელ, განზრახ და ისე, რომ ხელახლა უსაფრთხოდ გაუშვათ."
— მთელი მოხსენება, შეკუმშული
ხუთი სწრაფი შეკითხვა ETL vs ELT-ზე, პარტია vs სტრიმინგზე, იდემპოტენტურობასა და მონაცემთა ხარისხზე — მყისიერი პასუხი, ავტორიზაციის გარეშე.
ნავიგაცია ← → ღილაკებით ან სქროლით · უკან ბიბლიოთეკაში