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 du déplacement doivent se trouver dans le même emplacement externe, volume ou montage DBFS. Les déplacements entre buckets et entre conteneurs 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 obtenir la référence complète sur les options cleanSource et leurs default, consultez Nettoyer les fichiers traités avec Auto Loader.
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',
`cleanSource.moveDestination` => 's3://my-bucket/archive/landing/',
`cleanSource.retentionDuration` => '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 effectue le suivi des fichiers détectés dans l'emplacement de point de contrôle à l'aide de RocksDB afin de garantir que chaque fois un fichier est ingéré exactement une fois. Pour les flux d'ingestion à haut volume ou de longue durée, Databricks recommande la mise à niveau vers Databricks Runtime 15.4 LTS ou une version ultérieure. Dans ces versions, Auto Loader n'attend pas que l'intégralité de l'état 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 indéfiniment, vous pouvez également envisager d'utiliser l'option cloudFiles.maxFileAge pour expirer les événements de fichier 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 en tant qu'entrées de type tombstone. Par conséquent, il est possible que vous constatiez une augmentation temporaire de l'utilisation du stockage à mesure que les événements expirent, avant qu'elle 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 lors du premier start du stream, vous n'ingérez pas de données antérieures à cloudFiles.maxFileAge. Par conséquent, si vous souhaitez ingérer des données anciennes, ne définissez pas cette option lors du premier start de votre stream. Cependant, définissez cette option lors des exécutions ultérieures.
Trigger des remplissages réguliers à l'aide de cloudFiles.backfillInterval
Un remplissage rétroactif est un listage asynchrone des répertoires qu'Auto Loader exécute en parallèle avec la découverte normale des fichiers afin de récupérer ceux qui ont été omis. Bien que les systèmes de notification cloud délivrent des événements au moins une fois, un fichier peut tout de même être manqué. Un remplissage rétroactif périodique répertorie à nouveau le répertoire source afin que les fichiers manqués soient finissent par être découverts.
Définissez cloudFiles.backfillInterval sur une chaîne de durée telle que 1 day ou 1 week pour planifier des remplissages rétroactifs périodiques. Il n'y a pas de default. En mode de listage de répertoire et de notification de fichier classique, les remplissages rétroactifs ne s'exécutent que lorsque vous définissez cet intervalle.
Comportement des opérations de backfill :
- Basé sur le temps, et non sur les fichiers : Auto Loader déclenche le remplissage rétroactif suivant lorsque le temps écoulé depuis le dernier dépasse l'intervalle, en suivant l'heure du dernier remplissage dans le point de contrôle plutôt qu'en comparant les timestamps des fichiers. Pour confirmer le remplissage rétroactif le plus récent, utilisez les métriques
lastBackfillStartTimeMsetlastBackfillFinishTimeMs. Consultez Monitoring Structured Streaming queries on Databricks. - Only missed files are ingested : un remplissage rétroactif n'ingère que les fichiers qui n'ont pas encore été traités, et les nouveaux fichiers continuent d'arriver via le mode de découverte configuré. Il ignore les fichiers déjà ingérés en vérifiant leur état dans le point de contrôle, ce qui évite les doublons.
- Asynchrone : les remplissages rétroactifs s'exécutent en arrière-plan et ne bloquent pas le traitement par micro-batch.
Définissez un intervalle de remplissage rétroactif lorsque vous utilisez le mode de notification de fichier classique et que vous avez des exigences strictes en matière d'exhaustivité des données ou d'accord de niveau de service (SLA). Auto Loader répertorie ensuite la source selon cette cadence pour détecter les notifications manquées.
Ne définissez pas d'intervalle de remplissage rétroactif lors de l'utilisation de file events. Databricks remplit automatiquement ces emplacements externes avec une liste complète lors de la première activation des événements de fichiers, puis continue de les remplir environ toutes les 24 heures pendant qu'un stream ingère des données. Le paramètre d'intervalle n'est pas pris en charge avec les événements de fichiers et, le remplissage rétroactif automatique étant moins coûteux, Databricks recommande d'utiliser les événements de fichiers plutôt que de définir un intervalle manuel.
Chaque remplissage rétroactif correspond à une liste complète des répertoires ; par conséquent, son coût est proportionnel au nombre de fichiers dans le répertoire source et, en mode de listage de répertoires, entraîne des frais d'API LIST. Choisissez l'intervalle le plus long qui respecte toujours votre SLA de exhaustivité.
Évitez les listages de répertoires complets avec des événements de fichiers
Lors de l'utilisation d'événements de fichiers, exécutez vos flux Auto Loader au moins tous les 7 jours pour éviter un listage complet des répertoires. L'exécution de vos flux Auto Loader à cette fréquence garantit que la découverte des fichiers est incrémentielle.
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.