Surveiller et observer Auto Loader
Les pipelines Auto Loader nécessitent un monitoring actif pour détecter les problèmes tels que l’augmentation des retards, le drift de schéma, les données corrompues et les Streams bloqués avant qu’ils n’affectent les consommateurs en aval. Cette page explique comment monitorer les métriques clés, query l’état au niveau du fichier, créer des tableaux de bord d’observabilité et résoudre les problèmes courants.
Pour plus de détails sur la configuration de production, voir Configurer Auto Loader pour les charges de travail de production. Pour les bonnes pratiques de configuration, voir les bonnes pratiques d'Auto Loader.
Prérequis
Plusieurs workflows de monitoring sur cette page s'appuient sur cloud_files_state() pour observer l'état d'ingestion par fichier — y compris les query en attente, les calculs de latence et la détection de drift de schéma. cloud_files_state() est une fonction à valeur de table qui renvoie l'état d'ingestion au niveau du fichier pour un point de contrôle Auto Loader. Tous ses champs ne sont pas disponibles par default. La disponibilité dépend de votre version de Databricks Runtime et de la configuration :
- Databricks Runtime 18.2 et versions supérieures :
discovery_time,processed_timeetcommit_timesont disponibles automatiquement. Sur Databricks Runtime entre 16.4 et 18.1, Ces champs sont disponibles uniquement lorsquecloudFiles.cleanSourceest activé. - **Databricks Runtime 16.4
cloudFiles.cleanSourcearchive_time``archive_modeetmove_locationversions supérieures avec activé** :, et sont disponibles.
L'activation de cloudFiles.cleanSource entraîne une certaine surcharge de performance. Effectuez des tests de performance sur vos charges de travail dans un environnement de pré-production avant de l'activer en production.
De plus :
- Annotez les données ingérées avec la colonne
_metadata. Capturer au minimumfile_pathetfile_modification_time. Voir Colonne de métadonnées de fichier. - Activer les colonnes
_rescued_dataet_corrupt_record.
Principales métriques d’Auto Loader
Le tableau suivant résume les métriques les plus importantes à surveiller pour les pipelines Auto Loader. Ces métriques sont disponibles à partir des événements de progression StreamingQueryListener, avec des valeurs spécifiques à Auto Loader exposées sous la carte metrics de chaque source.
Métriques | Ce que cela vous dit |
|---|---|
| Nombre de fichiers dans le backlog en attente de traitement. |
| Taille du backlog de fichiers en octets |
| Profondeur de la file d'attente du cloud (mode de notification de fichier uniquement) |
| Lignes traitées par batch |
| Taux d'arrivée des données |
| Throughput de traitement |
| Où le temps est passé dans chaque batch |
À surveiller
Les modèles suivants indiquent que votre pipeline peut nécessiter une attention.
- Augmentation de
numFilesOutstanding: Le backlog s'accumule. Votre pipeline prend du retard sur les données entrantes. processedRowsPerSecond<inputRowsPerSecond: le pipeline traite les données plus lentement qu’elles n’arrivent.- Grand
durationMs.latestOffset: La découverte de fichiers est lente. Envisagez de passer aux événements de fichier. - Grand
durationMs.addBatch: Le traitement des données est lent. Envisagez de mettre à l'échelle le compute ou d'optimiser les transformations.
Pour la référence complète des indicateurs, consultez les indicateurs source d'Auto Loader.
État de la query au niveau du fichier avec cloud_files_state
La fonction à valeur de table cloud_files_state() fournit des informations détaillées sur chaque fichier découvert par Auto Loader. Les champs suivants sont disponibles. Les champs marqués comme nécessitant Databricks Runtime 16.4 ou version supérieure ou 18.2 ou version supérieure ne sont renseignés que dans les conditions décrites dans Prérequis.
Champ | Type | Description |
|---|---|---|
|
| Le chemin d'accès au fichier |
|
| La taille du fichier en octets |
|
| Lorsque le fichier a été créé |
|
| Lorsque l'Auto Loader a découvert le fichier (Databricks Runtime 16.4 et versions supérieures) |
|
| Lorsque Auto Loader a traité le fichier (Databricks Runtime 16,4 et versions supérieures) |
|
| Lorsque le fichier a été validé au point de contrôle (Databricks Runtime 16.4 et versions ultérieures) |
|
| Lorsque le fichier a été archivé (nécessite |
|
|
|
|
| Chemin de destination lorsque |
|
| État actuel de l'ingestion de fichiers |
Examiner l'état d'ingestion des fichiers
Les queries suivantes couvrent des scénarios de diagnostic courants.
Trouvez tous les fichiers non traités (le backlog actuel) :
SELECT * FROM cloud_files_state('path/to/checkpoint')
WHERE ingestion_state != 'COMMITTED';
Latence moyenne d'ingestion du compute (temps entre la création du fichier et le commit) :
SELECT avg(unix_timestamp(commit_time) - unix_timestamp(create_time)) AS avg_latency_seconds
FROM cloud_files_state('path/to/checkpoint')
WHERE commit_time IS NOT NULL AND create_time IS NOT NULL;
Rechercher les fichiers corrompus ou ignorés :
SELECT path, ingestion_state, size, create_time
FROM cloud_files_state('path/to/checkpoint')
WHERE ingestion_state LIKE 'SKIPPED%';
Suivre la progression de l'archivage (nécessite cloudFiles.cleanSource) :
SELECT archive_mode, count(*) AS file_count
FROM cloud_files_state('path/to/checkpoint')
GROUP BY archive_mode;
Trouvez les fichiers avec une latence élevée entre la découverte et le commit pour identifier les goulets d'étranglement :
SELECT
path,
size,
unix_timestamp(commit_time) - unix_timestamp(discovery_time) AS processing_latency_seconds,
unix_timestamp(commit_time) - unix_timestamp(create_time) AS end_to_end_latency_seconds
FROM cloud_files_state('path/to/checkpoint')
WHERE commit_time IS NOT NULL
ORDER BY end_to_end_latency_seconds DESC
LIMIT 20;
Pour la référence SQL complète, consultez la fonction tablecloud_files_state.
Surveiller Auto Loader dans les LakeFlow Pipelines
Databricks recommande d'utiliser les Lakeflow pipelines pour les pipelines de production Auto Loader. Pour tirer parti de ses capacités de monitoring intégrées :
-
Stockez le journal des événements LakeFlow Pipelines dans une table Delta afin qu'il puisse être interrogé pour les données d'observabilité. Configurez ceci via les paramètres avancés du pipeline ou l'API. Pour plus de détails, consultez le log des événements du pipeline.
-
Structurez votre pipeline pour l'observabilité. Un pipeline Auto Loader bien structuré dans les LakeFlow Pipelines comprend une vue
{table}_source(la définition de la source Auto Loader), une table de streaming{table}_bronze(ingestion de données brutes avec les colonnes_rescued_dataet_corrupt_record), uncorrupt_records_sinkqui met en quarantaine les lignes avec des données non analysables, et une vue propre{table}pour la consommation en aval. -
Définissez des attentes sur vos tables de streaming bronze pour surveiller le drift de schéma et la corruption des données.
_rescued_data IS NULLdétecte les modifications de schéma inattendues et_corrupt_record IS NULLdétecte les données inanalysables. Les Lakeflow Pipelines évaluent ces attentes à mesure que les données arrivent et génèrent une piste d'observabilité. Vous pouvez configurer les attentes pour avertir, supprimer des lignes ou faire échouer le pipeline.
Après avoir créé la vue event_log_raw pour votre pipeline, utilisez les queries suivantes pour les métriques spécifiques à Auto Loader.
Surveiller le throughput d’ingestion par flux :
SELECT
origin.flow_name,
origin.update_id,
timestamp,
TRY_CAST(details:flow_progress.metrics.num_output_rows AS BIGINT) AS rows_written
FROM event_log_raw
WHERE event_type = 'flow_progress'
ORDER BY timestamp DESC;
Surveillez le backlog de données par flux :
SELECT
origin.flow_name,
timestamp,
DOUBLE(details:flow_progress.metrics.backlog_bytes) AS backlog_bytes
FROM event_log_raw
WHERE event_type = 'flow_progress'
AND details:flow_progress.metrics.backlog_bytes IS NOT NULL
ORDER BY timestamp DESC;
Résumer les violations d'attentes pour détecter la drift de schéma et les données corrompues :
SELECT
origin.flow_name,
explode(from_json(
details:flow_progress.data_quality.expectations,
'array<struct<name:string, dataset:string, passed_records:bigint, failed_records:bigint>>'
)) AS expectation
FROM event_log_raw
WHERE event_type = 'flow_progress'
AND details:flow_progress.data_quality.expectations IS NOT NULL;
Pour des conseils généraux sur le monitoring des LakeFlow Pipelines, consultez Superviser les pipelines et Journal des événements du pipeline.
Surveiller Auto Loader avec Structured Streaming
Lorsque vous exécutez Auto Loader en dehors des Lakeflow pipelines, utilisez les approches de monitoring Structured Streaming suivantes.
- Implémentez un
StreamingQueryListenerpour capturer les métriques spécifiques à Auto Loader de chaque batch en lisant à partir desource.metrics.
from pyspark.sql.streaming import StreamingQueryListener
class AutoLoaderMonitor(StreamingQueryListener):
def onQueryStarted(self, event):
pass
def onQueryProgress(self, event):
for source in event.progress.sources:
if "CloudFilesSource" in source.description:
metrics = source.metrics
files_outstanding = metrics.get("numFilesOutstanding", "0")
bytes_outstanding = metrics.get("numBytesOutstanding", "0")
rows_per_sec = source.processedRowsPerSecond
# Push metrics to your monitoring system (for example, write to a Delta table)
def onQueryIdle(self, event):
pass
def onQueryTerminated(self, event):
pass
spark.streams.addListener(AutoLoaderMonitor())
La logique de traitement dans les écouteurs peut ralentir le traitement des query. Limitez le calcul dans les rappels de l'écouteur et évitez les écritures externes synchrones à cet endroit ; au lieu de cela, émettez de la télémétrie légère de manière asynchrone ou transmettez les métriques à un Job distinct pour la persistance.
-
Utilisez
numInputRows,inputRowsPerSecondetprocessedRowsPerSecondde la progression source pour calculer le throughput — fichiers par seconde et lignes par seconde pour chaque batch. -
Pour calculer la latence d'ingestion, comparez
create_timeetcommit_timedecloud_files_state()pour la latence de bout en bout. Pour la latence de traitement, utilisez la répartitiondurationMs(par exemple,latestOffset,addBatch, et d'autres phases de batch signalées) pour identifier la phase qui constitue le goulet d'étranglement. -
Utilisez
df.observe()pour définir des métriques de qualité des données en ligne directement sur le DataFrame de streaming. Les métriques sont visibles dans lesStreamingQueryListenerévénements de progression sousobservedMetrics.
from pyspark.sql.functions import count, lit, col
observed_df = df.observe(
"auto_loader_quality",
count(lit(1)).alias("total_rows"),
count(col("_rescued_data")).alias("rescued_rows"),
count(col("_corrupt_record")).alias("corrupt_rows")
)
- Utilisez
.queryName()pour attribuer un nom unique à chaque Stream, facilitant ainsi la distinction des Streams Auto Loader dans l'onglet Streaming de l'interface utilisateur Spark et dans les tableaux de bord de monitoring.
Pour la référence complète sur le monitoring Structured Streaming, consultez Monitoring des queries Structured Streaming sur Databricks.
Créer un tableau de bord d'observabilité
Combinez les données de plusieurs sources pour créer un tableau de bord d'observabilité complet pour vos pipelines Auto Loader. Ce tableau affiche des sources suggérées que vous pouvez utiliser pour structurer votre tableau de bord d'observabilité.
Source de données | Données d'observabilité |
|---|---|
| État d'ingestion au niveau du fichier : découverte, traitement, commit et Timestamp d'archivage par fichier |
Log des événements LakeFlow Pipelines | Historique d'exécution des pipelines, métriques de flux par batch et résultats des attentes en matière de qualité des données |
Tables de sortie du pipeline | Nombre de lignes et volume de données écrits par table ingérée |
Vous pouvez ensuite agréger les données d'observabilité dans des tables dédiées qui servent de base aux tableaux de bord et aux alertes :
- Récapitulez les statuts d'exécution de pipeline (succès ou échec) au fil du temps, dérivés des événements
event_type = 'update_progress'. - Mesures agrégées d'ingestion de fichiers (taille du backlog, throughput, latence par batch), dérivées des événements
cloud_files_state()etevent_type = 'flow_progress'. - Développer des statistiques de table en utilisant le nombre de lignes et le volume de données par table, dérivées de
num_output_rowsdans l'event Logs. - Recueillez les informations de debugging à partir des logs d'erreurs détaillés et des violations d'attente par mise à jour, dérivées des événements
event_type = 'flow_progress'avecdata_qualityrenseigné.
Ces tables agrégées peuvent alimenter un AI/BI dashboard et des alertes SQL. Les tableaux de bord recommandés comprennent la chronologie de l'état d'exécution du pipeline, la tendance de l'arriéré d'ingestion, la tendance du throughput, la distribution de la latence d'ingestion, les métriques de qualité des données, les événements d'évolution des schémas et l'état d'archivage des fichiers.
Surveiller les événements d'évolution des schémas
Utilisez les approches suivantes pour détecter les modifications de schéma telles qu'elles se produisent.
- Les valeurs non NULL dans
_rescued_datadans les décomptes de violation des attentes indiquent une drift de schéma. Interrogez le Logs des événements pourfailed_records > 0sur l'attenteno rescued data. - Les modifications apportées au répertoire
_schemasà l'intérieur ducloudFiles.schemaLocationconfiguré (ou uniquement à l'intérieur du point de contrôle lorsque l'emplacement du schéma n'est pas défini séparément) indiquent qu'une évolution des schémas s'est produite. Vous pouvez interroger ce répertoire à partir d'un Job de monitoring distinct. - Ne traitez pas un événement
onQueryTerminatedsuivi deonQueryStartedpour le même nom de stream comme une preuve suffisante d'évolution des schémas en soi. Les streams redémarrent pour de nombreuses raisons (redémarrages de clusters, déploiements de code, erreurs de stockage transitoires). Corrélez les redémarrages avec des signaux indépendants — changements de répertoire_schemasou violations d'attente_rescued_data— avant de conclure qu'une évolution des schémas s'est produite. - Utilisez
_metadata.file_pathpour identifier les fichiers ayant introduit des modifications de schéma. Joignez-le àcloud_files_state()sur le champpathpour corréler les modifications de schéma avec des fichiers et des batch spécifiques.
Utilisez cet exemple de query pour détecter une drift de schéma récente via des violations d'attentes :
SELECT
timestamp,
origin.flow_name,
exp.name AS expectation_name,
exp.failed_records
FROM (
SELECT
timestamp,
origin,
explode(from_json(
details:flow_progress.data_quality.expectations,
'array<struct<name:string, dataset:string, passed_records:bigint, failed_records:bigint>>'
)) AS exp
FROM event_log_raw
WHERE event_type = 'flow_progress'
AND details:flow_progress.data_quality.expectations IS NOT NULL
)
WHERE exp.name = '<rescued-data expectation name>'
AND exp.failed_records > 0
ORDER BY timestamp DESC;
Configurez des alertes pour les problèmes courants
Utilisez les alertes Databricks SQL ou les notifications de pipeline pour détecter les problèmes avant qu'ils n'affectent les consommateurs en aval.
Le code SQL suivant détecte un arriéré croissant et peut servir de base à une alerte Databricks SQL. Planifiez son exécution périodique (par exemple, toutes les 5 minutes) et alertez lorsque le résultat n'est pas vide.
-- Alert when backlog exceeds threshold or trends upward across recent batches
WITH recent_backlog AS (
SELECT
origin.flow_name,
timestamp,
DOUBLE(details:flow_progress.metrics.backlog_bytes) AS backlog_bytes,
ROW_NUMBER() OVER (PARTITION BY origin.flow_name ORDER BY timestamp DESC) AS rn
FROM event_log_raw
WHERE event_type = 'flow_progress'
AND details:flow_progress.metrics.backlog_bytes IS NOT NULL
)
SELECT flow_name, backlog_bytes, timestamp
FROM recent_backlog
WHERE rn = 1
AND backlog_bytes > 1073741824 -- alert when backlog exceeds 1 GB
Le tableau suivant récapitule les conditions d'alerte recommandées :
Que détecter. | Comment le détecter | Quand déclencher une alerte |
|---|---|---|
Backlog croissant |
| Augmentation soutenue sur plusieurs batchs |
Stream bloqué | Aucun événement de progression | Aucun événement pendant N minutes (basé sur l'intervalle de Trigger attendu) |
Latence d'ingestion élevée |
| Dépasse votre threshold SLA |
Dégradation de la qualité des données | Taux d'échec des attentes | Pourcentage croissant de lignes ne répondant pas aux attentes |
Événement d'évolution des schémas |
| Toutes les valeurs non NULL dans le nombre de violations d'attentes |
Découverte de fichiers lente |
| Considérablement supérieure à la base de référence. |
Dépanner les problèmes courants
Le tableau suivant décrit les problèmes courants de pipeline Auto Loader, leurs causes probables et les actions recommandées pour les résoudre.
Problème | Cause possible | Action recommandée |
|---|---|---|
Backlog qui augmente plus vite que le traitement | Compute sous-dimensionné, asymétrie des données ou limites de débit régulées | Monter en charge le compute, vérifiez l'asymétrie avec la Spark UI et examinez les paramètres |
Fichiers non découverts | Événements de fichier mal configurés, problème d'autorisations ou Stream non exécuté dans les 7 jours | Vérifiez les autorisations d'emplacement externe, contrôlez la configuration des événements de fichier dans l'interface utilisateur de Unity Catalog et assurez-vous que le stream s'exécute au moins tous les 7 jours pour éviter l'expiration de l'état de RocksDB |
Le Startup du Stream prend trop de temps | Large checkpoint state download (RocksDB) | Mettre à niveau vers Databricks Runtime 15.3 et versions ultérieures pour le chargement d'état asynchrone, ce qui réduit le temps de Startup d'environ 90 % |
Traitement des fichiers en double | Paramètres | Utilisez un |
Évolution des schémas entraînant des redémarrages de pipelines | Modifications de schéma fréquentes ou incompatibles | Examinez |
Données corrompues s'accumulant dans le stockage | Problèmes de qualité des données source | Vérifiez le récepteur de quarantaine |
| Exécutant sur Databricks Runtime inférieur à 18,2 sans | Mettre à niveau vers Databricks Runtime 18.2 ou version ultérieure, ou activer |
Pour un dépannage supplémentaire, voir la FAQ d'Auto Loader.