Aller au contenu principal

Bonnes pratiques d'Auto Loader

Cette page décrit les bonnes pratiques que vous pouvez appliquer pour configurer Auto Loader afin qu'il fonctionne de manière fiable, rentable et pour monter en charge pour votre cas d'utilisation.

Ces bonnes pratiques réduisent la charge opérationnelle et préviennent les problèmes courants difficiles à diagnostiquer en production tels que : les coûts API LIST inutiles liés aux analyses complètes de répertoires, la perte silencieuse de données due à la drift de schéma, et les redémarrages de pipelines causés par une mauvaise configuration des points de contrôle.

Pour plus de détails sur la configuration de production, voir Configurer Auto Loader pour les charges de travail de production. Pour le monitoring et l'observabilité, consultez Surveiller et observer Auto Loader.

Choisir le bon framework d'exécution

Le meilleur framework d'exécution pour votre cas d'utilisation dépend du niveau de contrôle dont vous avez besoin sur le pipeline et de la surcharge opérationnelle que vous souhaitez gérer. Pour la plupart des utilisateurs et des pipelines de production, Auto Loader avec les LakeFlow Pipelines convient bien. Cependant, si vous avez besoin d'un contrôle et d'une personnalisation maximaux, utilisez Auto Loader avec Structured Streaming. Pour la configuration la plus simple avec une expérience gérée, utilisez un connecteur LakeFlow géré lorsqu'il est disponible.

Les LakeFlow pipelines étendent Structured Streaming avec l'autoscaling, les contrôles de qualité des données, la gestion de l'évolution des schémas et le monitoring via le journal des événements Log. Databricks recommande les LakeFlow Pipelines pour la plupart des charges de travail d'ingestion en production.

Choisissez le bon type de planification et de Trigger

Le meilleur type de planification et de Trigger pour votre cas d'utilisation dépend de vos exigences en matière de latence et de vos modèles d'arrivée de fichiers. Pour la plupart des cas d'utilisation, Databricks recommande un Trigger d'arrivée de fichier avec les événements de fichier activés. Ceci permet une ingestion à faible latence et à faible coût, car le compute ne s'exécute que lorsque de nouveaux fichiers arrivent. Les trois types de Trigger diffèrent quant au moment et à la fréquence de start du pipeline :

  • Continu : le pipeline s'exécute sans s'arrêter. À n’utiliser que lorsqu’une latence inférieure à la seconde est une exigence stricte, car le compute continu coûte plus cher. Associer aux événements de fichier.
  • Trigger d'arrivée de fichiers : le pipeline start lorsque de nouveaux fichiers arrivent dans l'emplacement source. Idéal pour une latence faible à moyenne ou des modèles d'arrivée de fichiers irréguliers. Nécessite que les événements de fichier soient activés. Voir Trigger des Job lorsque de nouveaux fichiers arrivent.
  • Planifié : le pipeline s’exécute selon un calendrier temporel (par exemple, toutes les heures). À utiliser lorsque les exigences de latence sont souples (de quelques minutes à quelques heures). Compatible avec la liste des répertoires, mais les événements de fichier réduisent les coûts même en mode planifié en évitant les analyses complètes des répertoires.

Pour plus de détails sur l'utilisation de Trigger.AvailableNow pour la planification des batchs, consultez Utilisation de Trigger.AvailableNow et la limitation du débit.

Choisissez le bon mode de découverte de fichiers

Auto Loader prend en charge trois modes de découverte de fichiers avec différents compromis en termes de complexité de configuration, d'évolutivité et de coût.

Mode

Complexité de la configuration

Évolutivité

Coût

Quand utiliser

Événements de fichier (recommandé)

Faible (configuration des autorisations unique)

Millions de fichiers par heure

Le plus bas

default pour la plupart des charges de travail

Notification de fichier classique

Élevé (21+ options de configuration cloud)

Millions de fichiers par heure

Medium

Lorsque les événements de fichier ne sont pas disponibles

Liste de répertoires

Aucun

Limité par la taille du répertoire

Le plus élevé (coûts de l'API LIST)

Petits répertoires, backfills ponctuels ou lorsque les politiques de sécurité empêchent les événements de fichiers

Mode

Complexité de la configuration

Évolutivité

Coût

Quand utiliser

Événements de fichier (recommandé)

Faible (configuration des autorisations unique)

Millions de fichiers par heure

Le plus bas

default pour la plupart des charges de travail

Notification de fichier classique

Élevé (21+ options de configuration cloud)

Millions de fichiers par heure

Medium

Lorsque les événements de fichier ne sont pas disponibles

Liste de répertoires

Aucun

Limité par la taille du répertoire

Le plus élevé (coûts de l'API LIST)

Petits répertoires, backfills ponctuels ou lorsque les politiques de sécurité empêchent les événements de fichiers

Les événements de fichiers consolident les ressources de stockage cloud en utilisant un abonnement et une file d'attente par emplacement externe au lieu d'un par Stream. La différence de performance est significative à l'échelle : la liste des répertoires doit analyser l'intégralité du répertoire source à chaque trigger, ainsi, le temps d'ingestion augmente avec la taille du répertoire. Les événements de fichiers fournissent directement de nouvelles notifications de fichiers, ainsi, le temps d'ingestion reste faible, quel que soit le nombre d'objets dans le répertoire.

Activer les événements de fichier

Les événements de fichier nécessitent une autorisation de cloud unique et un emplacement externe configuré pour utiliser le service d'événements de fichier géré. Une fois configurés, tous les flux Auto Loader lisant depuis cet emplacement externe peuvent utiliser les événements de fichiers sans configuration supplémentaire.

  1. Accordez les autorisations cloud requises côté fournisseur de cloud. Les exigences varient selon le fournisseur de cloud. Consultez Configurer les événements de fichiers pour un emplacement externe.

  2. Définissez cloudFiles.useManagedFileEvents sur true dans votre query Auto Loader.

    Python
    df = (spark.readStream
    .format("cloudFiles")
    .option("cloudFiles.format", "json")
    .option("cloudFiles.useManagedFileEvents", "true")
    .load("/path/to/data/dir"))

    Pour connaître toutes les étapes de configuration, consultez Migrer vers Auto Loader avec les événements de fichier.

Lorsque vous ne pouvez pas utiliser les événements de fichier

Vous ne pourrez peut-être pas utiliser les événements de fichier dans les cas suivants :

  • L'emplacement externe n'est pas configuré avec des événements de fichier.
  • Les politiques de sécurité de l’organisation ne permettent pas d’activer les événements de fichiers sur un emplacement externe partagé.

Dans ces cas, utilisez le mode de notification de fichiers classique ou le mode de liste de répertoires. Pour une comparaison complète des modes de détection de fichiers, consultez Comparer les modes de détection de fichiers Auto Loader.

Gérer l'évolution des schémas

Auto Loader infère automatiquement le schéma, mais la façon dont vous configurez l'évolution des schémas affecte l'exhaustivité des données et la stabilité du pipeline. Utilisez la table suivante pour choisir une stratégie.

Scénario

Recommandation

Le schéma est connu et fixe.

Fournissez un schéma explicite avec .schema()

Le schéma est inconnu, des modifications additives sont attendues.

schemaEvolutionMode: addNewColumns

Le schéma est inconnu, des modifications de type sont attendues

schemaEvolutionMode: addNewColumnsWithTypeWidening

Contrat de schéma strict requis

schemaEvolutionMode: failOnNewColumns

Schéma arbitraire ou imprévisible

Ingérer en tant que type Variant

Scénario

Recommandation

Le schéma est connu et fixe.

Fournissez un schéma explicite avec .schema()

Le schéma est inconnu, des modifications additives sont attendues.

schemaEvolutionMode: addNewColumns

Le schéma est inconnu, des modifications de type sont attendues

schemaEvolutionMode: addNewColumnsWithTypeWidening

Contrat de schéma strict requis

schemaEvolutionMode: failOnNewColumns

Schéma arbitraire ou imprévisible

Ingérer en tant que type Variant

Après avoir choisi une stratégie, appliquez les pratiques suivantes pour affiner le comportement de l'évolution des schémas.

Utilisez des indications de schéma pour les types de champs connus

Utilisez l'option cloudFiles.schemaHints pour appliquer des types aux champs que vous connaissez à l'avance, tout en permettant l'inférence de schéma pour d'autres champs.

Python
df = (spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.schemaHints", "id long, amount double")
.load("/path/to/data/dir"))

Utiliser l'élargissement de type pour les changements de type compatibles

Le mode d'évolution des schémas addNewColumnsWithTypeWidening élargit automatiquement les types compatibles (par exemple, de int à long) au lieu de router les données vers la colonne _rescued_data. Cela évite d'avoir recours à des Jobs de post-traitement pour gérer les promotions de type simples.

Python
df = (spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "parquet")
.option("cloudFiles.schemaEvolutionMode", "addNewColumnsWithTypeWidening")
.load("/path/to/data/dir"))

Ingérer en tant que type Variant pour les schémas imprévisibles

Lorsque vos données ne sont pas conformes à un schéma spécifique ou que le schéma change en permanence, ingérez les données en tant que type Variant.

Python
df = (spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("singleVariantColumn", "data")
.load("/path/to/data/dir"))

Variant fournit un schéma à la lecture au moment de la query mais est moins efficace que l'interrogation de colonnes structurées. Pour connaître tous les mécanismes de l'inférence et de l'évolution de schéma, consultez Configurer l'inférence et l'évolution de schéma dans Auto Loader.

Gérer les données de mauvaise qualité et la qualité des données

Les pratiques suivantes vous aident à détecter, capturer et isoler les mauvaises données avant qu'elles ne se propagent aux couches en aval.

Activez _rescued_data et _corrupt_record

Auto Loader fournit deux colonnes pour capturer les données qui ne peuvent pas être analysées correctement.

  • _rescued_data capture les champs qui ne correspondent pas au schéma actuel. Il est ajouté automatiquement par Auto Loader.
  • _corrupt_record capture les lignes qui ne peuvent pas être analysées du tout. Activez-le en utilisant columnNameOfCorruptRecord:
Python
df = (spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.schemaHints", "_corrupt_record string")
.option("columnNameOfCorruptRecord", "_corrupt_record")
.load("/path/to/data/dir"))

Databricks recommande columnNameOfCorruptRecord plutôt que badRecordsPath pour éviter les conditions de concurrence potentielles qui peuvent manquer des enregistrements corrompus.

Utiliser les attentes des Lakeflow pipelines pour la surveillance

Définissez les attentes des Lakeflow Pipelines pour vérifier que _rescued_data et _corrupt_record sont NULL dans des conditions normales. Les valeurs non-NULL signalent un drift de schéma ou une corruption de données.

Python
import dlt

@dlt.table
@dlt.expect("no rescued data", "_rescued_data IS NULL")
@dlt.expect("no corrupt records", "_corrupt_record IS NULL")
def bronze_table():
return (spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.schemaHints", "_corrupt_record string")
.option("columnNameOfCorruptRecord", "_corrupt_record")
.load("/path/to/data/dir"))

Isolez les données corrompues

Isolez les lignes contenant des données non analysables dans un récepteur dédié à des fins d'enquête. Ceci empêche les données corrompues de se propager vers les couches en aval.

Python
import dlt

@dlt.table
def corrupt_records_sink():
return dlt.read_stream("bronze_table").where("_corrupt_record IS NOT NULL")

@dlt.view
def clean_table():
return dlt.read_stream("bronze_table").where("_corrupt_record IS NULL")

Annoter les données avec les métadonnées du fichier source

Incluez la colonne _metadata dans vos requêtes d'ingestion Auto Loader. Au minimum, capturez file_path et file_modification_time. Cela vous permet de retracer les problèmes de données jusqu'aux fichiers sources spécifiques et de les joindre à cloud_files_state() pour le cycle de vie complet du fichier.

Python
df = (spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.load("/path/to/data/dir")
.select("*", "_metadata.file_path", "_metadata.file_modification_time"))

Pour plus de détails, consultez colonne de métadonnées de fichier.

Optimiser les coûts et les performances

Les pratiques suivantes réduisent les trois principaux facteurs de coût pour Auto Loader : les appels LIST API cloud, la capacité de compute inactive et la croissance du stockage à long terme.

  • Utilisez les événements de fichier pour minimiser les coûts de l'API LIST : les événements de fichier offrent une découverte incrémentielle des fichiers, éliminant le besoin de listes complètes de répertoires à chaque exécution. C'est l'optimisation des coûts la plus impactante pour Auto Loader.

  • Utilisez les triggers d'arrivée de fichiers pour le traitement événementiel : les triggers d'arrivée de fichiers start votre pipeline uniquement lorsque de nouveaux fichiers arrivent, vous ne payez donc pas pour le compute inactif. Voir Trigger des Job lorsque de nouveaux fichiers arrivent.

  • Archiver les fichiers traités avec cloudFiles.cleanSource : Utilisez cloudFiles.cleanSource pour supprimer ou déplacer automatiquement les fichiers traités. Cela réduit les coûts de stockage et les coûts de listage de répertoires pour les flux à longue durée de vie. Pour plus de détails, voir Archivage des fichiers dans le répertoire source pour réduire les coûts.

    • Utilisez le mode delete pour supprimer les fichiers après l’ingestion.
    • Utilisez le mode move pour archiver les fichiers vers un emplacement différent à des fins de conformité ou d'audit.
    Python
    df = (spark.readStream
    .format("cloudFiles")
    .option("cloudFiles.format", "json")
    .option("cloudFiles.cleanSource", "delete")
    .load("/path/to/data/dir"))
attention

Ne pas activer cloudFiles.cleanSource si plusieurs streams Auto Loader ou d'autres clients lisent depuis le même répertoire source.

  • Tirez parti des améliorations de performances : Mettez à niveau vers la dernière version de Databricks Runtime ou utilisez le compute serverless pour bénéficier des récentes améliorations de performances d'Auto Loader.

Gestion des points de contrôle

Le point de contrôle stocke la progression et l'état du fichier du Stream. Mal configurer ou perdre le point de contrôle nécessite un redémarrage complet, traitez-le donc comme une infrastructure critique.

  • N'appliquez jamais de stratégies de cycle de vie des objets cloud aux emplacements de point de contrôle. Si les fichiers de point de contrôle sont supprimés, l'état du Stream est corrompu et vous devez redémarrer de zéro.
  • Utilisez des points de contrôle séparés pour chaque Stream et répertoire source.
  • Envisagez cloudFiles.maxFileAge pour les Stream à grand volume et de longue durée afin de limiter la croissance de l'état. Utilisez un paramètre conservateur (90 jours minimum recommandés). Définir cette valeur de manière trop agressive risque de retraiter des fichiers qu'Auto Loader a déjà ingérés s'ils tombent en dehors de la fenêtre.

Pour tous les détails, consultez le suivi des événements de fichiers.

Utiliser les volumes pour une découverte optimale des fichiers avec les événements de fichier

Pour des performances améliorées avec les événements de fichier, créez un volume externe pour chaque chemin ou sous-répertoire à partir duquel Auto Loader charge les données. Fournissez les chemins de volume (par exemple, /Volumes/catalog/schema/volume) à Auto Loader au lieu des chemins cloud (par exemple, s3://bucket/path). Ceci optimise la découverte de fichiers grâce à un modèle d'accès aux données optimisé.

Pour connaître les meilleures pratiques concernant les événements de fichier, consultez Bonnes pratiques pour Auto Loader avec les événements de fichier.