Të nderuar lexues, ditë të mirë!
Në këtë artikull, konsultanti kryesor i drejtimit të biznesit Big Data Solutions të kompanisë "Neoflex" përshkruan me detaje opcionet e ndërtimit të vitrinave me strukturë të ndryshueshme duke përdorur Apache Spark.
Në kuadër të një projekti për analizën e të dhënave, shpesh paraqitet detyra e ndërtimit të vitrinave mbi baza të dhënash pak të strukturuara.
Zakonisht, këto janë loge, ose përgjigje të ndryshme nga sisteme, të ruajtura në formën e JSON ose XML. Të dhënat eksportohen në Hadoop, pastaj nga ato duhet të ndërtohet një vitrinë. Organizimi i aksesit në vitrinë mund të bëhet, për shembull, përmes Impala.
Në këtë rast, skema e vitrinës së synuar nuk dihet paraprakisht. Për më tepër, skema nuk mund të përgatitet paraprakisht, sepse varet nga të dhënat, dhe ne merremi me këto të dhëna pak të strukturuara.
Për shembull, sot regjistrohet një përgjigje e tillë:
{source: "app1", error_code: ""}ndërsa nesër nga i njëjti sistem vjen një përgjigje e tillë:
{source: "app1", error_code: "error", description: "Gabim nĂ« rrjet"}Si pasojĂ«, nĂ« vitrinĂ« duhet tĂ« shtohet njĂ« fushĂ« tjetĂ«r â description, dhe se do tĂ« vijĂ« apo jo, askush nuk e di.
Detyra e krijimit të vitrinave mbi këto të dhëna është mjaft standarde, dhe Spark ka një sërë mjetesh për këtë. Për parsing të të dhënave burimore, ka mbështetje për JSON dhe XML, dhe për skema të panjohura paraprakisht, është parashikuar mbështetje për schemaEvolution.
Nga pamja e parë, zgjidhja duket e thjeshtë. Duhet të merret një dosje me JSON dhe të lexohet në dataframe. Spark do të krijojë skemën, të dhënat e përfshira do të shndërrohen në struktura. Më pas, gjithçka duhet të ruhet në parquet, që përkrahet gjithashtu në Impala, duke regjistruar vitrinën në Hive metastore.
Duket se gjithçka është e thjeshtë.
Megjithatë, nga shembujt e shkurtër në dokumentacion, nuk është e qartë se çfarë të bëhet me disa probleme në praktikë.
Në dokumentacion përshkruhet qasja jo për krijimin e vitrinës, por për leximin e JSON-it ose XML-it në dataframe.
I veçantë, thjesht jepen hapat për të lexuar dhe parser JSON-in:
df = spark.read.json(path...)Kjo është e mjaftueshme për ta bërë të dhënat të disponueshme për Spark.
NĂ« praktikĂ«, skenari Ă«shtĂ« shumĂ« mĂ« i komplikuar se thjesht tĂ« lexosh skedarĂ«t JSON nga dosja dhe tĂ« krijosh dataframe. Situata duket kĂ«shtu: tashmĂ« ekziston njĂ« vitrinĂ« e caktuar, çdo ditĂ« vijnĂ« tĂ« dhĂ«na tĂ« reja, duhet tâi shtohet vitrinĂ«s, pa harruar se skema mund tĂ« ndryshojĂ«.
Skema e zakonshme e ndërtimit të vitrinës është kështu:
Hapi 1. Të dhënat ngarkohen në Hadoop me një përditësim të përditshëm dhe ruhet në një particion të ri. Resultati është një dosje e ndarë në ditë me të dhënat e origjinës.
Hapi 2. Në procesin e ngarkimit inicial, kjo dosje lexon dhe analizohet me mjetet Spark. Dataframe që rezulton ruhet në një format të disponueshëm për analizë, për shembull, në formatin parquet, i cili pastaj mund të importohet në Impala. Kjo krijon një ndriçim të synuar me të gjitha të dhënat që janë akumuluar deri në atë moment.
Hapi 3. Krijohet një ngarkesë që do të përditësojë ndriçimin çdo ditë.
Shtrohet pyetja e ngarkimit inkremental, nevoja për ndarjen e ndriçimit dhe çështja e mbështetjes së skemës së përgjithshme të ndriçimit.
Le të japim një shembull. Le të themi se hapi i parë i ndërtimit të depozitës është implementuar dhe janë të vendosura skedarët JSON në dosje.
Krijimi i një dataframe prej tyre për ta ruajtur si një ndriçim, nuk paraqet probleme. Ky është hapi i parë që lehtë mund të gjendet në dokumentacionin e 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)Duket se gjithçka është në rregull.
Ne lexuam dhe analizuam JSON, tani ruajmë dataframe-n si parquet, duke regjistruar në Hive në çdo mënyrë të përshtatshme:
df.write.format(âparquetâ).option('path','').saveAsTable('') Kemi krijuar ndriçimin.
Por, ditën tjetër u shtuan të dhëna të reja nga burimi. Kemi një dosje me JSON dhe një ndriçim të krijuar mbi këtë dosje. Pas ngarkimit të dozës tjetër të të dhënave nga burimi, ndriçimi nuk ka të dhëna për një ditë.
Një zgjidhje logjike do të ishte ndarësimi i ndriçimit sipas ditëve, që do të lejojë shtimin e një partizioni të re çdo ditë. Ky mekanizëm është po ashtu i njohur, Spark lejon që të shkruhen partitë ndarazi.
Së pari, bëjmë ngarkimin inicial, duke ruajtur të dhënat ashtu siç është përshkruar më sipër, duke shtuar vetëm ndarjen. Ky veprim quhet inicializimi i ndriçimit dhe bëhet vetëm një herë:
df.write.partitionBy("date_load").mode("overwrite").parquet(dbpath + "/" + db + "/" + destTable)
Në ditën tjetër ngarkohet vetëm partizioni i ri:
df.coalesce(1).write.mode("overwrite").parquet(dbpath + "/" + db + "/" + destTable +"/date_load=" + date_load + "/")
Mbete vetëm të ri-regjistrohet në Hive për të përditësuar skemën.
Megjithatë, këtu dhe lindin problemet.
Problemi i parë. Një ditë ose tjetër, parquet i krijuar nuk do të mund të lexOHET. Kjo lidhet me mënyrën se si parquet dhe JSON i qasen fushave të zbrazëta ndryshe.
Le të marrim një situatë tipike. Për shembull, dje vjen JSON-i:
Dita 1: {"a": {"b": 1}},
ndërsa sot ky JSON duket kështu:
Dita 2: {"a": null}
Supozoni se kemi dy partitione të ndryshme, ku secila ka një rresht.
Kur lexojmë të dhënat origjinale tërësisht, Spark do të jetë në gjendje të përcaktojë tipin dhe do ta kuptojë se "a" është një fushë e tipit "strukturë", me një fushë të brendshme "b" të tipit INT. Por, nëse çdo partition është ruajtur veçmas, atëherë do të kemi parquet me skema të papajtueshme të partitioneve:
df1 (a: <struct>)
df2 (a: STRING NULLABLE)
Kjo situatĂ« Ă«shtĂ« e njohur mirĂ«, prandaj Ă«shtĂ« shtuar veçanĂ«risht njĂ« mundĂ«si â duke analizuar tĂ« dhĂ«nat origjinale, fshijmĂ« fushat e zbrazĂ«ta:
df = spark.read.json("...", dropFieldIfAllNull=True)
Në këtë rast, parquet do të përbëhet nga partitione që mund të lexohen së bashku.
Megjithatë, ata që e kanë bërë këtë në praktikë, do të qeshin me bitterness. Pse? Sepse me shumë mundësi do të shfaqen dy situata të tjera. Ose tre. Ose katër. E para, që do të ndodhi pothuajse me siguri, është se tipet numerike do të duken ndryshe në skedarë të ndryshëm JSON. Për shembull, {intField: 1} dhe {intField: 1.1}. Nëse këto fusha ndodhen në një partition, atëherë mërgimi i skemës do të lexojë gjithçka saktë, duke e çuar në tipin më të saktë. Por nëse janë në partitione të ndryshme, atëherë në një do të jetë intField: int, por në tjetrin intField: double.
Për të trajtuar këtë situatë, ekziston flamuri i mëposhtëm:
df = spark.read.json("...", dropFieldIfAllNull=True, primitivesAsString=True)
Tani kemi një dosje ku ndodhen partitionet që mund të lexohen në një dataframe të vetëm dhe një parquet valid të gjithë vitrinës. Po? Jo.
Duhet të kujtojmë se ne regjistruam tabelën në Hive. Hive nuk është i ndjeshëm ndaj regjistrit në emrat e fushave, ndërsa parquet është i ndjeshëm. Prandaj partitionet me skemat: field1: int dhe Field1: int janë të njëjta për Hive, por jo për Spark. Duhet të mos harrojmë të çojmë emrat e fushave në shkronja të vogla.
Pas kësaj, duket se gjithçka është mirë.
Megjithatë, nuk është gjithçka kaq e thjeshtë. Shfaqet problemi i dytë, gjithashtu i njohur mirë. Duke pasur parasysh se çdo partition i ri ruhet veçmas, dosja e partitioneve do të përmbajë skedarë shërbimi të Spark, për shembull, flamurin e suksesit të operacionit _SUCCESS. Kjo do të çojë në një gabim gjatë përpjekjes për parquet. Për të shmangur këtë, duhet të konfigurojmë konfigurimin, duke ndaluar Spark të shkruajë skedarë shërbimi në dosje:
hadoopConf = sc._jsc.hadoopConfiguration()
hadoopConf.set("parquet.enable.summary-metadata", "false")
hadoopConf.set("mapreduce.fileoutputcommitter.marksuccessfuljobs", "false")
Duket se tani çdo ditë po shtohet një pariticion e re parquet në dosjen e vitrinës përkatëse, ku ruhen të dhënat e parsezuara të ditës. Ne u kujdesëm paraprakisht që të mos kishte pariticione me konflikte të tipeve të dhënash.
Por, përpara nesh është një problem i tretë. Tani skema e përgjithshme nuk është e njohur, madje, në tabelën Hive ka një skemë të gabuar, pasi secila pariticion e re, me shumë mundësi, ka sjellë një distorcim në skemën.
Duhet të rregjistrojmë përsëri tabelën. Kjo mund të bëhet thjesht: të lexojmë sërish vitrinën parquet, të marrim skemën dhe të krijojmë mbi të DDL, me të cilin të regjistrojmë përsëri dosjen në Hive si një tabelë të jashtme, duke përditësuar skemën e vitrinës përkatëse.
Na paraqitet një problem i katërt. Kur regjistruam tabelën për herë të parë, ne u mbështetëm në Spark. Tani po e bëjmë këtë vetë, dhe duhet të mbajmë mend se fushat parquet mund të fillojnë me simbole, të papranueshme për Hive. Për shembull, Spark hedh jashtë rreshtat që nuk mund t'i parsezojë në fushën "corrupt_record". Një fushë e tillë nuk do të mund të regjistrohet në Hive pa u shkëputur.
Duke e ditur këtë, ne marrim skemën:
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)
Kodi ("_corrupt_record", "`_corrupt_record`") + " " + f[1].replace(":", "`:").replace("<", "<`").replace(",", ",`").replace("array<`", "array<") bën një DDL të sigurt, dmth përveç:
create table tname (_field1 string, 1field string)
Me emra fushe të tillë si "_field1, 1field", bëhet një DDL i sigurt, ku emrat e fushave janë të shkëputur: create table `tname` (`_field1` string, `1field` string).
Shqetësimi është: si të marrim saktësisht dataframe me skemën e plotë (në kodin pf)? Si mund të marrim këtë pf? Kjo është një problem i pestë. A duhet të lexojmë skemën e të gjitha pariticioneve nga dosja me skedarë parquet të vitrinës përkatëse? Ky është metoda më e sigurt, por e rëndë.
Schemi është tashmë në Hive. Për të marrë një schemë të re, duhet të kombinohet schemën e gjithë tabelës me një parti të re. Kështu që duhet të marrim schemën e tabelës nga Hive dhe ta kombinojmë atë me schemën e partisë së re. Kjo mund të bëhet duke lexuar metadatat testuese nga Hive, duke i ruajtur ato në një dosje përkohshme dhe duke lexuar me Spark të dy partitë njëherësh.
Në thelb, ka gjithçka që na nevojitet: schemën origjinale të tabelës në Hive dhe një parti të re. Të dhënat gjithashtu i kemi. Na mbetet vetëm të marrim një schemë të re, ku kombinohet schemën e vitrinës dhe fushat e reja nga partia e krijuar:
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/*")
Më pas krijojmë DDL-in për regjistrimin e tabelës, si në fragmentin e mëparshëm.
Nëse e gjithë zinxhirja funksionon siç duhet, konkretisht - u bë ngarkimi fillestar, dhe tabela e krijuar në Hive është e saktë, atëherë ne marrim një schemë të përditësuar të tabelës.
Dhe problemi i fundit Ă«shtĂ« se nuk mund ta shtoni thjesht njĂ« parti nĂ« tabelĂ«n Hive, pasi ajo do tĂ« prishet. ĂshtĂ« e nevojshme tĂ« bĂ«jmĂ« qĂ« Hive tĂ« rregullojĂ« strukturĂ«n e saj tĂ« partive:
from pyspark.sql import HiveContext
hc = HiveContext(spark)
hc.sql("MSCK REPAIR TABLE " + db + "." + destTable)
Një detyrë e thjeshtë për të lexuar JSON dhe për të krijuar një vitrinë mbi të, shndërrohet në përballjen me një numër të vështirësish të pacaktuara, zgjidhjet e të cilave duhet të kërkohen veç e veç. Dhe ndonëse këto zgjidhje janë të thjeshta, kërkon shumë kohë për t'i gjetur.
Për të realizuar ndërtimin e vitrinës, ishte e nevojshme:
- Të shtoni partinë në vitrinë, duke e hequr dosjet shërbyese
- Të merret me fushat e zbrazëta në të dhënat origjinale, të cilat Spark i kishte tipizuar
- Të konvertohen llojet e thjeshta në varg
- Të sjellim emrat e fushave në shkronja të vogla
- Të ndajmë eksportimin e të dhënave dhe regjistrimin e tabelës në Hive (krijimi i DDL)
- Të mos harrohet të shkruhen emrat e fushave, të cilat mund të jenë të papërshtatshme për Hive
- Të mësojmë si të përditësojmë regjistrimin e tabelës në Hive
Duke përmbledhur, vërejmë se zgjidhja për ndërtimin e vitrinës fsheh shumë pengesa. Prandaj, në rast të vështirësive në zbatim, është më mirë të kontaktoni një partner të përvojës me ekspertizë të suksesshme.
Faleminderit për leximin e këtij artikulli, shpresojmë se informacioni do të jetë i dobishëm.
Burimi: habr.com
Blini hosting tĂ« besueshĂ«m pĂ«r faqe interneti me mbrojtje nga DDoS, serverĂ« VPS VDS đ„ Blini hosting tĂ« besueshĂ«m pĂ«r faqe interneti me mbrojtje nga DDoS, serverĂ« VPS VDS | ProHoster
