Remarque
L’accès à cette page requiert une autorisation. Vous pouvez essayer de vous connecter ou de modifier des répertoires.
L’accès à cette page requiert une autorisation. Vous pouvez essayer de modifier des répertoires.
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_timeetcommit_timesont disponibles automatiquement. Sur Databricks Runtime 16.4 à 18.1, ces champs sont disponibles uniquement quandcloudFiles.cleanSourceils sont activés. -
Databricks Runtime 16.4 et versions ultérieures avec
cloudFiles.cleanSourceactivé :archive_time,archive_modeetmove_locationsont 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
_metadatacolonne. Capturez au minimumfile_pathetfile_modification_time. Consultez Colonne métadonnées de fichier. - Activez
_rescued_dataet_corrupt_recordcolonnes.
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.
-
numFilesOutstandingen 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}_sourcevue (la définition de la source Auto Loader), une{table}_bronzetable de streaming (ingestion de données brutes avec les colonnes_rescued_dataet_corrupt_record), unecorrupt_records_sinkqui 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 NULLdétecte les modifications inattendues du schéma et_corrupt_record IS NULLdé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
StreamingQueryListenerpour collecter des métriques spécifiques à Auto Loader pour chaque lot en lisant danssource.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,inputRowsPerSecondetprocessedRowsPerSecondde 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_timeetcommit_timeà partir decloud_files_state()pour la latence de bout en bout. Pour la latence de traitement, utilisez la décompositiondurationMs(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 progressionStreamingQueryListener, 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")
)
- 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()etevent_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_rowsdu 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 lesquelsdata_qualityest 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_datadans 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 > 0sur l’attenteno rescued data. - Les modifications apportées au
_schemasré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
onQueryTerminatedsuivi deonQueryStartedpour 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 —_schemasmodifications du répertoire ou_rescued_dataviolations des attentes — avant de conclure qu’une évolution du schéma a eu lieu. - Permet
_metadata.file_pathd’identifier les fichiers qui ont introduit des modifications de schéma. Joignez ces données àcloud_files_state()sur le champpathafin 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.