Configurez Auto Loader pour les charges de travail de production
Pour connaître les bonnes pratiques complètes de configuration d’Auto Loader, notamment la sélection du mode de découverte de fichiers, la gestion du schéma et le traitement de la qualité des données, consultez les bonnes pratiques d’Auto Loader.
Databricks recommande d'utiliser Auto Loader dans les LakeFlow Pipelines pour l'ingestion incrémentielle des données. Les LakeFlow Pipelines étendent les fonctionnalités d'Apache Spark Structured Streaming et vous permettent d'écrire seulement quelques lignes de Python ou de SQL déclaratif pour déployer une pipeline de données de qualité production avec :
- Infrastructure de compute à dimensionnement automatique pour des économies de coûts : Optimisez l'utilisation des clusters LakeFlow Pipelines avec le dimensionnement automatique
- Vérifications de la qualité des données avec des attentes : Gérer la qualité des données avec les attentes du pipeline
- Gestion automatique de l'évolution des schémas : configurez l'inférence et l'évolution des schémas dans Auto Loader
- Monitoring via des métriques dans le log des événements : Log des événements du pipeline
Databricks vous recommande également de suivre les bonnes pratiques de streaming pour exécuter Auto Loader en production. Consultez Considérations relatives à la production pour Structured Streaming.
Les Lakeflow Pipelines sont le moyen recommandé pour exécuter Auto Loader pour la plupart des ingestions en production. Si votre charge de travail n'a pas d'exigences de faible latence et que votre priorité est de minimiser le coût de calcul, vous pouvez à la place planifier Auto Loader comme un job de batch déclenché qui utilise Trigger.AvailableNow. Consultez Considérations sur les coûts.
Monitoring Auto Loader
Les sections suivantes décrivent comment surveiller Auto Loader en production, y compris les métriques, les Logs, les alertes et les workflows de dépannage courants. Pour une référence complète couvrant les modèles de tableau de bord, l'analyse de la latence et la détection de la drift de schéma, consultez Surveiller et observer Auto Loader.
Interrogation des fichiers découverts par Auto Loader
Auto Loader fournit une API SQL pour inspecter l'état d'un Stream. En utilisant la fonction cloud_files_state, vous pouvez trouver des métadonnées sur les fichiers qui ont été découverts par un Stream Auto Loader. Query cloud_files_state, en fournissant l’emplacement du point de contrôle associé à un Stream Auto Loader.
La fonction cloud_files_state est disponible dans Databricks Runtime 11.3 LTS et versions ultérieures.
SELECT * FROM cloud_files_state('path/to/checkpoint');
Écouter les mises à jour de Stream
Pour surveiller davantage les flux Auto Loader, Databricks recommande d'utiliser l'interface Streaming Query Listener d'Apache Spark. Consultez monitoring Structured Streaming queries sur Databricks.
Auto Loader signale des métriques au Streaming Query Listener à chaque batch. Vous pouvez voir le nombre de fichiers dans le backlog et la taille du backlog dans les métriques numFilesOutstanding et numBytesOutstanding sous l'onglet Données brutes du tableau de bord de progression de la query de streaming :
{
"sources": [
{
"description": "CloudFilesSource[/path/to/source]",
"metrics": {
"numFilesOutstanding": "238",
"numBytesOutstanding": "163939124006"
}
}
]
}
Lorsque vous utilisez le mode de notification de fichiers dans Databricks Runtime 10.4 LTS et versions ultérieures, les métriques incluent également le nombre approximatif d'événements de fichiers dans la file d'attente du cloud comme approximateQueueSize pour AWS et Azure.
Considérations sur les coûts
Lorsque vous exécutez Auto Loader, vos principales sources de coût sont les ressources de compute et la découverte de fichiers.
Si votre charge de travail n’a pas d’exigences de faible latence, vous pouvez réduire les coûts de compute en utilisant Lakeflow Jobs pour planifier Auto Loader en tant que batch jobs à l’aide de Trigger.AvailableNow au lieu de l’exécuter en continu. Consultez Configurer les intervalles de Trigger Structured Streaming. Ces batch jobs peuvent être déclenchés à l’aide de déclencheurs d’arrivée de fichiers afin de réduire davantage la latence entre l’arrivée des fichiers et le traitement.
Les coûts de découverte de fichiers peuvent prendre la forme de LIST opérations sur vos comptes de stockage en mode de listage de répertoires et de requêtes d'API sur le service d'abonnement et le service de file d'attente en mode de notification de fichiers. Les Trigger continus tels que Trigger.ProcessingTime sont particulièrement coûteux en mode de listage de répertoires, car Auto Loader liste en continu l'ensemble du répertoire pour trouver de nouveaux fichiers. Si votre charge de travail nécessite des triggers continus, Databricks recommande de choisir un mode de découverte de fichiers basé sur vos exigences de latence :
- Faible latence et simplicité : utilisez Auto Loader avec les événements de fichiers. Les événements de fichier ne nécessitent qu'une seule file d'attente par compartiment et utilisent une découverte incrémentielle lors des exécutions ultérieures. Pour plus d'informations, voir présentation d'Auto Loader avec les événements de fichiers.
- Applications très sensibles à la latence : Utilisez le mode de notification de fichiers classique. Le mode classique lit directement à partir de la file d'attente cloud sans le saut de mise en cache supplémentaire introduit par les événements de fichier. Dans ce mode, vous pouvez étiqueter les ressources créées par Auto Loader afin de suivre vos coûts à l'aide d'étiquettes de ressources. Pour plus de détails, voir Notification de fichier.
Rétention des données sources
Disponible dans Databricks Runtime 16.4 LTS et versions ultérieures.
Lorsque les fichiers s'accumulent dans votre répertoire source, les coûts de stockage augmentent et la découverte des fichiers ralentit, en particulier en mode de listage de répertoire. Auto Loader fournit l'option cloudFiles.cleanSource pour gérer automatiquement la rétention des fichiers en archivant ou en supprimant les fichiers après leur traitement.
Archivage de fichiers dans le répertoire source pour réduire les coûts
- Le paramètre
cloudFiles.cleanSourcesupprime ou déplace des fichiers dans le répertoire source. - Si vous utilisez
foreachBatchpour votre traitement de données, vos fichiers deviennent des candidats au déplacement ou à la suppression dès que votre opérationforeachBatchse termine avec succès, même si votre opération n'a consommé qu'un sous-ensemble des fichiers du batch.
Databricks recommande d'utiliser Auto Loader avec des événements de fichier pour réduire les coûts de découverte. Cela réduit également les coûts de compute, car la découverte est incrémentielle.
Si vous ne pouvez pas utiliser les événements de fichier et devez utiliser la liste de répertoires pour découvrir les fichiers, vous pouvez utiliser l'option cloudFiles.cleanSource pour archiver ou supprimer automatiquement les fichiers une fois qu'Auto Loader les a traités afin de réduire les coûts de découverte. Comme Auto Loader nettoie les fichiers de votre répertoire source après traitement, moins de fichiers doivent être répertoriés lors de la détection.
Lorsque vous utilisez cloudFiles.cleanSource avec l’option MOVE, tenez compte des exigences suivantes :
- Le répertoire source et le répertoire de destination doivent être situés dans le même bucket ou conteneur. Les déplacements d'un compartiment à l'autre et d'un conteneur à l'autre ne sont pas pris en charge et entraînent une erreur.
- La destination du déplacement peut être un chemin de volume (par exemple,
/Volumes/my_catalog/my_schema/my_volume/archive/). - Si vos répertoires source et de destination se trouvent au même emplacement externe, ils ne doivent pas avoir de répertoires frères contenant un stockage géré (par exemple, un volume géré ou un catalogue). Dans ces cas, Auto Loader est incapable d'obtenir les autorisations nécessaires pour écrire dans le répertoire de destination.
Databricks recommande d'utiliser cette option dans les cas suivants :
- Votre répertoire source accumule un grand nombre de fichiers au fil du temps.
- Vous devez conserver les fichiers traités à des fins de conformité ou d'audit (définissez
cloudFiles.cleanSourcesurMOVE). - Vous souhaitez réduire les coûts de stockage en supprimant les fichiers après ingestion (définissez
cloudFiles.cleanSourcesurDELETE). Lorsque vous utilisez le modeDELETE, Databricks recommande d'activer la gestion des versions sur le bucket afin que les suppressions d'Auto Loader agissent comme des suppressions logiques et soient disponibles en cas de mauvaise configuration. De plus, Databricks recommande de configurer des politiques de cycle de vie du cloud pour purger les versions anciennes et supprimées temporairement après une période de grâce spécifiée (telle que 60 ou 90 jours) en fonction de vos exigences de récupération.
Pour la référence complète sur les options cleanSource et leurs default, consultez cloudFiles.cleanSource.
Déplacement des fichiers traités vers un chemin de stockage froid
L'exemple suivant configure Auto Loader pour déplacer les fichiers traités vers un répertoire d'archive dans le même bucket après 14 jours. Vous pouvez appliquer une politique de cycle de vie cloud sur le chemin d'archivage pour transférer les fichiers vers des niveaux de stockage moins chers (par exemple, AWS S3 Glacier, Azure Cool/Archive ou GCS Coldline/Archive).
- Python
- SQL
# Step 1: Configure Auto Loader to move processed files to an archive path.
checkpoint = "/Volumes/my_catalog/my_schema/my_volume/checkpoints/ingest_stream"
archive_path = "s3://my-bucket/archive/landing/"
df = (spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.cleanSource", "MOVE")
.option("cloudFiles.cleanSource.moveDestination", archive_path)
.option("cloudFiles.cleanSource.retentionDuration", "14 days")
.option("cloudFiles.schemaLocation", checkpoint)
.load("s3://my-bucket/landing/")
)
# Step 2: Write to a Delta table.
(df.writeStream
.option("checkpointLocation", checkpoint)
.trigger(availableNow=True)
.toTable("my_catalog.my_schema.raw_events")
)
# Step 3 (outside Databricks): Set up a cloud lifecycle policy on the
# archive path to transition files to cold storage after a grace period.
# For example, in AWS you can configure an S3 Lifecycle rule to move
# objects under s3://my-bucket/archive/landing/ to S3 Glacier after
# 30 days.
-- Step 1: Configure Auto Loader to move processed files to an archive path
-- using a Lakeflow Declarative Pipeline.
CREATE OR REFRESH STREAMING TABLE raw_events
AS SELECT * FROM STREAM read_files(
's3://my-bucket/landing/',
format => 'json',
cleanSource => 'MOVE',
cleanSourceMoveDestination => 's3://my-bucket/archive/landing/',
cleanSourceRetentionDuration => '14 days'
);
-- Step 2 (outside Databricks): Set up a cloud lifecycle policy on the
-- archive path to transition files to cold storage.
-- For example, in AWS configure an S3 Lifecycle rule to move objects
-- under s3://my-bucket/archive/landing/ to S3 Glacier after 30 days.
Utilisation de Trigger.AvailableNow et limitation de débit
Disponible dans Databricks Runtime 10.4 LTS et versions ultérieures.
Auto Loader peut être programmé pour s'exécuter dans les Lakeflow Jobs comme un job par batch en utilisant Trigger.AvailableNow. Le AvailableNow Trigger indique à Auto Loader de traiter tous les fichiers arrivés **avant** l'heure de début de la query. Les nouveaux fichiers qui arrivent après le start du Stream sont ignorés jusqu'au prochain Trigger.
Avec Trigger.AvailableNow, la découverte de fichiers s'effectue de manière asynchrone avec le traitement des données et les données peuvent être traitées sur plusieurs micro-batches avec limitation de débit. Auto Loader par default traite un maximum de 1 000 fichiers par micro-batch. Vous pouvez configurer cloudFiles.maxFilesPerTrigger et cloudFiles.maxBytesPerTrigger pour déterminer le nombre de fichiers ou d'octets à traiter dans un micro-batch. La limite de fichiers est une limite stricte, mais la limite d'octets est une limite souple, ce qui signifie que davantage d'octets peuvent être traités que la valeur fournie maxBytesPerTrigger. Lorsque les options sont fournies ensemble, Auto Loader traite autant de fichiers nécessaires pour atteindre l'une des limites.
Emplacement du point de contrôle
L'emplacement du point de contrôle est utilisé pour stocker l'état et les informations de progression du Stream. Databricks recommande de définir l'emplacement du point de contrôle à un emplacement sans politique de cycle de vie d'objet cloud. Si les fichiers dans l'emplacement du point de contrôle sont nettoyés conformément à la politique, l'état du Stream est corrompu. Si cela se produit, vous devez redémarrer le Stream à partir de zéro.
Suivi des événements de fichiers
Auto Loader assure le suivi des fichiers découverts dans l'emplacement de point de contrôle en utilisant RocksDB pour fournir des garanties d'ingestion « exactement une fois ». Pour les flux d'ingestion à volume élevé ou de longue durée, Databricks recommande de passer à Databricks Runtime 15,4 LTS ou version ultérieure. Dans ces versions, Auto Loader n'attend pas que l'état entier de RocksDB soit downloadé avant que le Stream ne start, ce qui peut accélérer le temps de Startup du Stream.
Si vous souhaitez empêcher les états de fichier de croître sans limites, vous pouvez également envisager d'utiliser l'option cloudFiles.maxFileAge pour faire expirer les événements de fichier qui sont plus anciens qu'un certain âge. La valeur minimale que vous pouvez définir pour cloudFiles.maxFileAge est "14 days". Les suppressions dans RocksDB apparaissent comme des entrées de suppression. Par conséquent, l'utilisation du stockage pourrait augmenter temporairement à mesure que les événements expirent avant qu'il ne start à se stabiliser.
cloudFiles.maxFileAge est fourni comme mécanisme de contrôle des coûts pour les datasets à volume élevé. Le réglage trop agressif de cloudFiles.maxFileAge peut entraîner des problèmes de qualité des données, tels qu'une ingestion en double ou des fichiers manquants. Par conséquent, Databricks recommande un réglage conservateur pour cloudFiles.maxFileAge, tel que 90 jours, ce qui est similaire à ce que recommandent des solutions d'ingestion de données comparables.
Tenter de régler l'option cloudFiles.maxFileAge peut entraîner l'ignorance des fichiers non traités par Auto Loader ou l'expiration des fichiers déjà traités, puis leur retraitement, ce qui provoque des données en double. Voici quelques éléments à prendre en compte lors du choix d'un cloudFiles.maxFileAge:
- Si votre Stream redémarre après une longue période, les événements de notification de fichier qui sont extraits de la file d'attente et qui sont plus anciens que
cloudFiles.maxFileAgesont ignorés. De même, si vous utilisez la liste de répertoires, les fichiers qui pourraient être apparus pendant le temps d'arrêt et qui sont plus anciens quecloudFiles.maxFileAgesont ignorés. - Si vous utilisez le mode de listage de répertoires et utilisez
cloudFiles.maxFileAge, par exemple défini sur"1 month", vous arrêtez votre Stream et le redémarrez aveccloudFiles.maxFileAgedéfini sur"2 months", les fichiers datant de plus d'un mois, mais de moins de deux mois, sont retraités.
Si vous définissez cette option la première fois que vous start le Stream, vous n'ingérerez pas de données plus anciennes que cloudFiles.maxFileAge. Par conséquent, si vous souhaitez ingérer d'anciennes données, vous ne devez pas définir cette option lorsque vous start votre Stream pour la première fois. Cependant, vous devriez définir cette option lors des exécutions ultérieures.
Trigger des remplissages réguliers à l'aide de cloudFiles.backfillInterval
Dans de rares cas, des fichiers peuvent être manquants ou en retard lorsqu'on dépend uniquement de systèmes de notification, par exemple lorsque les limites de rétention des messages de notification sont atteintes. Si vous avez des exigences strictes en matière d'exhaustivité des données et de SLA, envisagez de définir cloudFiles.backfillInterval pour Trigger des remplissages asynchrones à un intervalle spécifié. Par exemple, réglez-le sur un jour pour les remplissages quotidiens, ou sur une semaine pour les remplissages hebdomadaires. Le déclenchement de remplissages réguliers n'entraîne pas de doublons.
Lorsque vous utilisez des événements de fichier, exécutez votre Stream au moins une fois tous les 7 jours.
Lorsque vous utilisez les événements de fichiers, exécutez vos Stream Auto Loader au moins une fois tous les 7 jours pour éviter une liste complète de répertoires. L'exécution fréquente de votre Stream Auto Loader garantira une découverte incrémentielle des fichiers.
Pour des meilleures pratiques complètes en matière d'événements de fichiers gérés, consultez Meilleures pratiques pour Auto Loader avec les événements de fichiers.