Spark schemaEvolution praktikas

Lugupidajad, tere pÀevast!

KĂ€esolevas artiklis kĂ€sitleb Big Data Solutions Ă€riĂŒksuse juhtiv nĂ”ustaja ettevĂ”ttest "Neoflex" ĂŒksikasjalikult vĂ”imalusi muuta muutuva struktuuriga vitriine Apache Sparki abil.

AndmeanalĂŒĂŒsiprojekti raames tekib sageli vajadus luua vitriine halvasti struktureeritud andmete pĂ”hjal.

Need on tavaliselt logid vĂ”i erinevate sĂŒsteemide vastused, mis salvestatakse JSON-i vĂ”i XML-i kujul. Andmed laaditakse Hadoopisse ja seejĂ€rel tuleb nendest vitriin luua. JuurdepÀÀsu loodud vitriinile saame korraldada nĂ€iteks lĂ€bi Impala.

Sel juhul ei ole sihtvitriini skeem eelnevalt teada. Veelgi enam, skeemi ei saa ka eelnevalt koostada, kuna see sÔltub andmetest ja meil on tegemist nende halvasti struktureeritud andmetega.

NÀiteks tÀnasel pÀeval logitakse selline vastus:

{source: "app1", error_code: ""}

ja homme tuleb samalt sĂŒsteemilt selline vastus:

{source: "app1", error_code: "error", description: "Network error"}

SeetĂ”ttu peaks vitriini lisanduma veel ĂŒks vĂ€li — description, ja kas see tuleku vĂ”i mitte, ei tea keegi.

Vitriini loomise ĂŒlesanne selliste andmete pĂ”hjal on ĂŒsna tavaline, ja Sparkil on selleks olemas mitmeid tööriistu. Algandmete analĂŒĂŒsimiseks on toetatud nii JSON kui ka XML, ning eelnevalt tundmatu skeemi puhul on olemas ka schemaEvolution toe.

Esmapilgul nÀeb lahendus lihtne vÀlja. Tuleb vÔtta kaust JSON-ide ja lugeda see dataframe'i. Spark loob skeemi, sisemised andmed muudab struktuuri. Edasi tuleb kÔik salvestada parquet'isse, mis on samuti toetatud Impalas, registreerides vitriini Hive metastore'is.

Tundub, et kÔik on lihtne.

Kuid lĂŒhikestest nĂ€idistest dokumendis ei ole selge, mida praktikas mitmete probleemidega ette vĂ”tta.

Dokumentatsioonis kirjeldatakse lÀhenemist mitte vitriini loomiseks, vaid JSON-i vÔi XML-i lugemiseks dataframe'i.

Konkreetselt, kirjeldatakse lihtsalt, kuidas lugeda ja analĂŒĂŒsida JSON-i:

df = spark.read.json(path...)

Seda on piisavalt, et andmed oleksid Sparki jaoks kergesti kÀttesaadavad.

Kuid praktikas on stsenaarium palju keerulisem kui lihtsalt lugeda JSON-failid kaustast ja luua dataframe. Olukord on selline: juba on olemas kindel vitriin, igapÀev tulevad uued andmed, need tuleb vitriini lisada, meeles pidades, et skeem vÔib varieeruda.

Tavaline vitriini ĂŒlesehitamise skeem on jĂ€rgmine:

Samm 1. Andmed laaditakse Hadoopisse, millele jÀrgneb igapÀevane tÀiendamine ja need koondatakse uude partitsiooni. Saame pÀevade kaupa partitsioneeritud kausta, kus on algandmed.

Samm 2. Initsialiseerimise kĂ€igus loetakse ja analĂŒĂŒsitakse seda kausta Sparkiga. Saadud dataframe salvestatakse analĂŒĂŒsiks sobivasse formaati, nĂ€iteks parquet, mille saab hiljem Impalasse importida. Nii luuakse sihtkapp, kus on kĂ”ik senikogutud andmed.

Samm 3. Luakse laadimine, mis igapÀevaselt uuendab vÔti.
KĂŒsimus on inkrementaalses laadimises, vitriini partitsioneerimises ja ĂŒldise vitriini skeemi toetamise vajaduses.

Toome nÀite. Oletame, et hoidlahoone esimene etapp on teostatud ja JSON-failide eksportimine kausta on seadistatud.

Nendest dataframe'i loomine, et hiljem salvestada see vitriinina, ei valmista probleeme. See on just see esimene samm, mida on lihtne leida Spark'i dokumentatsioonist:

df = spark.read.option("mergeSchema", True).json(".../*") 
df.printSchema()

root 
|-- a: long (nullable = true) 
|-- b: string (nullable = true) 
|-- c: struct (nullable = true) |    
|-- d: long (nullable = true)

Tundub, et kÔik on korras.

Me oleme JSON'i lugenud ja analĂŒĂŒsinud ning nĂŒĂŒd salvestame dataframe'i parquet'ina, registreerides Hive'is igal sobival viisil:

df.write.format(“parquet”).option('path','').saveAsTable('')

Saame vitriini.

Kuid jĂ€rgmisel pĂ€eval lisandusid andmed allikast. Meil on kaust JSON-idega ja vitriin, mis on loodud selle kausta pĂ”hjal. PĂ€rast jĂ€rgmise andmepartii laadimist allikast puuduvad vitriinis ĂŒhe pĂ€eva andmed.

Loogiliseks lahenduseks on partitsioneerida vitriin pÀevade kaupa, mis vÔimaldab igapÀevaselt lisada uue partitsiooni. Selle mehhanism on samuti hÀsti tuntud, Spark vÔimaldab partitsioone eraldi salvestada.

Esmalt teeme initsialiseerimise laadimise, salvestades andmed nagu eespool kirjeldatud, lisades ainult partitsioneerimise. See toiming on vitriini initsialiseerimine ja see tehakse ainult ĂŒks kord:

df.write.partitionBy("date_load").mode("overwrite").parquet(dbpath + "/" + db + "/" + destTable)

JÀrgmisel pÀeval laadime ainult uue partitsiooni:

df.coalesce(1).write.mode("overwrite").parquet(dbpath + "/" + db + "/" + destTable + "/date_load=" + date_load + "/")

JÀÀb ainult Hive'is uuesti registreerida, et skeemi vÀrskendada.
Kuid siin tekivad probleemid.

Esimene probleem. A varem vĂ”i hiljem ei saa loodud parquet'i enam lugeda. See on seotud sellega, kuidas parquet ja JSON kĂ€sitlevad tĂŒhje vĂ€lju erinevalt.

Vaadakem tĂŒĂŒpilist olukorda. NĂ€iteks, eile tuleb JSON:

PĂ€ev 1: {"a": {"b": 1}},

aga tÀna nÀeb see sama JSON vÀlja nii:

PĂ€ev 2: {"a": null}

Oletame, et meil on kaks erinevat partiid, kus kummaski on ĂŒks rida.
Kui me loeme algandmed tĂ€ielikult, suudab Spark tuvastada tĂŒĂŒbi ja mĂ”ista, et "a" on struktuuri tĂŒĂŒpi vĂ€li, mille sees on INT tĂŒĂŒpi vĂ€li "b". Kuid kui iga partii on salvestatud eraldi, tekib parquet, millel on mitteĂŒhtuvad partii skeemid:

df1 (a: <struct>)
df2 (a: STRING NULLABLE)

See olukord on hĂ€sti teada, seetĂ”ttu on spetsiaalselt lisatud valik - algandmete parsimisel tĂŒhjad vĂ€ljad eemaldada:

df = spark.read.json("...", dropFieldIfAllNull=True)

Sel juhul koosneb parquet partiiidest, mida saab koos lugeda.
Kuigi need, kes on seda praktikas teinud, naeratavad sellele kibedalt. Miks? Sest ilmselt tekib veel kaks olukorda. VĂ”i kolm. VĂ”i neli. Esimene, mis peaaegu kindlasti tekib, on see, et numbrilised tĂŒĂŒbid nĂ€evad erinevates JSON failides erinevad vĂ€lja. NĂ€iteks, {intField: 1} ja {intField: 1.1}. Kui sellised vĂ€ljad satuvad ĂŒhte partiisse, loeb skeemi mergenĂ”uet kĂ”ik Ă”igesti, tuues vĂ€lja kĂ”ige tĂ€psema tĂŒĂŒbi. Kuid kui need on erinevates partii, siis ĂŒhes on intField: int ja teises intField: double.

Selle olukorra töötlemiseks on jÀrgmine lipp:

df = spark.read.json("...", dropFieldIfAllNull=True, primitivesAsString=True)

NĂŒĂŒd on meil kaust, kus asuvad partiiid, mida saab lugeda ĂŒheks dataframe'iks ja kehtivaks parquet'iks kogu vitriini jaoks. Jah? Ei.

Peame meenutama, et registreerisime tabeli Hive'is. Hive ei ole tundlik suuruse suhtes vÀljade nimedes, kuid parquet on tundlik. SeetÔttu on partii skeemid: field1: int ja Field1: int Hive jaoks samad, kuid Sparkile mitte. Tuleb meeles pidada, et muudame vÀljade nimed vÀikeste tÀhtedega.

PÀrast seda tundub, et kÔik on korras.

Kuid asi ei ole nii lihtne. Tekib teine, samuti hÀsti tuntud probleem. Kuna iga uus partii salvestatakse eraldi, on partii kaustas Spark'i teenindusfailid, nÀiteks operatsiooni edukuse lipp _SUCCESS. See pÔhjustab viga, kui proovida parquet'it. Selle vÀltimiseks tuleb seadistada konfiguratsioon, keelates Spark'il teenindusfailide kirjutamise kausta:

hadoopConf = sc._jsc.hadoopConfiguration()
hadoopConf.set("parquet.enable.summary-metadata", "false")
hadoopConf.set("mapreduce.fileoutputcommitter.marksuccessfuljobs", "false")

Tundub, nĂŒĂŒd lisatakse igapĂ€evaselt sihtkausta uusi parquet partiisid, kus on pĂ€eva jooksul analĂŒĂŒsitud andmed. Oleme eelnevalt hoolitsenud selle eest, et andmetĂŒĂŒpide konfliktidega partiisid ei esineks.

Kuid nĂŒĂŒd seisame silmitsi kolmanda probleemiga. NĂŒĂŒd ei ole ĂŒldine skeem teada, veel enam, Hive'is on vale skeem, kuna iga uus partiia on tĂ”enĂ€oliselt skeemi moonutanud.

Tabel tuleb uuesti registreerida. Selle saab teha lihtsalt: loeme uuesti parquet vitriine, vÔtame skeemi ja loome selle pÔhjal DDL, millega registreerime kausta Hive'is uuena vÀlise tabelina, vÀrskendades sihtvitriini skeemi.

NĂŒĂŒd tekib neljas probleem. Kui registreerisime tabeli esmakordselt, tuginesime Sparkile. NĂŒĂŒd teeme seda ise ja tuleb meeles pidada, et parquet'i vĂ€ljad vĂ”ivad alata mĂ€rkidest, mis ei ole Hive'is lubatud. NĂ€iteks viskab Spark vĂ€lja read, mida ei suudeta "corrupt_record" vĂ€lja analĂŒĂŒsida. Sellist vĂ€lja ei saa Hive'is registreerida ilma, et see oleks escapesitud.

Teades seda, saame skeemi:

f_def = ""
for f in pf.dtypes:
  if f[0] != "date_load":
    f_def = f_def + "," + f[0].replace("_corrupt_record", "`_corrupt_record`") + " " + f[1].replace(":", "`:").replace("<", "<`").replace(",", ",`").replace("array<`", "array<") 
table_define = "CREATE EXTERNAL TABLE jsonevolvtable (" + f_def[1:] + " ) "
table_define = table_define + "PARTITIONED BY (date_load string) STORED AS PARQUET LOCATION '\/user\/admin\/testJson\/testSchemaEvolution\/pq\/'"
hc.sql("drop table if exists jsonevolvtable")
hc.sql(table_define)

Kood ("_corrupt_record", "`_corrupt_record`") + " " + f[1].replace(":", "`:").replace("<", "<`").replace(",", ",`").replace("array<`", "array<") teeb turvalise DDL, st selle asemel:

create table tname (_field1 string, 1field string)

Selliste vĂ€ljanimega nagu „_field1, 1field“ tehakse turvaline DDL, kus vĂ€lja nimed on escapesitud: create table `tname` (`_field1` string, `1field` string).

KĂŒsimus on: kuidas Ă”igesti saada dataframe tĂ€ieliku skeemiga (koodis pf)? Kuidas saada see pf? See on viies probleem. Kas lugeda skeemi kĂ”ikidest partiidest sihtkaustas olevatest parquet failidest? See meetod on kĂ”ige turvalisem, kuid keeruline.

Schema on juba Hive'is. Uue skeemi saamiseks tuleb ĂŒhendada kogu tabeli skeem ja uus partitsioon. See tĂ€hendab, et tuleb vĂ”tta tabeli skeem Hive'ist ja ĂŒhendada see uue partitsiooni skeemiga. Seda saab teha, lugedes testmetainfot Hive'ist, salvestades selle ajutisse kausta ja lugedes Spark'i abil mĂ”lemat partitsiooni korraga.

Sisuliselt on kĂ”ik vajalik olemas: algne tabeli skeem Hive'is ja uus partitsioon. Andmed on meil samuti olemas. JÀÀb vaid saada uus skeem, milles ĂŒhendatakse vitriiniskeem ja uued vĂ€ljad loodud partitsioonist:

from pyspark.sql import HiveContext
from pyspark.sql.functions import lit
hc = HiveContext(spark)
df = spark.read.json("...", dropFieldIfAllNull=True)
df.write.mode("overwrite").parquet(".../date_load=12-12-2019")
pe = hc.sql("select * from jsonevolvtable limit 1")
pe.write.mode("overwrite").parquet(".../fakePartiton/")
pf = spark.read.option("mergeSchema", True).parquet(".../date_load=12-12-2019/*", ".../fakePartiton/*")

SeejÀrel loome DDL tabeli registreerimiseks, nagu eelmises fragmentis.
Kui kogu ahel töötab Ă”igesti, nimelt — oli algne laadimine ja Hive'is on Ă”igesti loodud tabel, siis saame uuendatud tabeli skeemi.

Ja viimane probleem seisneb selles, et ei saa lihtsalt lisada partitsiooni Hive'i tabelisse, kuna see rikki lĂ€heb. Oluline on sundida Hive’i parandama oma partitsioonide struktuuri:

from pyspark.sql import HiveContext
hc = HiveContext(spark) 
hc.sql("MSCK REPAIR TABLE " + db + "." + destTable)

Lihtne ĂŒlesanne lugeda JSON ja luua selle alusel vitriin toob esile mitmeid varjatud raskusi, mille lahendused tuleb eraldi vĂ€lja otsida. Ja kuigi need lahendused on lihtsad, vĂ”tab nende leidmine palju aega.

Vitriini loomise elluviimiseks tuli:

  • Lisada partitsioonid vitriini, vabastudes abifailidest
  • Selgitada vĂ€lja tĂŒhjad vĂ€ljad algandmetes, mille Spark tĂŒĂŒpides
  • Muuta lihtsad tĂŒĂŒbid stringiks
  • Viia vĂ€ljade nimed vĂ€ikeste tĂ€htedega
  • Eraldada andmete eksport ja tabeli registreerimine Hive'is (DDL loomine)
  • Ärge unustage ekraanida vĂ€ljade nimesid, mis vĂ”ivad olla Hive'iga mitteĂŒhilduvad
  • Õppida, kuidas tabeli registreerimist Hive'is uuendada

KokkuvÔtteks tuleks mÀrkida, et vitriinide loomise lahendusel on palju varjatud takistusi. SeetÔttu on keerukate rakenduste puhul parem pöörduda kogenud partneri poole, kellel on edukas ekspertiis.

TĂ€nan teid selle artikli lugemise eest, loodame, et teave osutub kasulikuks.

Allikas: habr.com

Osta usaldusvÀÀrne hostimine veebilehtede jaoks DDoS-i kaitsega, VPS VDS serverid đŸ”„ Osta usaldusvÀÀrne hostimine veebilehtede jaoks DDoS-i kaitsega, VPS VDS serverid | ProHoster