Spark schemaEvolution praktikas

Lugupidamine, head pÀeva!

Selles artiklis kirjeldab Big Data Solutionsi Ă€risuuna juhtiv konsultant ettevĂ”ttes «Neoflex» ĂŒksikasjalikult muutuva struktuuriga vitriinide loomise vĂ”imalusi, kasutades Apache Spark'i.

AndmeanalĂŒĂŒsi projekti raames tekib sageli vajadus vitriinide loomiseks halvasti struktureeritud andmete pĂ”hjal.

Tavaliselt on need logid vĂ”i erinevate sĂŒsteemide vastused, mis salvestatakse JSON-i vĂ”i XML-vormingus. Andmed laaditakse Hadoopi, sealt tuleb seejĂ€rel vitriin luua. JuurdepÀÀsu loodud vitriinile saab korraldada nĂ€iteks lĂ€bi Impala.

Sellisel juhul ei ole sihtvitriini skeem eelnevalt teada. Veelgi enam, skeemi ei ole vĂ”imalik koostada ka ette, kuna see sĂ”ltub andmetest, millega me tegeleme — need on samad halvasti struktureeritud andmed.

NÀiteks tÀnapÀeval logitakse jÀrgmine vastus:

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

ja homme tuleb samalt sĂŒsteemilt jĂ€rgmine vastus:

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

Tulemuseks peaks vitriinile lisanduma veel ĂŒks vĂ€li — description, ja kas see tuleb vĂ”i mitte, ei tea keegi.

Andmete esitlemine on ĂŒsna tavaline ĂŒlesanne ning Sparkil on selleks mitmeid tööriistu. Algandmete parsimiseks toetatakse nii JSON-i kui ka XML-i ning ettearvamatute skeemide jaoks on olemas schemaEvolution tugi.

Esmapilgul tundub lahendus lihtne. Tuleb vÔtta JSON-iga kaust ja lugeda see dataframe'i. Spark loob skeemi, muutes sisemised andmed struktuurideks. Edasi tuleb kÔik salvestada parquet formaati, mida toetatakse ka Impalas, registreerides vitriini Hive metastore'is.

Tundub ju kÔik lihtne.

Kuid lĂŒhikestest nĂ€idetest dokumentatsioonis ei ole selge, kuidas tegeleda mitmete praktiliste probleemidega.

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

Konkreetsemalt tuuakse lihtsalt vÀlja, kuidas lugeda ja parsida JSON-i:

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

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

Praktikas on stsenaarium palju keerulisem kui lihtsalt lugeda JSON-faile kaustast ja luua dataframe. Situatsioon nÀeb vÀlja selline: olemas on kindel vitriin, iga pÀev saabub uusi andmeid, mida tuleb vitriini lisada, unustamata, et skeem vÔib erineda.

TĂŒĂŒpiline vitriini ehitamise skeem on jĂ€rgmine:

Samm 1. Andmed laaditakse Hadoopisse koos jÀrgneva igapÀevase laadimisega ning need paigutatakse uude partitsiooni. Tulemuseks on pÀevade kaupa partitsioneeritud kaust algandmetega.

Samm 2. Alglaadimise kĂ€igus loetakse ja analĂŒĂŒsitakse seda kausta Spark'i vahendite abil. Saadud dataframe salvestatakse analĂŒĂŒsimiseks sobivasse formaati, nĂ€iteks parquet, mida saab hiljem importida Impalasse. Nii luuakse eesmĂ€rgipĂ€rane vitriin kĂ”igi nende andmetega, mis selleks ajaks on kogunenud.

Samm 3. Luakse laadimise ĂŒlesanne, mis igal pĂ€eval vitriini uuendab.
KĂŒsimuseks on inkrementaalne laadimine, vitriini partitsioneerimise vajadus ja ĂŒldise vitriini skeemi toetamise kĂŒsimus.

Toome nÀite. Oletame, et esimene samm andmehoidla loomisel on ellu viidud ja JSON-failide vÀljavÔte on mÀÀratud kausta.

Nendest dataframe'i loomine, et seejÀrel sÀilitada vitriinina, ei ole probleem. See on see esimene samm, mille leidmine Spark'i dokumentatsioonist on lihtne:

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)

NÀib, et kÔik on hÀsti.

Oleme lugenud ja analĂŒĂŒsinud JSON-i, seejĂ€rel salvestame andmeraami parquet-vormingus, registreerides Hive'is ĂŒkskĂ”ik millise mugava meetodi abil:

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

Saame vitriini.

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

Loogiliseks lahenduseks oleks vitriini jagamine pÀevade kaupa, mis vÔimaldaks iga jÀrgmise pÀeva lisada uue partitsiooni. Selle mehhanism on samuti hÀsti tuntud, Spark vÔimaldab partitsioone eraldi salvestada.

Esiteks teeme initsialiseeriva laadimise, sĂ€ilitades andmed nagu eespool kirjeldatud, lisades ainult partitsioneerimise. Seda tegevust nimetatakse vitriini initsialiseerimiseks ja seda tehakse vaid ĂŒ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ÀÀnud on vaid Hive'is uuesti registreerida, et skeemi vÀrskendada.
Siin aga tekivad probleemid.

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

Vaatame tĂŒĂŒpilist olukorda. NĂ€iteks eile saabus JSON:

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

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

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

Oletame, et meil on kaks erinevat partitsiooni, milles kummaski on ĂŒks rida.
Kui loeme algandmeid koos, suudab Spark mÀÀrata tĂŒĂŒbi ning mĂ”ista, et „a” on struktuurivĂ€li, mille sees on vĂ€lisvĂ€lja „b” tĂŒĂŒp INT. Kuid kui iga partitsioon salvestati eraldi, saame parquet'iga, millel on kokkusobimatud partitsiooniskeemid:

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

See olukord on hĂ€sti tuntud, seetĂ”ttu on spetsiaalselt lisatud valik — algandmete töötlemisel eemaldada tĂŒhjad vĂ€ljad:

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

Sellisel juhul koosneb parquet partitsioonidest, mida saab ĂŒhiselt lugeda.
Kuigi need, kes seda praktikas tegid, naeratavad siin kibedalt. Miks? Sest tĂ”enĂ€oliselt tekib veel kaks situatsiooni. VĂ”i kolm. VĂ”i neli. Esimene, mis peaaegu kindlasti juhtub, on see, et numbritĂŒĂŒbid nĂ€evad erinevalt vĂ€lja erinevates JSON-failides. NĂ€iteks, {intField: 1} ja {intField: 1.1}. Kui sellised vĂ€ljad satuvad samasse partitsiooni, loeb skeemi fusioon kĂ”ik Ă”igesti, muutes selle kĂ”ige tĂ€psemaks tĂŒĂŒbiks. Ent kui need on erinevates partitsioonides, on ĂŒhes 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 on partitsioonid, mida saab lugeda ĂŒhte andmeraamistikku ja valideeritud parquet kogu vitriini jaoks. Jah? Ei.

Peame meeles pidama, et registreerisime tabeli Hive'is. Hive ei ole tundlik vÀljade nimede suuruse osas, kuid parquet on tundlik. SeetÔttu on partitsioonid skeemidega: field1: int ja Field1: int Hive'i jaoks samad, kuid Spark'i jaoks mitte. Peame meeles pidama, et vÀljade nimed tuleks muuta vÀiketÀhtedeks.

PÀrast seda tundub kÔik hÀsti.

Kuid mitte kÔik pole nii lihtne. TÔuseb teine, samuti hÀsti tuntud probleem. Kuna iga uus partitsioon salvestatakse eraldi, siis on partitsiooni kaustas Spark'i teenindusfailid, nÀiteks operatsiooni edukuse lipp _SUCCESS. See toob kaasa vea, kui proovida parquet'i. Selle vÀltimiseks tuleb seadistada konfigureerimine, keelates Spark'il teenindusfailide kirjutamise kausta:

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

Paistab, et nĂŒĂŒd lisandub sihtvitriini kausta igapĂ€evaselt uus parquet partitsioon, kus on pĂ€evased parsitud andmed. Oleme eelnevalt mures, et andmetĂŒĂŒpide konflikte ei tekiks.

Aga meie ees on kolmas probleem. NĂŒĂŒd ei ole ĂŒldine skeem teada, lisaks on Hive'is vale skeem, kuna iga uus partitsioon on tĂ”enĂ€oliselt skeemi moonutanud.

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

Meie ees seisab neljas probleem. Kui me tabelit esmakordselt registreerisime, toetusime Sparkile. NĂŒĂŒd teeme seda ise, ja peame meeles pidama, et parquet-failid vĂ”ivad alata mĂ€rkidega, mis Hive'is pole lubatud. NĂ€iteks eemaldab Spark read, mida ta ei suutnud "corrupt_record" vĂ€lja lugeda. Sellist vĂ€lja ei saa Hive'is registreerida ilma, et see peaks olema ekraanituks.

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 ohutu DDL, mis tÀhendab, et selle asemel:

create table tname (_field1 string, 1field string)

Selliste vÀljade nagu «_field1, 1field» korral luuakse ohutu DDL, kus vÀlja nimed on ekraanitud: create table `tname` (`_field1` string, `1field` string).

KĂŒsimus on, kuidas saada andmeraam, millel on tĂ€ielik skeem (pf koodis)? Kuidas saada see pf? See on viies probleem. Kas lugeda kĂ”igi partitsioonide skeemi sihivitrini parquet faili kaustast? See meetod on kĂ”ige turvalisem, kuid töömahukam.

Skeem on juba Hive'is. Uue skeemi saab hankida, kombineerides terve tabeli skeemi uue partitsiooniga. See tĂ€hendab, et tabeli skeem tuleb vĂ”tta Hive'ist ja ĂŒhendada see uue partitsiooni skeemiga. Seda saab teha, lugedes testmetainformatsiooni Hive'ist, salvestades need ajutisse kausta ja lugedes Sparkiga mĂ”lemad partitsioonid korraga.

PĂ”himĂ”tteliselt on kĂ”ik vajalik olemas: algne tabeli skeem Hive'is ja uus partitsioon. Andmed on meil samuti olemas. JÀÀb ĂŒle vaid saada uus skeem, milles kombineeritakse vitrini skeem 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 registreerimise, nagu eelnevas fragmendis.
Kui kogu ahel töötab Ôigesti, st on toimunud alglaadimine ja Hive'is on Ôigesti loodud tabel, siis saame vÀrskendatud tabeli skeemi.

Ja viimane probleem on see, et ei saa lihtsalt lisada partitsiooni Hive'i tabelisse, kuna see puruneb. On vajalik panna Hive 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 pĂ”hjal vitriin osutub mitmete varjatud raskuste ĂŒletamiseks, mille lahendusi tuleb otsida eraldi. Ja kuigi need lahendused on lihtsad, vĂ”tab nende leidmine palju aega.

Vitriini koostamiseks tuli:

  • Lisa vitriinile partitsioone, vabanedes abifailidest.
  • Selgitada vĂ€lja tĂŒhjad vĂ€ljad algandmetes, mille Spark tĂŒĂŒbis.
  • Viia lihtsad tĂŒĂŒbid stringideks.
  • Viia vĂ€ljade nimed vĂ€ikestesse tĂ€htedesse.
  • Eraldada andmete vĂ€ljaviimine ja tabeli registreerimine Hive'is (DDL loomine).
  • Ära unusta escape'ida vĂ€ljade nimed, mis vĂ”ivad Hive'iga ĂŒhilduvad olla.
  • Õppida, kuidas uuendada tabeli registreerimist Hive'is.

KokkuvÔttes tasub mÀrkida, et vitriinide ehitamise lahendusel on palju varjatud ohte. SeetÔttu on keeruliste olukordade tekkimisel parem pöörduda kogenud partneri poole, kellel on edukas kogemus.

AitÀh, et lugesite seda artiklit, loodame, et teave osutub kasulikuks.

Allikas: habr.com

Osta usaldusvÀÀrne veebihosting DDoS kaitsega, VPS VDS serverid đŸ”„ Osta usaldusvÀÀrne veebihosting DDoS kaitsega, VPS VDS serverid | ProHoster