Meilleures pratiques relatives au chargeur automatique

Cette page décrit les meilleures pratiques que vous pouvez appliquer pour configurer le chargeur automatique pour qu’il s’exécute de manière fiable, rentable et à grande échelle pour votre cas d’usage.

Ces meilleures pratiques réduisent la surcharge opérationnelle et empêchent les problèmes courants qui sont difficiles à diagnostiquer en production, tels que les coûts d’API inutiles LIST des analyses complètes des répertoires, la perte de données silencieuse de la dérive de schéma et les redémarrages de pipeline causés par une configuration incorrecte du point de contrôle.

Pour plus d’informations sur la configuration de production, consultez Configurer le chargeur automatique pour les charges de travail de production. Pour la surveillance et l’observabilité, consultez Surveiller et observer le chargeur automatique.

Choisir l’infrastructure d’exécution appropriée

La meilleure infrastructure d’exécution pour votre cas d’usage dépend de la quantité de contrôle dont vous avez besoin sur le pipeline et de la charge opérationnelle que vous souhaitez gérer. Pour la plupart des utilisateurs et des pipelines de production, le chargeur automatique avec les pipelines Lakeflow est adapté. Toutefois, si vous avez besoin d’un contrôle et d’une personnalisation maximums, utilisez le chargeur automatique avec Structured Streaming. Pour la configuration la plus simple avec une expérience managée, utilisez un connecteur LakeFlow managé lorsqu’il est disponible.

Les pipelines Lakeflow étendent Structured Streaming avec la mise à l’échelle automatique, les vérifications de qualité des données, la gestion de l’évolution du schéma et la surveillance via le journal des événements. Databricks recommande les pipelines Lakeflow pour la plupart des charges de travail d’ingestion en production.

Choisir le type de planification et de déclencheur approprié

Le meilleur type de planification et de déclencheur pour votre cas d’usage dépend des exigences de latence et des modèles d’arrivée des fichiers. Pour la plupart des cas d’usage, Databricks recommande un déclencheur d’arrivée de fichier avec les événements de fichier activés. Cela permet d’obtenir une ingestion à faible latence à faible coût, car le calcul s’exécute uniquement lorsque de nouveaux fichiers arrivent. Les trois types de déclencheurs diffèrent selon le moment et la fréquence à laquelle le pipeline démarre :

  • Continu : Le pipeline s’exécute sans s’arrêter. Utilisez uniquement lorsque la latence de sous-seconde est une exigence difficile, car les coûts de calcul continus sont plus élevés. Associez à des événements de fichiers.
  • Déclencheur d’arrivée de fichier : le pipeline démarre lorsque de nouveaux fichiers atterrissent à l’emplacement source. Idéal pour les modèles d’arrivée de fichiers irréguliers ou à latence faible à moyenne. Nécessite l’activation des événements de fichier. Voir Déclencher des tâches lorsque de nouveaux fichiers arrivent.
  • Planifié : le pipeline s’exécute selon une planification basée sur le temps (par exemple, toutes les heures). À utiliser lorsque les exigences de latence sont peu strictes (de quelques minutes à quelques heures). Fonctionne 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 d’informations sur l’utilisation Trigger.AvailableNow de la planification par lots, consultez Utilisation de Trigger.AvailableNow et limitation de débit.

Choisir le mode de découverte de fichiers approprié

Le chargeur automatique prend en charge trois modes de découverte de fichiers avec différents compromis dans la complexité de l’installation, l’extensibilité et le coût.

Mode Complexité de l’installation Scalability Coûts Quand utiliser
Événements de fichier (recommandé) Faible (configuration d’autorisation ponctuelle) Millions de fichiers par heure Le plus bas Valeur par défaut pour la plupart des charges de travail
Notification de fichier classique Élevé (21 options de configuration du cloud ou plus) Millions de fichiers par heure Moyenne Quand les événements de fichier ne sont pas disponibles
Liste de répertoires None Limité par la taille du répertoire Plus élevé (LIST coûts d’API) Petits répertoires, recharges rétroactives ponctuelles ou lorsque les stratégies de sécurité empêchent les événements de fichier

Les événements de fichier consolident les ressources de stockage cloud à l’aide d’un abonnement et d’une file d’attente par emplacement externe au lieu d’un par flux. La différence de performances est significative à grande échelle : la liste des répertoires doit analyser l’intégralité du répertoire source sur chaque déclencheur, de sorte que le temps d’ingestion augmente avec la taille du répertoire. Les événements de fichier fournissent directement de nouvelles notifications de fichier, de sorte que 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 octroi d’autorisations cloud unique et un emplacement externe configuré pour utiliser le service d’événements de fichiers managés. Une fois configurés, tous les flux de chargeur automatique lisant à partir de cet emplacement externe peuvent utiliser des événements de fichier 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 des événements de fichier pour un emplacement externe.

  2. Définissez cloudFiles.useManagedFileEvents sur true dans votre requête Auto Loader.

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

    Pour connaître les étapes d’installation complètes, consultez Migrer vers le chargeur automatique avec des événements de fichier.

Lorsque vous ne pouvez pas utiliser d’événements de fichier

Vous ne pouvez peut-être pas utiliser les événements de fichier lorsque :

  • L’emplacement externe n’est pas configuré avec les événements de fichier.
  • Les stratégies de sécurité de l’organisation n’autorisent pas l’activation des événements de fichier sur un emplacement externe partagé.

Dans ces cas, utilisez le mode de notification de fichier classique ou le mode de liste d’annuaires. Pour obtenir une comparaison complète des modes de détection de fichiers, consultez Comparer les modes de détection de fichiers du chargeur automatique.

Gérer l’évolution du schéma

Le chargeur automatique déduit automatiquement le schéma, mais la façon dont vous configurez l’évolution du schéma affecte l’exhaustivité des données et la stabilité du pipeline. Utilisez le tableau suivant pour choisir une stratégie.

Scénario Recommendation
Le schéma est connu et corrigé Fournir un schéma explicite avec .schema()
Le schéma est inconnu, les modifications additives attendues schemaEvolutionMode : addNewColumns
Le schéma est inconnu, des changements de type sont à prévoir schemaEvolutionMode : addNewColumnsWithTypeWidening
Contrat de schéma strict requis schemaEvolutionMode : failOnNewColumns
Schéma arbitraire ou imprévisible Ingérer comme type Variant

Une fois que vous avez choisi une stratégie, appliquez les pratiques suivantes pour affiner le comportement de l’évolution du schéma.

Utiliser des indicateurs de schéma pour les types de champs connus

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

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 des modifications de type compatibles

Le addNewColumnsWithTypeWidening mode d’évolution du schéma élargit automatiquement les types compatibles (par exemple, int vers long) au lieu de router les données vers la _rescued_data colonne. Cela évite d’avoir besoin de travaux de post-traitement pour gérer les promotions de type simple.

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

Importer comme 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 continu, ingérer les données en tant que Variant type.

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

Variant fournit un schéma en lecture au moment de la requête, mais est moins efficace que l’interrogation de colonnes structurées. Pour connaître la mécanique complète de l’inférence et de l’évolution du schéma, consultez Configurer l’inférence de schéma et l’évolution dans le chargeur automatique.

Gérer les données incorrectes et la qualité des données

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

Activer _rescued_data et _corrupt_record

Auto Loader fournit deux colonnes pour recueillir 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 le chargeur automatique.
  • _corrupt_record capture les lignes qui ne peuvent pas être analysées du tout. Activez-le à l’aide de columnNameOfCorruptRecord:
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 afin d’éviter d’éventuelles situations de concurrence pouvant empêcher la détection d’enregistrements corrompus.

Utiliser les attentes des pipelines Lakeflow pour la surveillance

Définissez des attentes dans les pipelines Lakeflow pour vérifier que _rescued_data et _corrupt_record sont NULL dans des conditions normales. Les valeurs non NULL signalent la dérive du schéma ou l’altération des données.

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

Isoler les données endommagées

Isolez les lignes contenant des données non analysables dans un récepteur dédié à des fins d’investigation. Cela empêche la propagation de données endommagées vers des couches en aval.

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 des données avec des métadonnées de 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 remonter les problèmes de données jusqu’à des fichiers sources spécifiques et d’effectuer une jointure avec cloud_files_state() sur l’ensemble du cycle de vie des fichiers.

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

Pour plus d’informations, consultez la colonne 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 le chargeur automatique : les appels d’API cloud LIST , le calcul inactif et la croissance du stockage à long terme.

  • Utiliser des événements de fichier pour réduire LIST Coûts de l’API : les événements de fichier fournissent une découverte de fichiers incrémentielle, ce qui élimine la nécessité de répertorier des répertoires complets sur chaque exécution. C’est l’optimisation des coûts la plus importante pour Auto Loader.

  • Utilisez des déclencheurs d’arrivée de fichiers pour le traitement piloté par les événements : les déclencheurs d’arrivée de fichier démarrent votre pipeline uniquement lorsque de nouveaux fichiers arrivent, de sorte que vous ne payez pas pour le calcul inactif. Voir Déclencher des tâches lorsque de nouveaux fichiers arrivent.

  • Archivez les fichiers traités avec cloudFiles.cleanSource : permet cloudFiles.cleanSource de supprimer ou de déplacer automatiquement des fichiers traités. Cela réduit les coûts de stockage et les coûts de référencement des annuaires pour les flux de longue durée. Pour plus d’informations, consultez Les fichiers d’archivage dans le répertoire source pour réduire les coûts.

    • Utilisez delete le mode pour supprimer des fichiers après l’ingestion.
    • Utilisez move le mode pour archiver des fichiers à un autre emplacement pour la conformité ou l’audit.
    df = (spark.readStream
      .format("cloudFiles")
      .option("cloudFiles.format", "json")
      .option("cloudFiles.cleanSource", "delete")
      .load("/path/to/data/dir"))
    

    Warning

    N’activez cloudFiles.cleanSource pas si plusieurs flux de chargeur automatique ou d’autres clients lisent à partir du même répertoire source.

  • Tirez parti des améliorations des performances : effectuez une mise à niveau vers la dernière version du runtime Databricks ou utilisez le calcul serverless pour tirer parti des améliorations récentes des performances du chargeur automatique.

Gestion des points de contrôle

Le point de contrôle stocke la progression du flux et l’état du fichier. La configuration incorrecte ou la perte du point de contrôle nécessite un redémarrage complet, de sorte qu’il s’agit d’une infrastructure critique.

  • N’appliquez jamais de stratégies de cycle de vie d’objet cloud aux emplacements de point de contrôle. Si les fichiers de point de contrôle sont supprimés, l’état du flux est endommagé et vous devez redémarrer à partir de zéro.
  • Utilisez des points de contrôle distincts pour chaque flux et répertoire source.
  • Envisagez d’utiliser cloudFiles.maxFileAge pour les flux de longue durée et à fort volume afin de limiter la croissance de l’état. Utilisez un paramètre conservateur (90 jours minimum recommandé). Définir cette valeur de manière trop agressive risque d’entraîner le retraitement de fichiers qu’Auto Loader a déjà ingérés s’ils tombent hors de la fenêtre.

Pour plus d’informations, consultez Le suivi des événements de fichier.

Utiliser des volumes pour une détection optimale des fichiers à l’aide des événements de fichier

Pour améliorer les performances avec les événements de fichier, créez un volume externe pour chaque chemin d’accès ou sous-répertoire à partir duquel le chargeur automatique est chargé. Fournissez à Auto Loader des chemins de volume (par exemple, /Volumes/catalog/schema/volume) au lieu de chemins de stockage cloud (par exemple, s3://bucket/path). Cela optimise la découverte de fichiers par le biais d’un modèle d’accès aux données optimisé.

Pour plus d’informations sur les meilleures pratiques relatives aux événements de fichier, consultez Les meilleures pratiques pour le chargeur automatique avec les événements de fichier.