Aller au contenu principal

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_time et commit_time sont disponibles automatiquement. Sur Databricks Runtime entre 16.4 et 18.1, Ces champs sont disponibles uniquement lorsque cloudFiles.cleanSource est activé.
  • **Databricks Runtime 16.4 cloudFiles.cleanSource archive_time``archive_modeet move_location versions 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 minimum file_path et file_modification_time. Voir Colonne de métadonnées de fichier.
  • Activer les colonnes _rescued_data et _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

numFilesOutstanding

Nombre de fichiers dans le backlog en attente de traitement.

numBytesOutstanding

Taille du backlog de fichiers en octets

approximateQueueSize

Profondeur de la file d'attente du cloud (mode de notification de fichier uniquement)

numInputRows

Lignes traitées par batch

inputRowsPerSecond

Taux d'arrivée des données

processedRowsPerSecond

Throughput de traitement

durationMs répartition

Où le temps est passé dans chaque batch

Métriques

Ce que cela vous dit

numFilesOutstanding

Nombre de fichiers dans le backlog en attente de traitement.

numBytesOutstanding

Taille du backlog de fichiers en octets

approximateQueueSize

Profondeur de la file d'attente du cloud (mode de notification de fichier uniquement)

numInputRows

Lignes traitées par batch

inputRowsPerSecond

Taux d'arrivée des données

processedRowsPerSecond

Throughput de traitement

durationMs répartition

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

path

STRING

Le chemin d'accès au fichier

size

BIGINT

La taille du fichier en octets

create_time

TIMESTAMP

Lorsque le fichier a été créé

discovery_time

TIMESTAMP

Lorsque l'Auto Loader a découvert le fichier (Databricks Runtime 16.4 et versions supérieures)

processed_time

TIMESTAMP

Lorsque Auto Loader a traité le fichier (Databricks Runtime 16,4 et versions supérieures)

commit_time

TIMESTAMP

Lorsque le fichier a été validé au point de contrôle (Databricks Runtime 16.4 et versions ultérieures)

archive_time

TIMESTAMP

Lorsque le fichier a été archivé (nécessite cloudFiles.cleanSource)

archive_mode

STRING

MOVE, DELETE, ou NULL (nécessite cloudFiles.cleanSource)

move_location

STRING

Chemin de destination lorsque cloudFiles.cleanSource est MOVE

ingestion_state

STRING

État actuel de l'ingestion de fichiers

Champ

Type

Description

path

STRING

Le chemin d'accès au fichier

size

BIGINT

La taille du fichier en octets

create_time

TIMESTAMP

Lorsque le fichier a été créé

discovery_time

TIMESTAMP

Lorsque l'Auto Loader a découvert le fichier (Databricks Runtime 16.4 et versions supérieures)

processed_time

TIMESTAMP

Lorsque Auto Loader a traité le fichier (Databricks Runtime 16,4 et versions supérieures)

commit_time

TIMESTAMP

Lorsque le fichier a été validé au point de contrôle (Databricks Runtime 16.4 et versions ultérieures)

archive_time

TIMESTAMP

Lorsque le fichier a été archivé (nécessite cloudFiles.cleanSource)

archive_mode

STRING

MOVE, DELETE, ou NULL (nécessite cloudFiles.cleanSource)

move_location

STRING

Chemin de destination lorsque cloudFiles.cleanSource est MOVE

ingestion_state

STRING

É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) :

SQL
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) :

SQL
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 :

SQL
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) :

SQL
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 :

SQL
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_data et _corrupt_record), un corrupt_records_sink qui 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 NULL détecte les modifications de schéma inattendues et _corrupt_record IS NULL dé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 :

SQL
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 :

SQL
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 :

SQL
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 StreamingQueryListener pour capturer les métriques spécifiques à Auto Loader de chaque batch en lisant à partir de source.metrics.
Python
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())
remarque

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, inputRowsPerSecond et processedRowsPerSecond de 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_time et commit_time de cloud_files_state() pour la latence de bout en bout. Pour la latence de traitement, utilisez la répartition durationMs (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 les StreamingQueryListener événements de progression sous observedMetrics.

Python
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é

cloud_files_state()

É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

Source de données

Données d'observabilité

cloud_files_state()

É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() et event_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_rows dans 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' avec data_quality renseigné.

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_data dans les décomptes de violation des attentes indiquent une drift de schéma. Interrogez le Logs des événements pour failed_records > 0 sur l'attente no rescued data.
  • Les modifications apportées au répertoire _schemas à l'intérieur du cloudFiles.schemaLocation configuré (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 onQueryTerminated suivi de onQueryStarted pour 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 _schemas ou violations d'attente _rescued_data — avant de conclure qu'une évolution des schémas s'est produite.
  • Utilisez _metadata.file_path pour identifier les fichiers ayant introduit des modifications de schéma. Joignez-le à cloud_files_state() sur le champ path pour 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 :

SQL
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.

SQL
-- 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

numFilesOutstanding tendance à la hausse

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

commit_time - create_time

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

_rescued_data IS NOT NULL

Toutes les valeurs non NULL dans le nombre de violations d'attentes

Découverte de fichiers lente

durationMs.latestOffset

Considérablement supérieure à la base de référence.

Que détecter.

Comment le détecter

Quand déclencher une alerte

Backlog croissant

numFilesOutstanding tendance à la hausse

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

commit_time - create_time

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

_rescued_data IS NOT NULL

Toutes les valeurs non NULL dans le nombre de violations d'attentes

Découverte de fichiers lente

durationMs.latestOffset

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 maxFilesPerTrigger pour contrôler la taille des batchs.

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 cloudFiles.maxFileAge agressifs ou corruption des points de contrôle.

Utilisez un maxFileAge conservateur (90+ jours minimum), vérifiez l'intégrité du point de contrôle et évitez les politiques de cycle de vie sur le stockage des points de contrôle.

Évolution des schémas entraînant des redémarrages de pipelines

Modifications de schéma fréquentes ou incompatibles

Examinez schemaEvolutionMode, passez à addNewColumnsWithTypeWidening pour les promotions de type, ou utilisez le type Variant pour les schémas très dynamiques

Données corrompues s'accumulant dans le stockage

Problèmes de qualité des données source

Vérifiez le récepteur de quarantaine _corrupt_record pour les modèles, examinez la génération des données sources et envisagez d'ajouter une validation en amont

discovery_time et commit_time non renseigné

Exécutant sur Databricks Runtime inférieur à 18,2 sans cleanSource

Mettre à niveau vers Databricks Runtime 18.2 ou version ultérieure, ou activer cloudFiles.cleanSource sur Databricks Runtime de 16.4 à 18.1

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 maxFilesPerTrigger pour contrôler la taille des batchs.

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 cloudFiles.maxFileAge agressifs ou corruption des points de contrôle.

Utilisez un maxFileAge conservateur (90+ jours minimum), vérifiez l'intégrité du point de contrôle et évitez les politiques de cycle de vie sur le stockage des points de contrôle.

Évolution des schémas entraînant des redémarrages de pipelines

Modifications de schéma fréquentes ou incompatibles

Examinez schemaEvolutionMode, passez à addNewColumnsWithTypeWidening pour les promotions de type, ou utilisez le type Variant pour les schémas très dynamiques

Données corrompues s'accumulant dans le stockage

Problèmes de qualité des données source

Vérifiez le récepteur de quarantaine _corrupt_record pour les modèles, examinez la génération des données sources et envisagez d'ajouter une validation en amont

discovery_time et commit_time non renseigné

Exécutant sur Databricks Runtime inférieur à 18,2 sans cleanSource

Mettre à niveau vers Databricks Runtime 18.2 ou version ultérieure, ou activer cloudFiles.cleanSource sur Databricks Runtime de 16.4 à 18.1

Pour un dépannage supplémentaire, voir la FAQ d'Auto Loader.