Spark schemaEvolution în practică

Stimați cititori, o zi bună!

În acest articol, consultantul principal al direcției de business Big Data Solutions de la compania „Neoflex” descrie în detaliu opțiunile de construire a vitrinelor cu structură variabilă, utilizând Apache Spark.

În cadrul unui proiect de analiză a datelor, apare adesea sarcina de a construi vitrine pe baza unor date slab structurate.

De obicei, acestea sunt loguri sau răspunsuri ale diferitelor sisteme, salvate sub formă de JSON sau XML. Datele sunt exportate în Hadoop, iar din acestea trebuie să construim o vitrină. Accesul la vitrina creată poate fi organizat, de exemplu, prin intermediul Impala.

În acest caz, schema vitrinei țintă nu este cunoscută dinainte. Mai mult, schema nu poate fi alcătuită anticipat, deoarece depinde de date, iar noi avem de-a face tocmai cu aceste date slab structurate.

De exemplu, astăzi este logat următorul răspuns:

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

iar mâine această aceeași sistem va returna un răspuns de genul:

{source: "app1", error_code: "error", description: "Eroare de rețea"}

Ca rezultat, în vitrină ar trebui să fie adăugat un nou câmp — description, iar dacă va veni sau nu, nimeni nu știe.

Sarcina de a crea o vitrină pe astfel de date este destul de standard, iar Spark are la dispoziție o serie de instrumente pentru aceasta. Există suport pentru parsarea datelor originale, atât JSON, cât și XML, iar pentru schema necunoscută dinainte, este prevăzut suport pentru schemaEvolution.

La prima vedere, soluția pare simplă. Trebuie să luăm un folder cu JSON și să-l citim într-un dataframe. Spark va crea schema, datele imbricate vor fi convertite în structuri. Apoi, totul trebuie salvat în parquet, care este de asemenea acceptat de Impala, înregistrând vitrina în Hive metastore.

Se pare că totul este simplu.

Cu toate acestea, din exemplele scurte din documentație nu este clar ce trebuie făcut cu o serie de probleme în practică.

Documentația descrie o abordare nu pentru crearea unei vitrine, ci pentru citirea JSON-ului sau XML-ului într-un dataframe.

Astfel, este prezentat cum să citim și să parsăm JSON:

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

Acest lucru este suficient pentru a face datele disponibile pentru Spark.

În practică, scenariul este mult mai complex decât simpla citire a fișierelor JSON dintr-un folder și crearea unui dataframe. Situația arată astfel: deja există o vitrină determinată, în fiecare zi sosesc date noi, care trebuie adăugate în vitrină, fără a uita că schema poate diferi.

Schema obișnuită de construire a vitrinei este următoarea:

Pasul 1. Datele sunt încărcate în Hadoop cu încărcări zilnice ulterioare și sunt stocate într-o nouă partiție. Se obține astfel un folder partitionat pe zile cu datele originale.

Pasul 2. În timpul încărcării inițiale, acest folder este citit și analizat prin intermediul Spark. DataFrame-ul obținut este salvat într-un format accesibil pentru analiză, de exemplu, în parquet, care poate fi ulterior importat în Impala. Astfel se creează un țăruș țintă cu toate datele acumulate până în acel moment.

Pasul 3. Se creează o încărcare care va actualiza în fiecare zi țărușul.
Apare întrebarea despre încărcarea incrementală, necesitatea partitionării țărușului și întrebarea despre suportarea schemei generale a țărușului.

Să dăm un exemplu. Să presupunem că primul pas de construcție a depozitului este implementat și exportul fișierelor JSON este configurat pentru folder.

Crearea unui DataFrame din acestea pentru a-l salva ulterior ca țăruș nu este o problemă. Acesta este primul pas pe care îl poți găsi ușor în documentația Spark:

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)

Până acum, totul este bine.

Am citit și am analizat JSON-ul, apoi salvăm DataFrame-ul ca parquet, înregistrându-l în Hive în orice mod convenabil:

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

Obținem țărușul.

Dar, în ziua următoare, au fost adăugate date noi din sursă. Avem un folder cu JSON, și un țăruș creat pe baza acestui folder. După încărcarea următoarei porții de date din sursă, în țăruș lipsesc date pentru o zi.

O soluție logică ar fi să partitionăm țărușul pe zile, ceea ce ar permite adăugarea unei noi partiții în fiecare zi. Mecanismul acestuia este, de asemenea, bine cunoscut, Spark permite scrierea partițiilor separat.

Mai întâi facem o încărcare inițială, salvând datele așa cum a fost descris mai sus, adăugând doar partitionarea. Această acțiune se numește inițializare a țărușului și se face doar o dată:

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

În ziua următoare, încărcăm doar noua partiție:

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

Rămâne doar să reînregistrăm în Hive pentru a actualiza schema.
Cu toate acestea, aici apar probleme.

Prima problemă. Mai devreme sau mai târziu, fișierul parquet rezultat nu va putea fi citit. Acest lucru se datorează modului diferit în care parquet și JSON abordează câmpurile goale.

Să luăm o situație tipică. De exemplu, ieri a sosit JSON-ul:

Ziua 1: {"a": {"b": 1}},

iar astăzi același JSON arată astfel:

Ziua 2: {"a": null}

Să presupunem că avem două partiții diferite, fiecare având o singură linie.
Atunci când citim datele originale în întregime, Spark va reuși să determine tipul și va înțelege că „a” este un câmp de tip „structură”, cu un câmp încorporat „b” de tip INT. Dar, dacă fiecare partiție a fost salvată separat, se obține un parquet cu scheme de partiții incompatibile:

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

Această situație este binecunoscută, motiv pentru care a fost adăugată o opțiune - în timpul analizei datelor originale, să se elimine câmpurile goale:

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

În acest caz, parquet-ul va consta din partiții care pot fi citite împreună.
Deși cei care au lucrat cu asta în practică vor zâmbi amar. De ce? Pentru că, cel mai probabil, vor apărea încă două situații. Sau trei. Sau patru. Prima, care va apărea aproape cu siguranță, este că tipurile numerice vor arăta diferit în diferite fișiere JSON. De exemplu, {intField: 1} și {intField: 1.1}. Dacă astfel de câmpuri ajung în aceeași partiție, atunci fuzionarea schemelor va citi totul corect, aducându-le la cel mai precis tip. Dar, dacă sunt în diferite, atunci într-una va fi intField: int, iar în cealaltă intField: double.

Pentru a gestiona această situație, există următoarea opțiune:

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

Acum avem un folder unde se află partițiile, care pot fi citite într-un singur dataframe și parquet valid pentru întreaga vitrină. Da? Nu.

Trebuie să ne amintim că am înregistrat tabela în Hive. Hive nu este sensibil la majuscule în numele câmpurilor, pe când parquet este. Prin urmare, partițiile cu schemele: field1: int, și Field1: int pentru Hive sunt identice, dar pentru Spark nu. Trebuie să nu uităm să convertim numele câmpurilor în litere mici.

După aceasta, pare că totul este în regulă.

Cu toate acestea, nu este atât de simplu. Apare a doua problemă, de asemenea binecunoscută. Deoarece fiecare nouă partiție este salvată separat, în folderul partiției vor exista fișiere de sistem Spark, de exemplu, un indicator al succesului operației _SUCCESS. Aceasta va duce la o eroare la încercarea de a citi parquet. Pentru a evita acest lucru, trebuie să configurăm setările, interzicând Spark să scrie fișiere de sistem în folder:

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

Se pare că, în fiecare zi, se adaugă o nouă partiție parquet în folderul vitrinei țintă, unde se află datele parsate pentru zi. Ne-am ocupat din timp de faptul că nu vor exista partiții cu conflicte de tipuri de date.

Dar, înainte de noi — a treia problemă. Acum schema generală nu este cunoscută, mai mult, în tabelul Hive există o schemă incorectă, deoarece fiecare nouă partiție a adus, cel mai probabil, o distorsionare a schemei.

Trebuie să înregistrăm din nou tabela. Acest lucru se poate face simplu: reluând citirea vitrinei parquet, luând schema și creând pe baza ei DDL, cu care să înregistrăm din nou folderul în Hive ca tabel extern, actualizând schema vitrinei țintă.

În fața noastră apare a patra problemă. Atunci când am înregistrat prima dată tabela, ne-am bazat pe Spark. Acum facem acest lucru noi înșine și trebuie să ne amintim că câmpurile parquet pot începe cu simboluri care nu sunt acceptate în Hive. De exemplu, Spark elimină șirurile pe care nu a reușit să le parseze în câmpul „corrupt_record”. Un astfel de câmp nu va putea fi înregistrat în Hive fără a fi scăpat.

Având în vedere acest lucru, obținem schema:

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)

Cod ("_corrupt_record", "`_corrupt_record`") + " " + f[1].replace(":", "`:").replace("<", "<`").replace(",", ",`").replace("array<`", "array<") face un DDL sigur, adică în loc de:

create table tname (_field1 string, 1field string)

Cu nume de câmpuri cum ar fi „_field1, 1field”, se face un DDL sigur, unde numele câmpurilor sunt escapate: create table `tname` (`_field1` string, `1field` string).

Se ridică întrebarea: cum obținem corect un dataframe cu schema completă (în codul pf)? Cum obținem acest pf? Aceasta este a cincea problemă. Să recitim schema tuturor partițiilor din folderul cu fișiere parquet ale vitrinei țintă? Aceasta este cea mai sigură metodă, dar grea.

Schema există deja în Hive. O nouă schemă poate fi obținută prin combinarea schemei întregii tabele cu noua partiție. Așadar, trebuie să preluăm schema tabelului din Hive și să o combinăm cu schema noii partiții. Acest lucru se poate realiza citind metadatele de test din Hive, salvându-le într-un folder temporar și citind ambele partiții simultan cu Spark.

Practic, avem tot ce ne trebuie: schema inițială a tabelului din Hive și noua partiție. De asemenea, avem datele. Rămâne doar să obținem noua schemă, în care se combină schema vitrinei și noile câmpuri din partiția creată:

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/*")

Apoi creăm DDL-ul pentru înregistrarea tabelului, așa cum am făcut în fragmentul anterior.
Dacă întreaga lanț funcționează corect — adică a fost o încărcare inițială, iar tabela a fost creată corect în Hive, atunci obținem schema actualizată a tabelului.

Și, în cele din urmă, problema este că nu se poate adăuga o partiție în tabela Hive atât de ușor, deoarece aceasta va fi ruptă. Este necesar să forțăm Hive să repare structura partițiilor sale:

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

O sarcină simplă de citire a JSON-ului și de creare a unei vitrine pe baza acestuia se transformă în depășirea unei serii de dificultăți implicite, soluții pentru care trebuie căutate separat. Și, deși aceste soluții sunt simple, căutarea lor consumă mult timp.

Pentru a implementa construirea vitrinei, a fost necesar să:

  • Adăugăm partiții în vitrină, scăpând de fișierele de sistem
  • Să clarificăm câmpurile goale din datele originale, care au fost tipizate de Spark
  • Să convertim tipurile simple în stringuri
  • Să aducem denumirile câmpurilor la litere mici
  • Să separăm exportul de date de înregistrarea tabelului în Hive (crearea DDL)
  • Să nu uităm să escamotez denumirile câmpurilor care ar putea fi incompatibile cu Hive
  • Să învățăm să actualizăm înregistrarea tabelului în Hive

În concluzie, să subliniem că soluția pentru construirea vitrinelor ascunde multe capcane. Prin urmare, în cazul în care apar dificultăți în implementare, este mai bine să ne adresăm unui partener experimentat cu expertiză de succes.

Vă mulțumim pentru citirea acestui articol, sperăm că informația va fi utilă.

Sursa: habr.com

Cumpără un hosting fiabil pentru site-uri cu protecție DDoS, servere VPS VDS 🔥 Cumpără un hosting fiabil pentru site-uri cu protecție DDoS, servere VPS VDS | ProHoster