Spark schemaEvolution in de praktijk

Beste lezers, goedemiddag!

In dit artikel beschrijft de hoofdadviseur van de businessunit Big Data Solutions van 'Neoflex' in detail de opties voor het bouwen van een variabele structuur vitrine met behulp van Apache Spark.

In het kader van een project voor data-analyse komt vaak de taak naar voren om vitrines te bouwen op basis van slecht gestructureerde gegevens.

Dit zijn meestal logs of antwoorden van verschillende systemen, opgeslagen in JSON- of XML-formaat. De gegevens worden geƫxporteerd naar Hadoop, waarna een vitrine moet worden gebouwd. Toegang tot de gemaakte vitrine kan bijvoorbeeld via Impala worden georganiseerd.

In dit geval is het schema van de doelvitrine vooraf niet bekend. Bovendien kan het schema niet van tevoren worden opgesteld, omdat dit afhangt van de gegevens en we met deze slecht gestructureerde gegevens te maken hebben.

Bijvoorbeeld, vandaag wordt het volgende antwoord gelogd:

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

en morgen ontvangt dezezelfde systeem het volgende antwoord:

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

Als resultaat moet er nog een veld — description — aan de vitrine worden toegevoegd, en of deze al dan niet wordt ontvangen, dat weet niemand.

De taak om een vitrine op dergelijke gegevens te creƫren is vrij standaard, en Spark heeft hiervoor een aantal tools. Voor het parseren van de oorspronkelijke gegevens is er ondersteuning voor zowel JSON als XML, en voor een schema dat vooraf niet bekend is, is er ondersteuning voor schemaEvolution.

Op het eerste gezicht lijkt de oplossing eenvoudig. Je moet gewoon de map met JSON nemen en deze in een dataframe lezen. Spark maakt een schema aan en zet geneste gegevens om in structuren. Vervolgens moet alles worden opgeslagen in parquet, dat ook door Impala wordt ondersteund, waarbij de vitrine in de Hive metastore wordt geregistreerd.

Het lijkt allemaal eenvoudig.

Echter, uit de korte voorbeelden in de documentatie is niet duidelijk wat je in de praktijk met een aantal problemen moet doen.

De documentatie beschrijft de benadering niet voor het creƫren van een vitrine, maar voor het lezen van JSON of XML in een dataframe.

Specifiek wordt er gewoon aangegeven hoe JSON te lezen en te parseren:

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

Dit is voldoende om de gegevens toegankelijk te maken voor Spark.

In de praktijk is het scenario veel complexer dan gewoon JSON-bestanden uit een map te lezen en een dataframe te creƫren. De situatie is als volgt: er is al een bepaalde vitrine, elke dag komen er nieuwe gegevens binnen die aan de vitrine moeten worden toegevoegd, zonder te vergeten dat het schema kan verschillen.

Het gebruikelijke schema voor het bouwen van een vitrine is als volgt:

Stap 1. Gegevens worden in Hadoop geladen met dagelijkse incrementele uploads en worden opgeslagen in een nieuwe partitie. Hierdoor ontstaat een partitie met gegevens per dag.

Stap 2. Tijdens de initiële upload wordt deze map gelezen en geparsed met Spark. De verkregen dataframe wordt opgeslagen in een formaat dat geschikt is voor analyse, zoals parquet, dat vervolgens in Impala kan worden geïmporteerd. Zo wordt een doelfront ingericht met alle gegevens die tot dat moment zijn verzameld.

Stap 3. Er wordt een upload gemaakt die elke dag de frontend bijwerkt.
Er komt een vraag op over incrementele uploads, de noodzaak van partionering van de frontend, en de ondersteuning van het algemene schema van de frontend.

Laten we een voorbeeld geven. Stel dat de eerste stap van de opslag is geĆÆmplementeerd, en dat de export van JSON-bestanden naar de map is ingesteld.

Het is geen probleem om een dataframe van hen te maken om vervolgens als frontend op te slaan. Dit is de eerste stap die makkelijk te vinden is in de Spark-documentatie:

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)

Het lijkt allemaal goed te gaan.

We hebben de JSON gelezen en geparsed, daarna slaan we de dataframe op als parquet, en registreren deze in Hive op een manier die ons het beste uitkomt:

df.write.format(ā€œparquetā€).option('path','').saveAsTable('')

We krijgen de frontend.

Maar, op de volgende dag zijn er nieuwe gegevens uit de bron toegevoegd. We hebben een map met JSON, en de frontend die is gemaakt op basis van deze map. Na het uploaden van de volgende gegevens uit de bron ontbreken er gegevens voor ƩƩn dag in de frontend.

Een logische oplossing zou zijn om de frontend per dag te partitioneren, wat het mogelijk maakt om elke volgende dag een nieuwe partitie toe te voegen. Dit mechanisme is ook goed bekend; Spark stelt ons in staat om partities afzonderlijk op te slaan.

Eerst doen we een initiƫle upload, waarbij we de gegevens opslaan zoals hierboven beschreven, met alleen partionering toegevoegd. Deze actie heet het initialiseren van de frontend en wordt slechts ƩƩn keer uitgevoerd:

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

Op de volgende dag uploaden we alleen de nieuwe partitie:

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

Het enige dat overblijft is om het opnieuw te registreren in Hive om het schema bij te werken.
Echter, hier ontstaan problemen.

Het eerste probleem. Vroeg of laat kan de resulterende parquet niet worden gelezen. Dit heeft te maken met de verschillende benaderingen van parquet en JSON voor lege velden.

Laten we een typische situatie bekijken. Bijvoorbeeld, gisteren komt de JSON binnen:

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

en vandaag ziet dezelfde JSON er zo uit:

Dag 2: {"a": null}

Stel dat we twee verschillende partities hebben, elk met ƩƩn regel.
Wanneer we de oorspronkelijke gegevens in hun geheel lezen, kan Spark het type bepalen en begrijpen dat "a" een veld is van het type "structuur", met een genest veld "b" van het type INT. Maar als elke partitie afzonderlijk werd opgeslagen, resulteert dit in parquet met onverenigbare partitieschema's:

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

Deze situatie is goed bekend, daarom is er speciaal een optie toegevoegd — bij het parseren van de oorspronkelijke gegevens leegt men lege velden:

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

In dit geval bestaat parquet uit partities die samen gelezen kunnen worden.
Hoewel degenen die dit in de praktijk hebben gedaan, hier bitter om zullen lachen. Waarom? Omdat er waarschijnlijk nog twee situaties zullen ontstaan. Of drie. Of vier. De eerste die bijna zeker zal optreden, is dat numerieke types er in verschillende JSON-bestanden anders uitzien. Bijvoorbeeld, {intField: 1} en {intField: 1.1}. Als dergelijke velden in ƩƩn partitie voorkomen, zal de schema-samenvoeging alles correct lezen, leidend tot het meest precieze type. Maar als ze in verschillende zijn, zal in de ene intField: int zijn en in de andere intField: double.

Voor de verwerking van deze situatie is de volgende vlag beschikbaar:

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

Nu hebben we een map waarin de partities zich bevinden, die we in een enkele dataframe en een geldige parquet van de gehele vitrine kunnen lezen. Toch?

Nee. We moeten niet vergeten dat we de tabel in Hive hebben geregistreerd. Hive is niet gevoelig voor hoofdletters in veldnamen, terwijl parquet dat wel is. Daarom zijn partities met schema's: field1: int, en Field1: int voor Hive hetzelfde, maar voor Spark niet. We moeten niet vergeten de veldnamen naar kleine letters om te zetten.

Het lijkt erop dat na dit alles alles in orde is.

Toch is niet alles zo eenvoudig. Er ontstaat een tweede, ook goed bekende probleem. Aangezien elke nieuwe partitie afzonderlijk wordt opgeslagen, zullen in de partitiemap systeembestanden van Spark liggen, zoals de succesflag van de operatie _SUCCESS. Dit zal leiden tot een fout bij de poging om parquet te maken. Om dit te vermijden, moet de configuratie worden ingesteld om Spark te verbieden systeembestanden naar de map te schrijven:

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

Het lijkt erop dat er nu elke dag een nieuwe parquet-partitie aan de map van de doelvitrine wordt toegevoegd, waar de geparsed gegevens van de dag zich bevinden. We hebben ons voorbereid om conflicterende datatypes te voorkomen.

Maar nu hebben we een derde probleem. Het algemene schema is onbekend, bovendien is er in de Hive-tabel een verkeerd schema, aangezien elke nieuwe partitie waarschijnlijk een vervorming in het schema heeft geĆÆntroduceerd.

We moeten de tabel opnieuw registreren. Dit kan eenvoudig worden gedaan: lees de parquet-vitrine opnieuw, neem het schema en maak op basis daarvan een DDL aan, waarmee we de map opnieuw registreren in Hive als een externe tabel en het schema van de doelvitrine bijwerken.

Daarbij komt nog een vierde probleem. Toen we de tabel voor de eerste keer registreerden, vertrouwden we op Spark. Nu doen we het zelf, en we moeten ons herinneren dat parquet-velden kunnen beginnen met tekens die niet zijn toegestaan in Hive. Bijvoorbeeld, Spark verwijdert rijen die het niet heeft kunnen parsen in het veld "corrupt_record". Zo'n veld kan niet in Hive worden geregistreerd zonder te escapen.

Wetende dit, verkrijgen we het 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)

Code ("_corrupt_record", "`_corrupt_record`") + " " + f[1].replace(":", "`:").replace("<", "<`").replace(",", ",`").replace("array<`", "array<") maakt een veilige DDL, dus in plaats van:

create table tname (_field1 string, 1field string)

Met veldnamen zoals "_field1, 1field" wordt er een veilige DDL gemaakt, waarin de veldnamen worden geescaped: create table `tname` (`_field1` string, `1field` string).

De vraag rijst: hoe krijg je correct een dataframe met een volledig schema (in de code pf)? Hoe krijg je deze pf? Dit is een vijfde probleem. Moet je het schema van alle partities in de map met parquet-bestanden van de doelvitrine opnieuw lezen? Dit is de veiligste methode, maar ook de zwaarste.

Het schema is al aanwezig in Hive. Een nieuw schema kan worden verkregen door het schema van de hele tabel te combineren met het nieuwe partition. Dit betekent dat we het schema van de tabel uit Hive moeten halen en het combineren met het schema van de nieuwe partition. Dit kan worden gedaan door de testmetadata uit Hive te lezen, deze op te slaan in een tijdelijke map en met Spark beide partitions tegelijkertijd te lezen.

In wezen is alles wat nodig is aanwezig: het oorspronkelijke schema van de tabel in Hive en de nieuwe partition. We hebben ook de gegevens. We moeten alleen het nieuwe schema verkrijgen waarin het schema van de vitrine en de nieuwe velden van de aangemaakte partition worden gecombineerd:

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

Vervolgens creƫren we de DDL voor het registreren van de tabel, zoals in het vorige fragment.
Als de gehele keten correct werkt, namelijk - er was een initiƫle laadoperatie en de tabel is correct aangemaakt in Hive, dan krijgen we het bijgewerkte schema van de tabel.

En ten slotte is het probleem dat je niet zomaar een partition aan een Hive-tabel kunt toevoegen, omdat deze zou kunnen breken. Het is noodzakelijk om Hive te dwingen de structuur van de partitions te repareren:

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

Een eenvoudige taak van het lezen van JSON en het creƫren van een vitrine daarop, resulteert in het overwinnen van verschillende onduidelijke moeilijkheden, waarvoor we afzonderlijke oplossingen moeten zoeken. En hoewel deze oplossingen simpel zijn, kost het veel tijd om ze te vinden.

Om de bouw van de vitrine te realiseren, moesten we:

  • Partitions aan de vitrine toevoegen, waarbij we ons ontdeden van servicebestanden
  • Begrijpen hoe om te gaan met lege velden in de oorspronkelijke gegevens, die door Spark werden getypeerd
  • Eenvoudige types omzetten naar strings
  • Veldnamen naar kleine letters converteren
  • Het uitvoeren van gegevensuitvoer en het registreren van de tabel in Hive (creatie van DDL) scheiden
  • Vergeet niet om veldnamen te escapen die mogelijk niet compatibel zijn met Hive
  • Leer hoe je de registratie van de tabel in Hive kunt bijwerken

Samenvattend, het is belangrijk op te merken dat oplossingen voor het bouwen van vitrines veel verborgen valkuilen met zich meebrengen. Daarom is het beter om bij moeilijkheden in de uitvoering contact op te nemen met een ervaren partner met succesvolle expertise.

Bedankt voor het lezen van dit artikel. We hopen dat de informatie nuttig zal zijn.

Bron: habr.com

Koop betrouwbare webhosting met bescherming tegen DDoS, VPS VDS servers šŸ”„ Koop betrouwbare webhosting met bescherming tegen DDoS, VPS VDS servers | ProHoster