Aller au contenu principal

Modèles courants de chargement de données

Auto Loader simplifie un certain nombre de tâches courantes d’ingestion de données. Cette référence rapide fournit des exemples pour plusieurs modèles populaires.

Ingérer les données depuis le stockage d’objets cloud en tant que variant

Auto Loader peut charger toutes les données des sources de fichiers prises en charge sous la forme d'une seule colonne VARIANT dans une table cible. Parce que VARIANT est flexible aux changements de schéma et de type et maintient la sensibilité à la casse et les valeurs NULL présentes dans la source de données, ce modèle est robuste à la plupart des scénarios d'ingestion. Pour plus de détails, consultez Ingérer des données depuis le stockage d'objets cloud en tant que variante.

Filtrage des répertoires ou des fichiers à l'aide de modèles glob

Les modèles Glob peuvent être utilisés pour filtrer les répertoires et les fichiers lorsqu'ils sont fournis dans le chemin.

Modèle

Description

?

Correspond à n'importe quel caractère unique.

*

Correspond à zéro ou plusieurs caractères

[abc]

Correspond à un seul caractère du jeu de caractères {a,b,c}.

[a-z]

Correspond à un seul caractère de la plage de caractères {a…z}.

[^a]

Correspond à un seul caractère qui n'appartient pas à un jeu de caractères ou à une plage {a}. Notez que le caractère ^ doit apparaître immédiatement à droite du crochet ouvrant.

{ab,cd}

Correspond à une chaîne de l'ensemble de chaînes {ab, cd}.

{ab,c{de, fh}}

Correspond à une chaîne de caractères de l'ensemble de chaînes de caractères {ab, cde, cfh}.

Modèle

Description

?

Correspond à n'importe quel caractère unique.

*

Correspond à zéro ou plusieurs caractères

[abc]

Correspond à un seul caractère du jeu de caractères {a,b,c}.

[a-z]

Correspond à un seul caractère de la plage de caractères {a…z}.

[^a]

Correspond à un seul caractère qui n'appartient pas à un jeu de caractères ou à une plage {a}. Notez que le caractère ^ doit apparaître immédiatement à droite du crochet ouvrant.

{ab,cd}

Correspond à une chaîne de l'ensemble de chaînes {ab, cd}.

{ab,c{de, fh}}

Correspond à une chaîne de caractères de l'ensemble de chaînes de caractères {ab, cde, cfh}.

Utilisez le path pour fournir des modèles de préfixe, par exemple :

Python
df = spark.readStream.format("cloudFiles") \
.option("cloudFiles.format", <format>) \
.schema(schema) \
.load("/Volumes/catalog_name/schema_name/volume_name/*/files")

Vous devez utiliser l'option pathGlobFilter pour fournir explicitement des modèles de suffixe. Le path fournit uniquement un filtre de préfixe. Par exemple, si vous souhaitez analyser uniquement png fichiers dans un répertoire contenant des fichiers avec des suffixes différents, vous pouvez faire ce qui suit :

Python
df = spark.readStream.format("cloudFiles") \
.option("cloudFiles.format", "binaryFile") \
.option("pathGlobfilter", "*.png") \
.load("/Volumes/catalog_name/schema_name/volume_name/path")
remarque

Le comportement de globbing par default d'Auto Loader est différent du comportement par default des autres sources de fichiers Spark. Ajoutez .option("cloudFiles.useStrictGlobber", "true") à votre lecture pour utiliser le globbing qui correspond au comportement Spark par default pour les sources de fichiers. Reportez-vous au tableau suivant pour en savoir plus sur le globbing :

Modèle

Chemin d'accès au fichier

default globber

globber strict

/a/b

/a/b/c/file.txt

Oui

Oui

/a/b

/a/b_dir/c/file.txt

Non

Non

/a/b

/a/b.txt

Non

Non

/a/b/

/a/b.txt

Non

Non

/a/*/c/

/a/b/c/file.txt

Oui

Oui

/a/*/c/

/a/b/c/d/file.txt

Oui

Oui

/a/*/c/

/a/b/x/y/c/file.txt

Oui

Non

/a/*/c

/a/b/c_file.txt

Oui

Non

/a/*/c/

/a/b/c_file.txt

Oui

Non

/a/*/c/

/a/*/cookie/file.txt

Oui

Non

/a/b*

/a/b.txt

Oui

Oui

/a/b*

/a/b/file.txt

Oui

Oui

/a/{0.txt,1.txt}

/a/0.txt

Oui

Oui

/a/*/{0.txt,1.txt}

/a/0.txt

Non

Non

/a/b/[cde-h]/i/

/a/b/c/i/file.txt

Oui

Oui

Modèle

Chemin d'accès au fichier

default globber

globber strict

/a/b

/a/b/c/file.txt

Oui

Oui

/a/b

/a/b_dir/c/file.txt

Non

Non

/a/b

/a/b.txt

Non

Non

/a/b/

/a/b.txt

Non

Non

/a/*/c/

/a/b/c/file.txt

Oui

Oui

/a/*/c/

/a/b/c/d/file.txt

Oui

Oui

/a/*/c/

/a/b/x/y/c/file.txt

Oui

Non

/a/*/c

/a/b/c_file.txt

Oui

Non

/a/*/c/

/a/b/c_file.txt

Oui

Non

/a/*/c/

/a/*/cookie/file.txt

Oui

Non

/a/b*

/a/b.txt

Oui

Oui

/a/b*

/a/b/file.txt

Oui

Oui

/a/{0.txt,1.txt}

/a/0.txt

Oui

Oui

/a/*/{0.txt,1.txt}

/a/0.txt

Non

Non

/a/b/[cde-h]/i/

/a/b/c/i/file.txt

Oui

Oui

Permettre un ETL facile

Un moyen facile d'intégrer vos données dans Delta Lake sans perdre de données consiste à utiliser le modèle suivant et à activer l'inférence de schéma avec Auto Loader. Databricks vous recommande d'exécuter le code suivant dans un Job Databricks afin qu'il redémarre automatiquement votre Stream lorsque le schéma de vos données source change. By default, le schéma est déduit comme des types de chaîne, toutes les erreurs d'analyse (il ne devrait y en avoir aucune si tout reste une chaîne) iront à _rescued_data, et toutes les nouvelles colonnes feront échouer le Stream et feront évoluer le schéma.

Python
spark.readStream.format("cloudFiles") \
.option("cloudFiles.format", "json") \
.option("cloudFiles.schemaLocation", "<path-to-schema-location>") \
.load("/Volumes/catalog_name/schema_name/volume_name/source_data") \
.writeStream \
.option("mergeSchema", "true") \
.option("checkpointLocation", "<path-to-checkpoint>") \
.start("<path_to_target>")

Prévenir la perte de données dans des données bien structurées

Lorsque vous connaissez votre schéma, mais que vous souhaitez capturer des données inattendues, Databricks recommande d’utiliser le rescuedDataColumn.

Python
spark.readStream.format("cloudFiles") \
.schema(expected_schema) \
.option("cloudFiles.format", "json") \
# will collect all new fields as well as data type mismatches in _rescued_data
.option("cloudFiles.schemaEvolutionMode", "rescue") \
.load("/Volumes/catalog_name/schema_name/volume_name/source_data") \
.writeStream \
.option("checkpointLocation", "<path-to-checkpoint>") \
.start("<path_to_target>")

Si vous souhaitez que votre stream cesse de traiter si un nouveau champ est introduit qui ne correspond pas à votre schéma, vous pouvez ajouter :

Python
.option("cloudFiles.schemaEvolutionMode", "failOnNewColumns")

Activer des pipelines de données semi-structurées flexibles

Lorsque vous recevez des données d'un fournisseur qui introduit de nouvelles colonnes dans les informations qu'il fournit, vous pouvez ne pas savoir exactement quand il le fait, ou vous pouvez ne pas avoir la bande passante nécessaire pour mettre à jour votre pipeline de données. Vous pouvez maintenant tirer parti de l'évolution des schémas pour redémarrer le Stream et laisser Auto Loader mettre à jour automatiquement le schéma inféré. Vous pouvez également exploiter schemaHints pour certains des champs « sans schéma » que le fournisseur peut fournir.

Python
spark.readStream.format("cloudFiles") \
.option("cloudFiles.format", "json") \
# will ensure that the headers column gets processed as a map
.option("cloudFiles.schemaHints",
"headers map<string,string>, statusCode SHORT") \
.load("/Volumes/catalog_name/schema_name/volume_name/api/requests") \
.writeStream \
.option("mergeSchema", "true") \
.option("checkpointLocation", "<path-to-checkpoint>") \
.start("<path_to_target>")

Transformer des données JSON imbriquées

Puisque Auto Loader infère les colonnes JSON de niveau supérieur en tant que chaînes, il peut vous rester des objets JSON imbriqués qui nécessitent des transformations supplémentaires. Vous pouvez utiliser les APIs d'accès aux données semi-structurées pour transformer davantage le contenu JSON complexe.

Python
spark.readStream.format("cloudFiles") \
.option("cloudFiles.format", "json") \
# The schema location directory keeps track of your data schema over time
.option("cloudFiles.schemaLocation", "<path-to-checkpoint>") \
.load("/Volumes/catalog_name/schema_name/volume_name/nested_json") \
.selectExpr(
"*",
"tags:page.name", # extracts {"tags":{"page":{"name":...
"tags:page.id::int", # extracts {"tags":{"page":{"id":... and casts to int
"tags:eventType" # extracts {"tags":{"eventType":...}}
)

Inférer les données JSON imbriquées

Lorsque vous avez des données imbriquées, vous pouvez utiliser l'option cloudFiles.inferColumnTypes pour inférer la structure imbriquée de vos données et d'autres types de colonnes.

Python
spark.readStream.format("cloudFiles") \
.option("cloudFiles.format", "json") \
# The schema location directory keeps track of your data schema over time
.option("cloudFiles.schemaLocation", "<path-to-checkpoint>") \
.option("cloudFiles.inferColumnTypes", "true") \
.load("/Volumes/catalog_name/schema_name/volume_name/nested_json")

Charger les fichiers CSV sans en-têtes

L'exemple suivant montre comment charger des fichiers CSV sans en-têtes à l'aide d'Auto Loader. Utilisez rescuedDataColumn pour capturer toutes les données qui ne correspondent pas au schéma fourni.

Python
df = spark.readStream.format("cloudFiles") \
.option("cloudFiles.format", "csv") \
.option("rescuedDataColumn", "_rescued_data") \ # ensure that you don't lose data
.schema(<schema>) \ # provide a schema here for the files
.load(<path>)

Appliquer un schéma aux fichiers CSV avec en-têtes

L'exemple suivant montre comment appliquer un schéma aux fichiers CSV qui incluent des en-têtes. Utilisez rescuedDataColumn pour capturer toutes les données qui ne correspondent pas au schéma fourni.

Python
df = spark.readStream.format("cloudFiles") \
.option("cloudFiles.format", "csv") \
.option("header", "true") \
.option("rescuedDataColumn", "_rescued_data") \ # makes sure that you don't lose data
.schema(<schema>) \ # provide a schema here for the files
.load(<path>)

Ingérer des données d'image ou binaires dans Delta Lake pour le ML

Une fois les données stockées dans Delta Lake, vous pouvez exécuter une inférence distribuée sur les données. Consultez Effectuer une inférence distribuée à l’aide de pandas UDF.

Python
spark.readStream.format("cloudFiles") \
.option("cloudFiles.format", "binaryFile") \
.load("/Volumes/catalog_name/schema_name/volume_name/images") \
.writeStream \
.option("checkpointLocation", "<path-to-checkpoint>") \
.start("<path_to_target>")

Syntaxe d'Auto Loader pour les LakeFlow Pipelines

LakeFlow Pipelines provide slightly modified Python syntax for Auto Loader and add SQL support for Auto Loader. Les exemples suivants utilisent Auto Loader pour créer des dataset à partir de fichiers JSON à l'aide du dataset d'exemple de réservation de voyage Wanderbricks :

Python
@dp.table
def booking_updates():
return (
spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("multiLine", "true")
.load("/Volumes/my_catalog/my_schema/my_volume/wanderbricks/booking_updates")
)

@dp.table
def reviews():
return (
spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("multiLine", "true")
.load("/Volumes/my_catalog/my_schema/my_volume/wanderbricks/reviews")
)

Vous pouvez utiliser les options de format prises en charge pour Auto Loader. Les options pour read_files sont des paires clé-valeur. Pour plus de détails sur les formats et options pris en charge, consultez Options.

SQL
CREATE OR REFRESH STREAMING TABLE my_table
AS SELECT *
FROM STREAM read_files(
"/Volumes/my_volume/path/to/files/*",
option-key => option-value,
...
)

L’exemple suivant lit des fichiers JSON multilignes avec l’inférence du type de colonne activée :

SQL
CREATE OR REFRESH STREAMING TABLE booking_updates
AS SELECT * FROM STREAM read_files(
"/Volumes/my_catalog/my_schema/my_volume/wanderbricks/booking_updates",
format => "json",
multiLine => true,
inferColumnTypes => true
)

Vous pouvez utiliser le schema pour spécifier le format manuellement ; vous devez spécifier le schema pour les formats qui ne prennent pas en charge l'inférence de schéma:

Python
@dp.table
def booking_updates_raw():
return (
spark.readStream.format("cloudFiles")
.schema("booking_id LONG, booking_update_id LONG, user_id LONG, property_id LONG, status STRING, guests_count INT, total_amount DOUBLE, check_in DATE, check_out DATE, created_at TIMESTAMP, updated_at TIMESTAMP")
.option("cloudFiles.format", "json")
.option("multiLine", "true")
.load("/Volumes/my_catalog/my_schema/my_volume/wanderbricks/booking_updates")
)
remarque

Les LakeFlow Pipelines configurent et gèrent automatiquement les répertoires de schéma et de point de contrôle lors de l’utilisation d’Auto Loader pour lire les fichiers. Toutefois, si vous configurez manuellement l’un de ces répertoires, l’exécution d’un refresh complet n’affecte pas le contenu des répertoires configurés. Databricks recommande d’utiliser les répertoires configurés automatiquement afin d’éviter les effets secondaires inattendus pendant le traitement.

Ressources supplémentaires