Surveiller et observer le chargeur automatique

Les pipelines de chargeur automatique nécessitent une surveillance active pour détecter les problèmes tels que les backlogs croissants, la dérive de schéma, les données endommagées et les flux bloqués avant qu’ils n’affectent les consommateurs en aval. Cette page explique comment surveiller les métriques clés, interroger l’état au niveau du fichier, générer des tableaux de bord d’observabilité et résoudre les problèmes courants.

Pour plus d’informations sur la configuration de production, consultez Configurer le chargeur automatique pour les charges de travail de production. Pour connaître les meilleures pratiques de configuration, consultez les meilleures pratiques relatives au chargeur automatique.

Prerequisites

Plusieurs processus de surveillance de cette page s’appuient sur cloud_files_state() pour observer l’état d’ingestion de chaque fichier, notamment les requêtes sur l’arriéré, les calculs de latence et la détection de dérive du schéma. cloud_files_state() est une fonction table qui retourne l’état d’ingestion au niveau des fichiers pour un point de contrôle Auto Loader. Tous ses champs ne sont pas disponibles par défaut. La disponibilité dépend de la version et de la configuration de Databricks Runtime :

  • Databricks Runtime 18.2 et versions ultérieures : discovery_time, processed_timeet commit_time sont disponibles automatiquement. Sur Databricks Runtime 16.4 à 18.1, ces champs sont disponibles uniquement quand cloudFiles.cleanSource ils sont activés.
  • Databricks Runtime 16.4 et versions ultérieures avec cloudFiles.cleanSource activé : archive_time, archive_modeet move_location sont disponibles.

L’activation cloudFiles.cleanSource présente une surcharge de performances. Évaluez ses performances par rapport à vos charges de travail dans un environnement de préproduction avant de l’activer en production.

Plus:

  • Annotez les données ingérées avec la _metadata colonne. Capturez au minimum file_path et file_modification_time. Consultez Colonne métadonnées de fichier.
  • Activez _rescued_data et _corrupt_record colonnes.

Principales métriques du chargeur automatique

Le tableau suivant récapitule les métriques les plus importantes à surveiller pour les pipelines de chargeur automatique. Ces métriques sont disponibles à partir des événements de progression StreamingQueryListener, les valeurs propres à Auto Loader étant exposées dans le mappage metrics de chaque source.

Unité de mesure Ce qu’il 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 fichiers uniquement)
numInputRows Lignes traitées par lot
inputRowsPerSecond Taux d’arrivée des données
processedRowsPerSecond Débit de traitement
durationMs panne Répartition du temps dans chaque lot

Ce à quoi faire attention

Les modèles suivants indiquent que votre pipeline peut avoir besoin d’attention.

  • numFilesOutstanding en augmentation : le backlog s’accumule. Votre pipeline prend du retard par rapport aux données entrantes.
  • processedRowsPerSecond < inputRowsPerSecond: Le pipeline traite les données plus lentement qu’elles n’arrivent.
  • Grande taille durationMs.latestOffset : la détection des fichiers est lente. Envisagez de passer aux événements de fichiers.
  • Large durationMs.addBatch: Le traitement des données est lent. Envisagez de mettre à l’échelle le calcul ou d’optimiser les transformations.

Pour obtenir la référence complète des métriques, consultez les métriques sources du chargeur automatique.

Interroger l’état au niveau du fichier avec cloud_files_state

La cloud_files_state() fonction renvoyant une table fournit des informations détaillées sur chaque fichier détecté par Auto Loader. Les champs suivants sont disponibles. Les champs marqués comme nécessitant Databricks Runtime 16.4 et versions ultérieures ou 18.2 et ultérieures ne sont renseignés que dans les conditions décrites dans les conditions préalables.

Champ Type Description
path STRING Chemin d’accès du fichier
size BIGINT Taille du fichier en octets
create_time TIMESTAMP Lorsque le fichier a été créé
discovery_time TIMESTAMP Lorsque le chargeur automatique a découvert le fichier (Databricks Runtime 16.4 et versions ultérieures)
processed_time TIMESTAMP Lorsque le chargeur automatique a traité le fichier (Databricks Runtime 16.4 et versions ultérieures)
commit_time TIMESTAMP Lorsque le fichier a été validé sur le 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, DELETEou NULL (nécessite cloudFiles.cleanSource)
move_location STRING Chemin d’accès de destination quand cloudFiles.cleanSource est MOVE
ingestion_state STRING État d’ingestion de fichier actuel

Examiner l’état d’ingestion de fichier

Les requêtes suivantes couvrent les scénarios de diagnostic courants.

Recherchez tous les fichiers non traités (le backlog actuel) :

SELECT * FROM cloud_files_state('path/to/checkpoint')
WHERE ingestion_state != 'COMMITTED';

Latence d’ingestion moyenne de calcul (délai de création de fichier à validation) :

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;

Recherchez les fichiers endommagés 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;

Recherchez les fichiers présentant une latence élevée entre la détection et la validation afin d’identifier les goulots 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 consulter la référence SQL complète, consultez la fonction table cloud_files_state.

Surveiller le chargeur automatique dans les pipelines Lakeflow

Databricks recommande d’utiliser des pipelines Lakeflow pour les pipelines de chargeur automatique de production. Pour tirer parti de ses fonctionnalités de supervision intégrées :

  • Stockez le journal des événements des pipelines Lakeflow dans une table Delta afin qu’il puisse être interrogé pour les données d’observabilité. Configurez-le via les paramètres avancés du pipeline ou l’API. Pour plus d’informations, consultez le journal des événements pipeline.

  • Structurez votre pipeline pour l’observabilité. Un pipeline Auto Loader bien structuré dans les pipelines Lakeflow comprend une {table}_source vue (la définition de la source Auto Loader), une {table}_bronze table de streaming (ingestion de données brutes avec les colonnes _rescued_data et _corrupt_record), une corrupt_records_sink qui place en quarantaine les lignes contenant des données non analysables, et une {table} vue propre pour l’exploitation en aval.

  • Définissez les attentes sur vos tables de streaming bronze pour détecter la dérive de schéma et les corruptions de données. _rescued_data IS NULL détecte les modifications inattendues du schéma et _corrupt_record IS NULL détecte les données non exploitables. Les pipelines Lakeflow évaluent ces attentes à mesure que les données arrivent et génèrent une piste d’observabilité. Vous pouvez configurer des règles pour émettre un avertissement, supprimer des lignes ou mettre le pipeline en échec.

Après avoir créé la event_log_raw vue de votre pipeline, utilisez les requêtes suivantes pour les métriques spécifiques au chargeur automatique.

Surveiller le débit 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 des 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ésumez les violations attendues pour détecter la dérive de schéma et les données endommagées :

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 la surveillance des pipelines Lakeflow, consultez Surveiller les pipelines et Journal des événements du pipeline.

Surveiller le chargeur automatique avec Structured Streaming

Lorsque vous exécutez Auto Loader en dehors des pipelines Lakeflow, utilisez les approches suivantes de surveillance de Structured Streaming.

  • Implémentez un StreamingQueryListener pour collecter des métriques spécifiques à Auto Loader pour chaque lot en lisant dans source.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())

Note

La logique de traitement dans les écouteurs peut ralentir le traitement des requêtes. Limitez les traitements dans les fonctions de rappel d’écouteur et évitez d’y effectuer des écritures synchrones vers des systèmes externes ; à la place, émettez une télémétrie légère de façon asynchrone ou confiez les métriques à une tâche distincte chargée de leur persistance.

  • Utilisez numInputRows, inputRowsPerSecondet processedRowsPerSecond de la progression de la source pour calculer le débit ( fichiers par seconde et lignes par seconde pour chaque lot).

  • Pour calculer la latence d’ingestion, comparez create_time et commit_time à partir de cloud_files_state() pour la latence de bout en bout. Pour la latence de traitement, utilisez la décomposition durationMs (par exemple, latestOffset, addBatchet d’autres phases de traitement signalées) afin d’identifier le goulot d’étranglement.

  • Utilisez df.observe() pour définir des métriques de qualité de données inline directement sur le DataFrame de streaming. Les métriques sont visibles dans les événements de progression StreamingQueryListener, sous observedMetrics.

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")
)
  • Permet .queryName() d’attribuer un nom unique à chaque flux, ce qui facilite la distinction entre les flux de chargeur automatique dans l’onglet Streaming de l’interface utilisateur Spark et dans les tableaux de bord de surveillance.

Pour obtenir la référence complète sur la surveillance de Structured Streaming, consultez Surveiller les requêtes Structured Streaming sur Azure 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 de chargeur automatique. Ce tableau affiche certaines 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, validation et horodatages d’archivage par fichier
Journal des événements des pipelines Lakeflow Historique des exécutions de pipeline, métriques de flux par lot et résultats des attentes de qualité des données
Tables de sortie du pipeline Nombres de lignes et volume de données écrites par table ingérée

Vous pouvez ensuite agréger des données d’observabilité dans des tables dédiées qui servent de base pour les tableaux de bord et les alertes :

  • Résumez les états d’exécution du pipeline (réussite ou échec) au fil du temps, dérivés d’événements event_type = 'update_progress' .
  • Métriques agrégées d’ingestion de fichiers (taille de la file d’attente, débit, latence par lot), calculées à partir des événements cloud_files_state() et event_type = 'flow_progress'.
  • Développez des statistiques de table à l’aide du nombre de lignes et du volume de données par table, dérivées num_output_rows du journal des événements.
  • Collectez les informations de débogage à partir des journaux d’erreurs détaillés et des violations d’attentes pour chaque mise à jour, dérivées d’événements event_type = 'flow_progress' dans lesquels data_quality est renseigné.

Ces tables agrégées peuvent alimenter un tableau de bord IA/BI et des alertes SQL. Les panneaux de tableau de bord recommandés incluent la chronologie de l’état de l’exécution du pipeline, la tendance du backlog d’ingestion, la tendance du débit, la distribution de latence d’ingestion, les métriques de qualité des données, les événements d’évolution du schéma et l’état d’archivage des fichiers.

Surveiller les événements d’évolution du schéma

Utilisez les approches suivantes pour détecter les modifications de schéma à mesure qu’elles se produisent.

  • Des valeurs non nulles dans _rescued_data dans le nombre de violations d’attentes indiquent une dérive du schéma. Interrogez le journal d’événements pour les cas où failed_records > 0 sur l’attente no rescued data.
  • Les modifications apportées au _schemas répertoire à l’intérieur du point de contrôle configuré cloudFiles.schemaLocation (ou à l’intérieur du point de contrôle uniquement lorsque l’emplacement du schéma n’est pas défini séparément) indiquent que l’évolution du schéma s’est produite. Vous pouvez scruter ce répertoire depuis une tâche de surveillance distincte.
  • Ne traitez pas un événement onQueryTerminated suivi de onQueryStarted pour le même nom de flux comme une preuve suffisante de l’évolution du schéma à lui seul. Les flux redémarrent pour de nombreuses raisons (redémarrages du cluster, déploiements de code, erreurs de stockage temporaires). Corrélez les redémarrages avec des signaux indépendants — _schemas modifications du répertoire ou _rescued_data violations des attentes — avant de conclure qu’une évolution du schéma a eu lieu.
  • Permet _metadata.file_path d’identifier les fichiers qui ont introduit des modifications de schéma. Joignez ces données à cloud_files_state() sur le champ path afin de corréler les modifications du schéma à des fichiers et lots spécifiques.

Utilisez cet exemple de requête pour détecter la dérive de schéma récente via des violations d’attente :

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;

Configurer des alertes pour les problèmes courants

Utilisez des alertes Databricks SQL ou des 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 backlog croissant et peut être utilisé comme base pour une alerte Databricks SQL. Planifiez l’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 la détecter Quand alerter
Backlog en augmentation numFilesOutstanding tendance à la hausse Augmentation soutenue sur plusieurs lots
Flux bloqué Aucun événement de progression Aucun événement pour N minutes (en fonction de l’intervalle de déclencheur attendu)
Latence d’ingestion élevée commit_time - create_time Dépasse le seuil défini par votre SLA
Dégradation de la qualité des données Taux d’échec attendu Pourcentage croissant de lignes ne répondant pas aux attentes
Événement d’évolution du schéma _rescued_data IS NOT NULL Toute valeur non NULL dans le nombre de violations des attentes
Détection lente de fichiers durationMs.latestOffset Plus élevé que la base de référence

Résoudre les problèmes courants

Le tableau suivant décrit les problèmes courants de pipeline du chargeur automatique, leurs causes probables et les actions recommandées pour les résoudre.

Problème Cause potentielle Action recommandée
Le backlog augmente plus rapidement que le traitement Ressources de calcul sous-dimensionnées, déséquilibre des données ou limites de débit bridées Mettez à l’échelle les ressources de calcul, vérifiez s’il y a une asymétrie des données avec l’interface utilisateur Spark et examinez les paramètres maxFilesPerTrigger pour contrôler la taille des lots
Fichiers non détectés Événements de fichier mal configurés, problèmes d’autorisations ou flux non exécutés dans les 7 jours Vérifiez les autorisations d’emplacement externe, vérifiez la configuration des événements de fichier dans l’interface utilisateur du catalogue Unity et vérifiez que le flux s’exécute au moins tous les 7 jours pour éviter l’expiration de l’état RocksDB
Le lancement du flux prend trop de temps Téléchargement volumineux de l’état du point de contrôle (RocksDB) Mise à niveau vers Databricks Runtime 15.3 et versions ultérieures pour le chargement d’état asynchrone, ce qui réduit le temps de démarrage d’environ 90%
Traitement des fichiers dupliqués Paramètres agressifs cloudFiles.maxFileAge ou altération des points de contrôle Utilisez un paramètre conservateur de type maxFileAge (au moins 90 jours), vérifiez l’intégrité du point de contrôle et évitez les stratégies de cycle de vie appliquées au stockage des points de contrôle
Évolution du schéma entraînant des redémarrages du pipeline Modifications fréquentes ou incompatibles du schéma Examinez schemaEvolutionMode, passez à addNewColumnsWithTypeWidening pour les promotions de types, ou utilisez le type Variant pour les schémas hautement dynamiques
Accumulation de données corrompues dans le récepteur Problèmes de qualité des données sources Vérifiez le récepteur de quarantaine _corrupt_record afin d’identifier des schémas récurrents, examinez la manière dont les données sources sont générées et envisagez d’ajouter une validation en amont
discovery_time et commit_time non renseignés S’exécute sur Databricks Runtime version antérieure à 18.2 sans cleanSource Mettre à niveau vers Databricks Runtime 18.2 et versions ultérieures ou activer cloudFiles.cleanSource sur Databricks Runtime 16.4 à 18.1

Pour plus d’informations sur la résolution des problèmes, consultez faq sur le chargeur automatique.