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 |
| Correspond à un seul caractère du jeu de caractères {a,b,c}. |
| Correspond à un seul caractère de la plage de caractères {a…z}. |
| Correspond à un seul caractère qui n'appartient pas à un jeu de caractères ou à une plage {a}. Notez que le caractère |
| Correspond à une chaîne de l'ensemble de chaînes {ab, cd}. |
| 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
- Scala
df = spark.readStream.format("cloudFiles") \
.option("cloudFiles.format", <format>) \
.schema(schema) \
.load("/Volumes/catalog_name/schema_name/volume_name/*/files")
val 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
- Scala
df = spark.readStream.format("cloudFiles") \
.option("cloudFiles.format", "binaryFile") \
.option("pathGlobfilter", "*.png") \
.load("/Volumes/catalog_name/schema_name/volume_name/path")
val df = spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "binaryFile")
.option("pathGlobfilter", "*.png")
.load("/Volumes/catalog_name/schema_name/volume_name/path")
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 |
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
- Scala
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>")
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
- Scala
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>")
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 :
.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
- Scala
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>")
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
- Scala
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":...}}
)
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
- Scala
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")
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
- Scala
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>)
val df = spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "csv")
.option("rescuedDataColumn", "_rescued_data") // makes sure 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
- Scala
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>)
val 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
- Scala
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>")
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
- SQL
@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")
)
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
)
CREATE OR REFRESH STREAMING TABLE reviews
AS SELECT * FROM STREAM read_files(
"/Volumes/my_catalog/my_schema/my_volume/wanderbricks/reviews",
format => "json",
multiLine => true
)
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.
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 :
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
- SQL
@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")
)
CREATE OR REFRESH STREAMING TABLE booking_updates_raw
AS SELECT *
FROM STREAM read_files(
"/Volumes/my_catalog/my_schema/my_volume/wanderbricks/booking_updates",
format => "json",
multiLine => true,
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"
)
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.