Aller au contenu principal

Observabilité dans Databricks pour les Jobs, LakeFlow Pipelines et Lakeflow Connect

Le monitoring des performances, du coût et de la santé de vos applications de streaming est essentiel pour créer des pipelines ETL fiables et efficaces. Databricks offre un riche ensemble de fonctionnalités d'observabilité sur les Jobs, les Lakeflow pipelines et Lakeflow Connect pour aider à diagnostiquer les goulots d'étranglement, à optimiser les performances et à gérer l'utilisation des ressources et les coûts.

Ces bonnes pratiques couvrent les domaines suivants :

  • Mesures clés des performances de streaming
  • Schémas des logs d'événements et exemples de requêtes.
  • Monitoring des query de streaming
  • Observabilité des coûts à l'aide des tables système
  • Exportation des logs et des métriques vers des outils externes.

Métriques clés pour l'observabilité du streaming

Lors de l'exploitation de pipelines de streaming, surveillez les métriques clés suivantes :

Métriques

Objectif

Rétroaction

Surveille le nombre de fichiers et les décalages (tailles). Aide à identifier les goulots d'étranglement et s'assure que le système peut gérer les données entrantes sans prendre de retard.

throughput

Suit le nombre de messages traités par micro-batch. Évaluez l’efficacité du pipeline et assurez-vous qu’il suit le rythme de l’ingestion des données.

Durée

Mesure la durée moyenne d'un micro-batch. Indique la vitesse de traitement et permet de régler les intervalles de batch.

Latence

Indique le nombre d'enregistrements/messages traités au fil du temps. Aide à comprendre les retards des pipelines de bout en bout et à optimiser pour des latences plus faibles.

Utilisation des clusters

Reflète l'utilisation du CPU et de la mémoire (%). Assure une utilisation efficace des ressources et aide à monter en charge les clusters pour répondre aux demandes de traitement.

Réseau

Mesure les données transférées et reçues. Utile pour identifier les goulets d'étranglement du réseau et améliorer les performances de transfert de données.

Point de contrôle

Identifie les données traitées et les décalages. Assure la cohérence et permet la tolérance aux pannes pendant les défaillances.

Coût

Affiche les coûts horaires, quotidiens et mensuels d’une application de streaming. Aide à la budgétisation et à l'optimisation des Ressources.

Traçabilité

Affiche les datasets et les couches créés dans l'application de streaming. Facilite la transformation, le suivi, l'assurance qualité et le debugging des données.

Métriques

Objectif

Rétroaction

Surveille le nombre de fichiers et les décalages (tailles). Aide à identifier les goulots d'étranglement et s'assure que le système peut gérer les données entrantes sans prendre de retard.

throughput

Suit le nombre de messages traités par micro-batch. Évaluez l’efficacité du pipeline et assurez-vous qu’il suit le rythme de l’ingestion des données.

Durée

Mesure la durée moyenne d'un micro-batch. Indique la vitesse de traitement et permet de régler les intervalles de batch.

Latence

Indique le nombre d'enregistrements/messages traités au fil du temps. Aide à comprendre les retards des pipelines de bout en bout et à optimiser pour des latences plus faibles.

Utilisation des clusters

Reflète l'utilisation du CPU et de la mémoire (%). Assure une utilisation efficace des ressources et aide à monter en charge les clusters pour répondre aux demandes de traitement.

Réseau

Mesure les données transférées et reçues. Utile pour identifier les goulets d'étranglement du réseau et améliorer les performances de transfert de données.

Point de contrôle

Identifie les données traitées et les décalages. Assure la cohérence et permet la tolérance aux pannes pendant les défaillances.

Coût

Affiche les coûts horaires, quotidiens et mensuels d’une application de streaming. Aide à la budgétisation et à l'optimisation des Ressources.

Traçabilité

Affiche les datasets et les couches créés dans l'application de streaming. Facilite la transformation, le suivi, l'assurance qualité et le debugging des données.

Logs et métriques de cluster

Les logs de clusters Databricks et les métriques fournissent des insights détaillés sur les performances et l'utilisation des clusters. Ces Logs et métriques incluent des informations sur le processeur, la mémoire, les E/S disque, le trafic réseau et d’autres métriques système. Le monitoring de ces métriques est crucial pour optimiser les performances des clusters, gérer efficacement les ressources et résoudre les problèmes.

Les logs de clusters Databricks et les métriques offrent des informations détaillées sur les performances des clusters et l'utilisation des ressources. Celles-ci incluent l'utilisation du CPU et de la mémoire, les E/S disque et le trafic réseau. Le monitoring de ces métriques est essentielle pour :

  • Optimisation des performances des clusters.
  • Gérer les ressources efficacement.
  • Dépannage des problèmes opérationnels.

Les métriques peuvent être exploitées via l'interface utilisateur de Databricks ou exportées vers des outils de monitoring personnels. Voir Notebook exemple : Datadog metrics.

Spark UI

Le Spark UI affiche des informations détaillées sur la progression des jobs et des étapes, y compris le nombre de tâches terminées, en attente et échouées. Cela vous aide à comprendre le flux d'exécution et à identifier les goulots d'étranglement.

Pour les applications de streaming, la tab Streaming affiche des métriques telles que le taux d'entrée, le taux de traitement et la durée du batch. Cela vous aide à surveiller les performances de vos Jobs de streaming et à identifier tout problème d'ingestion ou de traitement des données.

Consultez Debugging avec la Spark UI pour plus d'informations.

Indicateurs de compute

Les métriques de compute vous aident à comprendre l'utilisation des clusters. Lorsque votre Job s'exécute, vous pouvez voir comment il monte en charge et comment vos ressources sont affectées. Vous pourrez détecter une pression de la mémoire qui pourrait entraîner des erreurs OOM ou une pression du CPU qui pourrait provoquer des retards importants. Voici les métriques spécifiques que vous verrez :

  • Répartition de la charge du serveur : Utilisation du CPU de chaque nœud au cours de la dernière minute.
  • Utilisation du CPU : Le pourcentage de temps passé par le CPU dans divers modes (par exemple, utilisateur, système, inactif et iowait).
  • Utilisation de la mémoire : utilisation totale de la mémoire par chaque mode (par exemple, utilisée, libre, tampon et mise en cache).
  • Utilisation du swap de mémoire : utilisation totale du swap de mémoire.
  • Espace libre sur le système de fichiers : Utilisation totale du système de fichiers par chaque point de montage.
  • Throughput réseau : Le nombre d'octets reçus et transmis par chaque appareil via le réseau.
  • Nombre de nœuds actifs : Le nombre de nœuds actifs à chaque Timestamp pour le compute donné.

Voir les graphiques de métriques matérielles pour plus d'informations.

Tables système

monitoring des coûts

Les tables système Databricks offrent une approche structurée pour surveiller les coûts et les performances des Jobs. Ces tables comprennent :

  • Détails de l'exécution du Job.
  • Utilisation des ressources.
  • Coûts associés.

Utilisez ces tables pour comprendre la santé opérationnelle et l'impact financier.

Exigences

Pour utiliser les tables système pour le monitoring des coûts :

  • Un administrateur de compte doit activer le system.lakeflow schema.
  • Les utilisateurs doivent soit :
    • Soyez à la fois administrateur du métastore et administrateur du compte, ou
    • Disposez des autorisations USE et SELECT sur les schémas système.

Exemple de requête : Jobs les plus coûteux (30 derniers jours)

Cette query identifie les Jobs les plus coûteux au cours des 30 derniers jours, contribuant à l'analyse des coûts et à l'optimisation.

SQL
WITH list_cost_per_job AS (
SELECT
t1.workspace_id,
t1.usage_metadata.job_id,
COUNT(DISTINCT t1.usage_metadata.job_run_id) AS runs,
SUM(t1.usage_quantity * list_prices.pricing.default) AS list_cost,
FIRST(identity_metadata.run_as, true) AS run_as,
FIRST(t1.custom_tags, true) AS custom_tags,
MAX(t1.usage_end_time) AS last_seen_date
FROM system.billing.usage t1
INNER JOIN system.billing.list_prices list_prices ON
t1.cloud = list_prices.cloud AND
t1.sku_name = list_prices.sku_name AND
t1.usage_start_time >= list_prices.price_start_time AND
(t1.usage_end_time <= list_prices.price_end_time OR list_prices.price_end_time IS NULL)
WHERE
t1.billing_origin_product = "JOBS"
AND t1.usage_date >= CURRENT_DATE() - INTERVAL 30 DAY
GROUP BY ALL
),
most_recent_jobs AS (
SELECT
*,
ROW_NUMBER() OVER(PARTITION BY workspace_id, job_id ORDER BY change_time DESC) AS rn
FROM
system.lakeflow.jobs QUALIFY rn=1
)
SELECT
t2.name,
t1.job_id,
t1.workspace_id,
t1.runs,
t1.run_as,
SUM(list_cost) AS list_cost,
t1.last_seen_date
FROM list_cost_per_job t1
LEFT JOIN most_recent_jobs t2 USING (workspace_id, job_id)
GROUP BY ALL
ORDER BY list_cost DESC

LakeFlow Pipelines

Le journal des événements des LakeFlow Pipelines consigne un enregistrement complet de tous les événements de pipeline, y compris :

  • Logs d'audit.
  • Contrôle de la qualité des données.
  • Progression du pipeline.
  • data lineage.

Le journal des événements est automatiquement activé pour tous les LakeFlow Pipelines et est accessible via :

  • Interface utilisateur du pipeline : afficher les Logs directement.
  • API Pipelines : Accès programmatique.
  • Direct query : query la table de Logs d'événements.

Pour plus d'information, consultez le schéma du log d'événements pour les LakeFlow Pipelines.

Exemples de requêtes

Ces exemples de requêtes aident à surveiller les performances et la santé des pipelines en fournissant des métriques clés telles que la durée de batch, le throughput, la contre-pression et l'utilisation des ressources.

Durée moyenne du batch

Cette query calcule la durée moyenne des batchs traités par le pipeline.

SQL
SELECT
(max_t - min_t) / batch_count as avg_batch_duration_seconds,
batch_count,
min_t,
max_t,
date_hr,
message
FROM
-- /60 for minutes
(
SELECT
count(*) as batch_count,
unix_timestamp(
min(timestamp)
) as min_t,
unix_timestamp(
max(timestamp)
) as max_t,
date_format(timestamp, 'yyyy-MM-dd:HH') as date_hr,
message
FROM
event_log
WHERE
event_type = 'flow_progress'
AND level = 'METRICS'
GROUP BY
date_hr,
message
)
ORDER BY
date_hr desc

Throughput moyen

Cette query calcule le throughput moyen du pipeline en termes de lignes traitées par seconde.

SQL
SELECT
(max_t - min_t) / total_rows as avg_throughput_rps,
total_rows,
min_t,
max_t,
date_hr,
message
FROM
-- /60 for minutes
(
SELECT
sum(
details:flow_progress:metrics:num_output_rows
) as total_rows,
unix_timestamp(
min(timestamp)
) as min_t,
unix_timestamp(
max(timestamp)
) as max_t,
date_format(timestamp, 'yyyy-MM-dd:HH') as date_hr,
message
FROM
event_log
WHERE
event_type = 'flow_progress'
AND level = 'METRICS'
GROUP BY
date_hr,
message
)
ORDER BY
date_hr desc

Contre-pression

Cette query mesure la contre-pression du pipeline en vérifiant le backlog de données.

SQL
SELECT
timestamp,
DOUBLE(
details:flow_progress:metrics:backlog_bytes
) AS backlog_bytes,
DOUBLE(
details:flow_progress:metrics:backlog_files
) AS backlog_files
FROM
event_log
WHERE
event_type = 'flow_progress'

Utilisation des clusters et des emplacements

Cette requête fournit des insights sur l'utilisation des clusters ou des emplacements utilisés par le pipeline.

SQL
SELECT
date_trunc("hour", timestamp) AS hour,
AVG (
DOUBLE (
details:cluster_resources:num_task_slots
)
) AS num_task_slots,
AVG (
DOUBLE (
details:cluster_resources:avg_num_task_slots
)
) AS avg_num_task_slots,
AVG (
DOUBLE (
details:cluster_resources:num_executors
)
) AS num_executors,
AVG (
DOUBLE (
details:cluster_resources:avg_task_slot_utilization
)
) AS avg_utilization,
AVG (
DOUBLE (
details:cluster_resources:avg_num_queued_tasks
)
) AS queue_size
FROM
event_log
WHERE
details : cluster_resources : avg_num_queued_tasks IS NOT NULL
AND origin.update_id = '${latest_update_id}'
GROUP BY
1;

Jobs

Vous pouvez surveiller les queries de streaming dans les jobs via le Streaming Query Listener.

Joignez un écouteur à la session Spark pour activer le Streaming Query Listener dans Databricks. Ce moniteur surveille la progression et les métriques de vos queries de streaming. Il peut être utilisé pour envoyer des métriques à des outils de monitoring externes ou pour les enregistrer dans des Logs pour une analyse approfondie.

Exemple : Exporter des métriques vers des outils de monitoring externes.

remarque

Ceci est disponible dans Databricks Runtime 11.3 LTS et versions ultérieures pour Python et Scala.

Vous pouvez exporter les métriques de streaming vers des services externes pour les alertes ou la création de tableaux de bord en utilisant l'interface StreamingQueryListener.

Voici un exemple basique de la façon d'implémenter un écouteur :

Python
from pyspark.sql.streaming import StreamingQueryListener

class MyListener(StreamingQueryListener):
def onQueryStarted(self, event):
print("Query started: ", event.id)

def onQueryProgress(self, event):
print("Query made progress: ", event.progress)

def onQueryTerminated(self, event):
print("Query terminated: ", event.id)

spark.streams.addListener(MyListener())

Exemple : Utilisez le récepteur de requêtes au sein de Databricks.

Voici un exemple de log d'événements StreamingQueryListener pour une query de streaming de Kafka vers Delta Lake :

JSON
{
"id": "210f4746-7caa-4a51-bd08-87cabb45bdbe",
"runId": "42a2f990-c463-4a9c-9aae-95d6990e63f4",
"timestamp": "2024-05-15T21:57:50.782Z",
"batchId": 0,
"batchDuration": 3601,
"numInputRows": 20,
"inputRowsPerSecond": 0.0,
"processedRowsPerSecond": 5.55401277422938,
"durationMs": {
"addBatch": 1544,
"commitBatch": 686,
"commitOffsets": 27,
"getBatch": 12,
"latestOffset": 577,
"queryPlanning": 105,
"triggerExecution": 3600,
"walCommit": 34
},
"stateOperators": [
{
"operatorName": "symmetricHashJoin",
"numRowsTotal": 20,
"numRowsUpdated": 20,
"allUpdatesTimeMs": 473,
"numRowsRemoved": 0,
"allRemovalsTimeMs": 0,
"commitTimeMs": 277,
"memoryUsedBytes": 13120,
"numRowsDroppedByWatermark": 0,
"numShufflePartitions": 5,
"numStateStoreInstances": 20,
"customMetrics": {
"loadedMapCacheHitCount": 0,
"loadedMapCacheMissCount": 0,
"stateOnCurrentVersionSizeBytes": 5280
}
}
],
"sources": [
{
"description": "KafkaV2[Subscribe[topic-1]]",
"numInputRows": 10,
"inputRowsPerSecond": 0.0,
"processedRowsPerSecond": 2.77700638711469,
"metrics": {
"avgOffsetsBehindLatest": "0.0",
"estimatedTotalBytesBehindLatest": "0.0",
"maxOffsetsBehindLatest": "0",
"minOffsetsBehindLatest": "0"
}
},
{
"description": "DeltaSource[file:/tmp/spark-1b7cb042-bab8-4469-bb2f-733c15141081]",
"numInputRows": 10,
"inputRowsPerSecond": 0.0,
"processedRowsPerSecond": 2.77700638711469,
"metrics": {
"numBytesOutstanding": "0",
"numFilesOutstanding": "0"
}
}
]
}

Pour plus d'exemples, voir : Exemples.

Mesures de progression des query

Les métriques de progression des requêtes sont essentielles pour le monitoring des performances et de la santé de vos requêtes de streaming. Ces métriques comprennent le nombre de lignes d'entrée, les taux de traitement et diverses durées liées à l'exécution de la query. Vous pouvez observer ces métriques en attachant un StreamingQueryListener à la session Spark. Le récepteur émet des événements contenant ces métriques à la fin de chaque époque de streaming.

Par exemple, vous pouvez accéder aux métriques à l'aide de la map StreamingQueryProgress.observedMetrics dans la méthode onQueryProgress du listener. Cela vous permet de suivre et d'analyser les performances de vos queries de streaming en temps réel.

Python
class MyListener(StreamingQueryListener):
def onQueryProgress(self, event):
print("Query made progress: ", event.progress.observedMetrics)