Chers lecteurs, bonjour !
Dans cet article, un consultant principal de la direction commerciale Big Data Solutions de la société « Neoflex » décrit en détail les options de construction de vitrines à structure variable en utilisant Apache Spark.
Dans le cadre d'un projet d'analyse de données, il se pose souvent la question de la construction de vitrines à partir de données faiblement structurées.
Il s'agit généralement de journaux ou de réponses de divers systèmes, stockés sous forme de JSON ou XML. Les données sont exportées vers Hadoop, puis une vitrine doit être construite à partir de celles-ci. Nous pouvons organiser l'accès à la vitrine créée, par exemple, via Impala.
Dans ce cas, le schéma de la vitrine cible est préalablement inconnu. De plus, le schéma ne peut pas être établi à l'avance, car il dépend des données, et nous traitons ces données faiblement structurées.
Par exemple, aujourd'hui, une réponse comme celle-ci est enregistrée :
{source: "app1", error_code: ""}et demain, le même système envoie cette réponse :
{source: "app1", error_code: "error", description: "Erreur de réseau"}En conséquence, un nouveau champ — description — doit être ajouté à la vitrine, et personne ne sait s'il sera présent ou non.
La tâche de création d'une vitrine à partir de telles données est assez standard, et Spark dispose de plusieurs outils pour cela. Pour le parsing des données sources, le support de JSON et XML est disponible, et pour un schéma inconnu à l'avance, il existe un support pour la schemaEvolution.
À première vue, la solution semble simple. Il suffit de prendre un dossier contenant des JSON et de le lire dans un dataframe. Spark créera le schéma, et les données imbriquées seront transformées en structures. Ensuite, tout doit être enregistré au format parquet, qui est également pris en charge par Impala, en enregistrant la vitrine dans le metastore Hive.
Il semble donc que tout soit simple.
Cependant, les courts exemples dans la documentation ne clarifient pas comment traiter un certain nombre de problèmes en pratique.
La documentation décrit une approche non pas pour créer une vitrine, mais pour lire JSON ou XML dans un dataframe.
À savoir, il suffit de montrer comment lire et parser un JSON :
df = spark.read.json(path...)C'est suffisant pour rendre les données accessibles à Spark.
En pratique, le scénario est beaucoup plus complexe que de simplement lire des fichiers JSON d'un dossier et de créer un dataframe. La situation est la suivante : une vitrine est déjà définie, de nouvelles données arrivent chaque jour, elles doivent être ajoutées à la vitrine en n'oubliant pas que le schéma peut différer.
Le schéma habituel de construction de vitrine est le suivant :
Étape 1. Les données sont chargées dans Hadoop avec un chargement quotidien ultérieur et sont placées dans une nouvelle partition. Cela crée un dossier partitionné par jour avec les données d'origine.
Étape 2. Au cours du chargement initial, ce dossier est lu et analysé à l'aide de Spark. Le dataframe obtenu est sauvegardé dans un format accessible pour l'analyse, par exemple en parquet, qui peut ensuite être importé dans Impala. Cela crée une vitrine cible avec toutes les données accumulées jusqu'à ce moment.
Étape 3. Une tâche est créée qui mettra à jour la vitrine tous les jours.
La question de la charge incrémentielle se pose, ainsi que la nécessité de partitionner la vitrine et de soutenir le schéma général de celle-ci.
Prenons un exemple. Supposons que la première étape de construction du stockage soit réalisée et que l'exportation de fichiers JSON dans le dossier soit configurée.
Créer un dataframe à partir d'eux pour ensuite l'enregistrer comme vitrine ne pose pas de problème. C'est cette première étape que l'on peut facilement trouver dans la documentation de 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)Tout semble en ordre.
Nous avons lu et analysé le JSON, ensuite nous sauvegardons le dataframe au format parquet, en l'enregistrant dans Hive de la manière la plus simple :
df.write.format(“parquet”).option('path','').saveAsTable('') Nous obtenons la vitrine.
Mais, le lendemain, de nouvelles données ont été ajoutées à partir de la source. Nous avons un dossier avec JSON, et une vitrine créée à partir de ce dossier. Après le chargement de la prochaine tranche de données de la source, il manque des données pour un jour dans la vitrine.
La solution logique serait de partitionner la vitrine par jour, ce qui permettra chaque jour d'ajouter une nouvelle partition. Ce mécanisme est également bien connu, Spark permet d'écrire des partitions séparément.
Tout d'abord, nous effectuons un chargement initial, en sauvegardant les données comme décrit ci-dessus, en ajoutant simplement la partition. Cette action s'appelle l'initialisation de la vitrine et ne se fait qu'une seule fois :
df.write.partitionBy("date_load").mode("overwrite").parquet(dbpath + "/" + db + "/" + destTable)
Le lendemain, nous ne chargeons que la nouvelle partition :
df.coalesce(1).write.mode("overwrite").parquet(dbpath + "/" + db + "/" + destTable +"/date_load=" + date_load + "/")
Il ne reste plus qu'à réenregistrer dans Hive pour mettre à jour le schéma.
Cependant, c'est là que les problèmes surviennent.
Le premier problème. Tôt ou tard, le parquet obtenu ne sera plus lisible. Cela est dû à la manière dont le parquet et le JSON abordent les champs vides différemment.
Considérons une situation typique. Par exemple, hier, nous recevons ce JSON :
Jour 1 : {"a": {"b": 1}},
et aujourd'hui, ce même JSON ressemble à ceci :
Jour 2 : {"a": null}
Supposons que nous ayons deux partitions différentes, chacune contenant une ligne.
Lorsque nous lisons les données source en entier, Spark saura déterminer le type et comprendra que « a » est un champ de type « structure », avec un champ imbriqué « b » de type INT. Mais, si chaque partition a été sauvegardée séparément, cela entraîne des parquet avec des schémas de partition incompatibles :
df1 (a: <struct>)
df2 (a: STRING NULLABLE)
Cette situation est bien connue, c'est pourquoi une option a été spécialement ajoutée — lors de l'analyse des données source, supprimer les champs vides :
df = spark.read.json("...", dropFieldIfAllNull=True)
Dans ce cas, le parquet se composera de partitions qui pourront être lues ensemble.
Cependant, ceux qui l'ont fait dans la pratique vont ici sourire amèrement. Pourquoi ? Parce qu'il est probable qu'il y ait encore deux situations. Ou trois. Ou quatre. La première, qui se produira presque certainement, est que les types numériques apparaîtront différemment dans différents fichiers JSON. Par exemple, {intField: 1} et {intField: 1.1}. Si de tels champs se retrouvent dans une même partition, la fusion des schémas lira tout correctement, en les adaptant au type le plus précis. En revanche, s'ils se trouvent dans des partitions différentes, alors, dans l'une, intField sera de type int, et dans l'autre, intField sera de type double.
Pour gérer cette situation, il existe le drapeau suivant :
df = spark.read.json("...", dropFieldIfAllNull=True, primitivesAsString=True)
Nous avons maintenant un dossier contenant des partitions que nous pouvons lire dans un seul dataframe et un parquet valide de l'intégralité de la vitrine. Oui ? Non.
Il faut se rappeler que nous avons enregistré la table dans Hive. Hive n'est pas sensible à la casse des noms de champs, tandis que parquet l'est. Par conséquent, les partitions avec des schémas : field1: int, et Field1: int sont identiques pour Hive, mais pas pour Spark. Il ne faut pas oublier de mettre les noms des champs en minuscules.
Après cela, il semble que tout soit en ordre.
Cependant, tout n'est pas si simple. Une deuxième, également bien connue, se pose. Comme chaque nouvelle partition est sauvegardée séparément, le dossier de partition contiendra des fichiers de service Spark, comme le drapeau de succès de l'opération _SUCCESS. Cela entraînera une erreur lors de la tentative de parquet. Pour éviter cela, il faut configurer la configuration pour interdire à Spark d'écrire des fichiers de service dans le dossier :
hadoopConf = sc._jsc.hadoopConfiguration()
hadoopConf.set("parquet.enable.summary-metadata", "false")
hadoopConf.set("mapreduce.fileoutputcommitter.marksuccessfuljobs", "false")
Il semble qu'une nouvelle partition parquet avec les données analysées du jour soit ajoutée chaque jour dans le dossier de la vitrine cible. Nous avons déjà pris soin d'éviter les partitions avec des conflits de types de données.
Cependant, nous faisons face à un troisième problème. Le schéma global n'est plus connu, et de plus, dans Hive, la table a un schéma incorrect, car chaque nouvelle partition a probablement déformé le schéma.
Il est nécessaire de réenregistrer la table. Cela peut être fait simplement : lire à nouveau les vitrines parquet, obtenir le schéma et créer à partir de celui-ci DDL afin de réenregistrer le dossier dans Hive comme une table externe, en mettant à jour le schéma de la vitrine cible.
Nous rencontrons un quatrième problème. Lorsque nous avons enregistré la table pour la première fois, nous nous sommes appuyés sur Spark. Maintenant, nous le faisons nous-mêmes, et il faut se rappeler que les champs parquet peuvent commencer par des caractères non valides pour Hive. Par exemple, Spark ignore les lignes qui n'ont pas pu être analysées dans le champ «corrupt_record». Ce champ ne pourra pas être enregistré dans Hive sans que cela soit échappé.
Sachant cela, nous obtenons le schéma :
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<» fait un DDL sécurisé, c'est-à-dire qu'au lieu de :
create table tname (_field1 string, 1field string)
Avec des noms de champs tels que «_field1, 1field», un DDL sécurisé est créé, où les noms des champs sont échappés : create table `tname` (`_field1` string, `1field` string).
La question se pose : comment obtenir correctement un dataframe avec le schéma complet (dans le code pf) ? Comment obtenir ce pf ? C'est un cinquième problème. Relire le schéma de toutes les partitions du dossier contenant les fichiers parquet de la vitrine cible ? C'est la méthode la plus sûre, mais elle est lourde.
Le schéma est déjà présent dans Hive. Pour obtenir un nouveau schéma, il faut combiner le schéma de toute la table avec celui de la nouvelle partition. Il est donc nécessaire de récupérer le schéma de la table depuis Hive et de le combiner avec le schéma de la nouvelle partition. Ceci peut être réalisé en lisant les métadonnées de test depuis Hive, en les sauvegardant dans un dossier temporaire et en lisant les deux partitions simultanément à l'aide de Spark.
En fait, tout ce dont nous avons besoin est là : le schéma original de la table dans Hive et la nouvelle partition. Nous avons également les données. Il reste seulement à obtenir le nouveau schéma, qui combine le schéma de la vitrine et les nouveaux champs de la partition créée :
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/*")
Nous créons ensuite le DDL d'enregistrement de la table, comme dans le fragment précédent.
Si toute la chaîne fonctionne correctement, à savoir — s'il y a eu un chargement initial et si la table a été correctement créée dans Hive, alors nous obtenons le schéma mis à jour de la table.
Enfin, le problème réside dans le fait qu'il n'est pas si simple d'ajouter une partition à une table Hive, car cela va casser la table. Il faut donc forcer Hive à réparer sa structure de partitions :
from pyspark.sql import HiveContext
hc = HiveContext(spark)
hc.sql("MSCK REPAIR TABLE " + db + "." + destTable)
Une tâche aussi simple que la lecture de JSON et la création d'une vitrine se transforme en un parcours d'obstacles où il faut trouver des solutions à plusieurs difficultés implicites. Bien que ces solutions soient simples, le temps nécessaire pour les trouver est considérable.
Pour réaliser la construction de la vitrine, il a fallu :
- Ajouter des partitions à la vitrine, en éliminant les fichiers temporaires
- S'attaquer aux champs vides dans les données d'origine, que Spark a typés
- Convertir les types simples en chaîne
- Mettre les noms des champs en lettres minuscules
- Séparer l'extraction des données et l'enregistrement de la table dans Hive (création de DDL)
- Ne pas oublier d'échapper les noms de champs qui pourraient être incompatibles avec Hive
- Apprendre à mettre à jour l'enregistrement de la table dans Hive
En résumé, il convient de noter que la solution pour construire des vitrines présente de nombreux pièges. Par conséquent, en cas de difficultés d'implémentation, il est préférable de faire appel à un partenaire expérimenté disposant d'une expertise réussie.
Merci d'avoir lu cet article, nous espérons que l'information vous sera utile.
Source : habr.com
Acheter un hébergement fiable pour les sites avec protection DDoS, serveurs VPS VDS 🔥 Acheter un hébergement fiable pour les sites avec protection DDoS, serveurs VPS VDS | ProHoster
