Tere, Habr! Esitlen teile tõlget artiklist autoriteks on Burak Yavuz, Brenner Heintz ja Denny Lee, kes valmistasid selle ette kursuse käivitamise eel OTUS-e poolt.

Andmed, nagu ka meie kogemused, kogunevad ja arenevad pidevalt. Et mitte maha jääda, peavad meie vaimsed mudelid maailmast kohanduma uute andmetega, millest mõned sisaldavad uusi mõõtmeid — uusi viise asjade jälgimiseks, millest me enne teadlikud ei olnud. Need vaimsed mudelid ei erine palju tabeli skeemidest, mis määravad, kuidas me uut teavet klassifitseerime ja töötleme.
See viib meid skeemihalduse küsimuseni. Aja jooksul, kui ärivajadused ja nõudmised muutuvad, muutub ka teie andmete struktuur. Delta Lake võimaldab uusi mõõtmeid andmete muutumisel hõlpsalt rakendada. Kasutajatel on juurdepääs lihtsatele semantikatele oma tabelite skeemide haldamiseks. Need tööriistad sisaldavad skeemi sundimise (Schema Enforcement), mis kaitseb kasutajaid nende tabelite kogemata määrdumise eest vigade või mittevajalike andmetega, samuti skeemi evolutsiooni (Schema Evolution), mis võimaldab automaatselt lisada uusi veerge väärtuslike andmetega vastavatesse kohtadesse. Selles artiklis uurime neid tööriistu põhjalikult.
Tabeli skeemide mõistmine
Iga DataFrame Apache Sparkis sisaldab skeemi, mis määratleb andmete kuju, nagu andmetüübid, veerud ja метаданные. Delta Lake'i abil salvestatakse tabeli skeem JSON-formaadis tehinguajalugu.
Mis on skeemi sundimine?
Skeemi sundimine (Schema Enforcement), tuntud ka kui skeemi valideerimine (Schema Validation), on Delta Lake'is kaitsemehhanism, mis tagab andmete kvaliteedi, lükates tagasi kirjed, mis ei vasta tabeli skeemile. Nagu populaarses restoranis registratuuris tööle asunud hostess, kes võtab vastu ainult eelnevalt broneeritud kohti, kontrollib ta, kas iga tabelisse sisestatava andmeveeru kohta on vastav oodatud veergude nimekiri (teisisõnu, kas igaühe jaoks on olemas «broneering»), ning lükkab tagasi kõik kirjed veergude kohta, mida nimekirjas ei ole.
Kuidas skeemi sundimine töötab?
Delta Lake kasutab salvestamisel skeemi kontrollimist, mis tähendab, et kõik uued kirjed tabelisse kontrollitakse sihttabeli skeemiga ühilduvuse osas salvestamise ajal. Kui skeem ei ühildu, tühistab Delta Lake täielikult tehingu (andmed ei salvestata) ja loob erandi, et teavitada kasutajat vastuolust.
Kirje ühilduvuse määramiseks kasutab Delta Lake järgmisi reegleid. Salvestatav DataFrame:
- ei tohi sisaldada lisa veerge, mida sihttabeli skeemis pole. Ja vastupidi, kõik on korras, kui sisendandmed ei sisalda täielikult kõiki tabeleid — nendele veergudele lihtsalt omistatakse nullväärtused.
- ei tohi omada veergude andmetüüpide erinevusi sihttabeli veergude andmetüüpides. Kui sihttabeli veerg sisaldab andmeid StringType, kuid vastav veerg DataFrame'is sisaldab andmeid IntegerType, sunnib skeemi rakendamine erandi tekkimist ja takistab salvestamisoperatsiooni.
- ei tohi sisaldada veergude nimesid, mis erinevad ainult suur- ja väiketähtedest. See tähendab, et te ei saa ühes tabelis määrata veerge nimedega 'Foo' ja 'foo'. Kuigi Spark'i saab kasutada suuruse tundlikus või tundmatuks (vaikimisi) režiimis, säilitab Delta Lake suuruse, kuid on skeemi salvestamise kontekstis suuruse suhtes tundetu. Parquet on salvestamise ja veeru teabe tagastamise ajal suuruse suhtes tundlik. Võimalike vigade, andmekahjustuste või nende kadumise vältimiseks (millega me oleme isiklikult kokku puutunud Databricks'is) oleme otsustanud selle piirangu lisada.
Selle illustreerimiseks vaatame, mis juhtub allolevas koodis, kui püüame lisada mõned hiljuti loodud veerud Delta Lake'i tabelisse, mis ei ole veel nende vastuvõtmiseks seadistatud.
# Сгенерируем DataFrame ссуд, который мы добавим в нашу таблицу Delta Lake
loans = sql("""
SELECT addr_state, CAST(rand(10)*count as bigint) AS count,
CAST(rand(10) * 10000 * count AS double) AS amount
FROM loan_by_state_delta
""")
# Вывести исходную схему DataFrame
original_loans.printSchema()
root
|-- addr_state: string (nullable = true)
|-- count: integer (nullable = true)
# Вывести новую схему DataFrame
loans.printSchema()
root
|-- addr_state: string (nullable = true)
|-- count: integer (nullable = true)
|-- amount: double (nullable = true) # new column
# Попытка добавить новый DataFrame (с новым столбцом) в существующую таблицу
loans.write.format("delta")
.mode("append")
.save(DELTALAKE_PATH)
Returns:
A schema mismatch detected when writing to the Delta table.
To enable schema migration, please set:
'.option("mergeSchema", "true")'
Table schema:
root
-- addr_state: string (nullable = true)
-- count: long (nullable = true)
Data schema:
root
-- addr_state: string (nullable = true)
-- count: long (nullable = true)
-- amount: double (nullable = true)
If Table ACLs are enabled, these options will be ignored. Please use the ALTER TABLE command for changing the schema.Uute veergude automaatse lisamise asemel kehtestab Delta Lake skeemi ja peatab salvestamise. Et aidata määrata, milline veerg (või nende hulk) põhjustab vastuolu, väljastab Spark mõlemad skeemid jälgimisteatega võrdlemiseks.
Mis kasu on skeemi sundrakendamisest?
Kuna sundrakendamine skeem on piisavalt range kontroll, on see suurepärane tööriist, mida kasutada puhta, täielikult muudetud andmestiku väravavahtina, mis on tootmiseks või tarbimiseks valmis. Seda rakendatakse tavaliselt tabelitele, mis esitavad andmeid otse:
- Masinõppe algoritmid
- BI armatuurlaudadele
- Andmeanalüüsile ja visualiseerimistööriistadele
- Igas produktsioonisüsteemis, mis nõuab rangeid struktureeritud, rangelt tüpiseeritud semantilisi skeeme.
Et valmistada oma andmeid ette selleks viimaseks takistuseks, kasutavad paljud kasutajad lihtsat "multi-hop" arhitektuuri, mis järk-järgult lisab struktuuri nende tabelitesse. Kui soovite sellest rohkem teada, võite tutvuda artikliga
Muidugi võib sundrakendamist skeemid kasutada igas teie andmetöödeldusse, kuid pidage meeles, et voogedastuse kirjutamine tabelisse võib sellisel juhul olla frustreeriv, kuna näiteks võite unustada, et olete lisanud veel ühe veeru sisendandmetesse.
Andmete vedeliku vältimine
Selles etapis võite küsida, mis põhjustab sellist elevust? Lõppude lõpuks, mõnikord võib ootamatu "skeemi mittevastavuse" viga segada teie töövoogu, eriti kui olete Delta Lake'is uus. Miks mitte lihtsalt lasta skeemil muutuda nii, nagu on vajalik oma DataFrame'i kirjutamiseks, vaatamata kõigele?
Nagu vana ütlus ütleb, "unts ennetust maksab naela ravi." Teatud hetkel, kui te ei hooli oma skeemi rakendamisest, toovad andmetüübi ühilduvuse probleemid oma vastikud pead välja - esmapilgul homogeenne toorse andmeallikate all võivad peituda äärmuslikud juhud, kahjustatud veerud, valesti vormistatud kaardistused või muud hirmutavad asjad, mis kätkivad öiseid õudusunenägusid. Parim lähenemine on need vaenlased väravas peatada - sundrakendamise skeemiga - ja tegeleda nendega avatult, mitte hiljem, kui nad hakkavad hiilima teie töötava koodi tumedatesse sügavustesse.
Sunnitud vahevormi rakendamine annab kindlustunde, et teie tabeli skeem ei muutu, kui te ise ei kinnita muudatusi. See takistab andmete "lahjendamist", mis võib juhtuda, kui uusi veerge lisatakse nii sageli, et varem väärtuslikud, kompaktsed tabelid kaotavad oma väärtuse ja kasulikkuse andmemere tõttu. Sunni rakendamine kutsub teid olema teadlik, seadma kõrgeid standardeid ja ootama kõrget kvaliteeti, võimaldades teil jääda kohusetundlikuks ja teie tabelid puhtaks.
Kui otsustate edasise kaalumise käigus, et teie jaoks on tegelikult vajalik uus veerg lisada – pole probleemi, allpool on ühe rea lahendus. Lahendus on skeemi evolutsioon!
Mis on skeemi evolutsioon?
Skeemi evolutsioon on funktsioon, mis võimaldab kasutajatel hõlpsasti muuta tabeli skeemi vastavalt andmetele, mis aja jooksul muutuvad. Seda kasutatakse kõige sagedamini lisamise või ülekirjutamise operatsiooni käitlemisel, et automaatselt kohandada skeemi ühe või mitme uue veeru lisamiseks.
Kuidas skeemi evolutsioon toimib?
Eelmise peatüki näidet järgides saavad arendajad hõlpsasti kasutada skeemi evolutsiooni, et lisada uusi veerge, mis olid varem skeemiga sobimatuse tõttu tagasilükatud. Skeemi evolutsioon aktiveeritakse, lisades .option('mergeSchema', 'true') teie Spark käsule .write või .writeStream.
# Добавьте параметр mergeSchema
loans.write.format("delta")
.option("mergeSchema", "true")
.mode("append")
.save(DELTALAKE_SILVER_PATH)Graafiku nägemiseks käivitage järgmine Spark SQL päring
# Создайте график с новым столбцом, чтобы подтвердить, что запись прошла успешно
%sql
SELECT addr_state, sum(`amount`) AS amount
FROM loan_by_state_delta
GROUP BY addr_state
ORDER BY sum(`amount`)
DESC LIMIT 10 
Alternatiivselt saate selle valiku seada kogu Spark seansi jaoks, lisades spark.databricks.delta.schema.autoMerge = True Spark seadistuses. Aga olge ettevaatlik, sest sunnitud skeemi rakendamine ei hoiataks teid enam soovimatute skeemiga mittevastavuste eest.
Käitumise päringusse lisamine mergeSchema, kõik veerud, mis on olemas DataFrames, kuid puuduvad sihttabelis, lisatakse automaatselt skeemi lõpuks kirje tehingu raames. Samuti võivad lisanduda sisemise väljad, mis lisatakse samuti vastavate veergude struktuuri lõppu.
Kuupäeva insenerid ja teadlased saavad kasutada seda valikut, et lisada oma olemasolevatesse masinõppe tootmistabelitesse uusi veerge (nt hiljuti jälgitav mõõdik või selle kuu müügiväärtuste veerg) ilma, et rikuks olemasolevaid mudeleid, mis põhinevad vanadel veergudel.
Järgnevad skeemi muutmise tüübid on lubatud skeemi evolutsiooni raames tabeli lisamisel või ülekirjutamisel:
- Uute veergude lisamine (see on kõige levinum stsenaarium)
- Andmetüüpide muutmine NullType -> mõni teine tüüp või tõstmine ByteType -> ShortType -> IntegerType
Teised muutused, mis ei ole lubatud skeemi evolutsiooni raames, nõuavad, et skeem ja andmed kirjutataks üle, lisades .option("overwriteSchema", "true"). Näiteks, kui veerg "Foo" oli algselt integer ja uus skeem oleks stringi andmetüüp, siis kõik Parquet failid (andmed) tuleks ümber kirjutada. Selliste muudatuste hulka kuuluvad:
- veeru eemaldamine
- olemasoleva veeru andmetüübi muutmine (kohapeal)
- veergude ümbernimetamine, mis erinevad vaid suur- ja väiketähtedes (nt "Foo" ja "foo")
Lõpuks, koos järgmise Spark 3.0 väljaandmisega toetatakse täielikult selget DDL (kasutades ALTER TABLE), mis võimaldab kasutajatel teha järgmisi toiminguid tabelite skeemidega:
- veergude lisamine
- veergude kommentaaride muutmine
- tabeli omaduste seadistamine, mis määravad tabeli käitumise, näiteks tehinguajaloogile säilitamise kestuse seadmine.
Mis on skeemi evolutsiooni kasu?
Skeemi evolutsiooni saab kasutada alati, kui te kavatsete muuta oma tabeli skeemi (vastupidiselt juhtumitele, kui olete kogemata lisanud oma DataFrame'i veerge, mida seal olema ei peaks). See on kõige lihtsam viis oma skeemi migreerimiseks, kuna see lisab automaatselt õiged veergude nimed ja andmetüübid ilma, et neid tuleks selgelt deklareerida.
Kokkuvõte
Sunnitud rakendamine skeem muudab kõik uued veerud või muud skeemimuudatused, mis ei ühildu teie tabeliga, automaatselt kehtetuks. Nende kõrgete standardite seadmine ja säilitamine võimaldab analüütikutel ja inseneridel usaldada, et nende andmed on kõrgeimal tasemel terviklikkuses, arutledes selle üle selgelt ja arusaadavalt, mis võimaldab neil teha tõhusamaid äriotsuseid.
Teiselt poolt täiendab skeemi evolutsioon sunnitud rakendamist, muutes selle lihtsamaks ennustatud automaatsete skeemimuudatuste puhul. Lõppude lõpuks ei tohiks see olla keeruline — lisada veerg.
Sunnitud rakendamine on yang, samas kui skeemi evolutsioon on yin. Koos kasutades muudavad need funktsioonid müra summutamise ja signaali seadistamise lihtsamaks nagu kunagi varem.
Tahame samuti tänada Mukula Murti ja Pranav Anandi nende panuse eest selle artikli kallal.
Teised artiklid sellest sarjast:

Seotud artiklid
Allikas: habr.com
