Monitoring des requêtes Structured Streaming sur Databricks
Databricks fournit un monitoring intégré pour les applications Structured Streaming via la Spark UI sous l'onglet Streaming .
Distinguer les Structured Streaming queries dans la Spark UI
Donnez à vos Streams un nom de query unique en ajoutant .queryName(<query-name>) à votre code writeStream pour distinguer facilement les métriques qui appartiennent à quel Stream dans la Spark UI.
Transférer les métriques Structured Streaming vers des services externes
Les métriques de streaming peuvent être transférées vers des services externes pour les cas d'usage d'alerte ou de création de tableaux de bord en utilisant l'interface Streaming Query Listener d'Apache Spark. Dans Databricks Runtime 11.3 LTS et versions ultérieures, StreamingQueryListener est disponible en Python et Scala.
Les limitations suivantes s'appliquent aux workloads utilisant les modes d'accès compute compatibles Unity Catalog :
StreamingQueryListenernécessite Databricks Runtime 15.1 ou supérieur pour utiliser les informations d'identification ou interagir avec les objets gérés par Unity Catalog sur les compute avec le mode d'accès dédié.StreamingQueryListenernécessite Databricks Runtime 16.1 ou une version ultérieure pour les charges de travail Scala configurées en mode d’accès standard (anciennement mode d’accès partagé).
La latence de traitement avec les écouteurs peut affecter de manière significative les vitesses de traitement des query. Il est conseillé de limiter la logique de traitement dans ces écouteurs et d'opter pour l'écriture dans des systèmes à réponse rapide comme Kafka pour l'efficacité.
Si la query ne dispose d'aucune donnée disponible sur la source et attend de nouvelles données, un message onQueryIdle est transmis au récepteur de query en streaming. Un message onQueryProgress n'est livré qu'à la fin du batch de query en streaming. Si la query a traité des données pendant longtemps, il est possible que ni les événements onQueryIdle ni les événements onQueryProgress ne soient envoyés, mais la query est toujours saine et continue de traiter les données.
Le code suivant fournit des exemples de base de la syntaxe pour implémenter un listener :
- Scala
- Python
import org.apache.spark.sql.streaming.StreamingQueryListener
import org.apache.spark.sql.streaming.StreamingQueryListener._
val myListener = new StreamingQueryListener {
/**
* Called when a query is started.
* @note This is called synchronously with
* [[org.apache.spark.sql.streaming.DataStreamWriter `DataStreamWriter.start()`]].
* `onQueryStart` calls on all listeners before
* `DataStreamWriter.start()` returns the corresponding [[StreamingQuery]].
* Do not block this method, as it blocks your query.
*/
def onQueryStarted(event: QueryStartedEvent): Unit = {}
/**
* Called when there is some status update (ingestion rate updated, etc.)
*
* @note This method is asynchronous. The status in [[StreamingQuery]] returns the
* latest status, regardless of when this method is called. The status of [[StreamingQuery]]
* may change before or when you process the event. For example, you may find [[StreamingQuery]]
* terminates when processing `QueryProgressEvent`.
*/
def onQueryProgress(event: QueryProgressEvent): Unit = {}
/**
* Called when the query is idle and waiting for new data to process.
*/
def onQueryIdle(event: QueryProgressEvent): Unit = {}
/**
* Called when a query is stopped, with or without error.
*/
def onQueryTerminated(event: QueryTerminatedEvent): Unit = {}
}
class MyListener(StreamingQueryListener):
def onQueryStarted(self, event):
"""
Called when a query is started.
Parameters
----------
event: :class:`pyspark.sql.streaming.listener.QueryStartedEvent`
The properties are available as the same as Scala API.
Notes
-----
This is called synchronously with
meth:`pyspark.sql.streaming.DataStreamWriter.start`,
that is, ``onQueryStart`` will be called on all listeners before
``DataStreamWriter.start()`` returns the corresponding
:class:`pyspark.sql.streaming.StreamingQuery`.
Do not block in this method as it will block your query.
"""
pass
def onQueryProgress(self, event):
"""
Called when there is some status update (ingestion rate updated, etc.)
Parameters
----------
event: :class:`pyspark.sql.streaming.listener.QueryProgressEvent`
The properties are available as the same as Scala API.
Notes
-----
This method is asynchronous. The status in
:class:`pyspark.sql.streaming.StreamingQuery` returns the
most recent status, regardless of when this method is called. The status
of :class:`pyspark.sql.streaming.StreamingQuery`.
may change before or when you process the event.
For example, you may find :class:`StreamingQuery`
terminates when processing `QueryProgressEvent`.
"""
pass
def onQueryIdle(self, event):
"""
Called when the query is idle and waiting for new data to process.
"""
pass
def onQueryTerminated(self, event):
"""
Called when a query is stopped, with or without error.
Parameters
----------
event: :class:`pyspark.sql.streaming.listener.QueryTerminatedEvent`
The properties are available as the same as Scala API.
"""
pass
my_listener = MyListener()
Définir des métriques observables dans le Structured Streaming
Les métriques observables sont des fonctions d'agrégation arbitraires nommées qui peuvent être définies sur une query (DataFrame). Dès que l'exécution d'un DataFrame atteint un point de complétion (c'est-a-dire, termine une query batch ou atteint une époque de streaming), un événement nommé est émis qui contient les métriques pour les données traitées depuis le dernier point de complétion.
Vous pouvez observer ces métriques en attachant un écouteur à la session Spark. Le programme d'écoute dépend du mode d'exécution :
-
Mode batch : Utilisez
QueryExecutionListener.QueryExecutionListenerest appelée lorsque la query est terminée. Accédez aux métriques à l'aide de la carteQueryExecution.observedMetrics. -
Streaming ou micro-batch :
StreamingQueryListenerutilisez.StreamingQueryListenerest appelée lorsque la query de streaming termine une époque. Accédez aux métriques à l'aide de la carteStreamingQueryProgress.observedMetrics. Databricks ne prend pas en charge le modecontinuousTrigger pour le streaming.
Par exemple :
- Scala
- Python
// Observe row count (rc) and error row count (erc) in the streaming Dataset
val observed_ds = ds.observe("my_event", count(lit(1)).as("rc"), count($"error").as("erc"))
observed_ds.writeStream.format("...").start()
// Monitor the metrics using a listener
spark.streams.addListener(new StreamingQueryListener() {
override def onQueryProgress(event: QueryProgressEvent): Unit = {
event.progress.observedMetrics.get("my_event").foreach { row =>
// Trigger if the number of errors exceeds 5 percent
val num_rows = row.getAs[Long]("rc")
val num_error_rows = row.getAs[Long]("erc")
val ratio = num_error_rows.toDouble / num_rows
if (ratio > 0.05) {
// Trigger alert
}
}
}
})
# Observe metric
observed_df = df.observe("metric", count(lit(1)).as("cnt"), count(col("error")).as("malformed"))
observed_df.writeStream.format("...").start()
# Define my listener.
class MyListener(StreamingQueryListener):
def onQueryStarted(self, event):
print(f"'{event.name}' [{event.id}] got started!")
def onQueryProgress(self, event):
row = event.progress.observedMetrics.get("metric")
if row is not None:
if row.malformed / row.cnt > 0.5:
print("ALERT! Ouch! there are too many malformed "
f"records {row.malformed} out of {row.cnt}!")
else:
print(f"{row.cnt} rows processed!")
def onQueryTerminated(self, event):
print(f"{event.id} got terminated!")
# Add my listener.
spark.streams.addListener(MyListener())
Mapper les identifiants de table de métriques Unity Catalog, Delta Lake et Structured Streaming
Les métriques de Structured Streaming utilisent le champ reservoirId à plusieurs endroits pour l'identité unique d'une table Delta Lake utilisée comme source pour une query de streaming.
Le champ reservoirId mappe l'identifiant unique stocké par la table Delta Lake dans les Logs de transactions Delta. Cet ID ne correspond pas à la valeur tableId attribuée par Unity Catalog et affichée dans Catalog Explorer.
Utilisez la syntaxe suivante pour examiner l'identifiant de table pour une table Delta Lake. Ceci fonctionne pour les tables gérées Unity Catalog, les tables externes Unity Catalog et toutes les tables Delta Lake du Hive metastore :
DESCRIBE DETAIL <table-name>
Le champ id affiché dans les résultats est l'identifiant qui correspond au reservoirId dans les métriques de streaming.
Métriques de l'objet StreamingQueryListener
Champs | Description |
|---|---|
| Un ID de query unique qui persiste après les redémarrages. |
| Un ID de query unique pour chaque start/redémarrage. Consultez StreamingQuery.runId(). |
| Le nom de la query spécifié par l'utilisateur. Le nom est nul si aucun nom n'est spécifié. |
| Le Timestamp de l'exécution du micro-batch. |
| Un ID unique pour le batch de données en cours de traitement. En cas de nouvelles tentatives après un échec, un ID de batch donné peut être exécuté plus d'une fois. De même, lorsqu'il n'y a pas de données à traiter, l'ID du batch n'est pas incrémenté. |
| La durée de traitement d'une opération de batch, en millisecondes. |
| Nombre total (toutes sources confondues) d'enregistrements traités dans un trigger. |
| Le taux agrégé (toutes sources confondues) des données entrantes. |
| Le taux agrégé (toutes sources confondues) auquel Spark traite les données. |
StreamingQueryListener définit également les champs suivants qui contiennent des objets que vous pouvez examiner pour les métriques clients et les détails de progression de la source.
Champs | Description |
|---|---|
| Type : |
| Type : |
| Type : |
| Type : |
| Type : |
| Type : |
Objet durationMs
Type d'objet : ju.Map[String, JLong]
Information sur le temps nécessaire pour achever les différentes étapes du processus d'exécution des micro-batchs.
Champs | Description |
|---|---|
| La durée d'exécution du micro-batch. Cela exclut le temps que Spark prend pour planifier le micro-batch. |
| Le temps nécessaire pour récupérer les métadonnées concernant les offsets de la source. |
| Le dernier décalage consommé pour le micro-batch. Cet objet de progression fait référence au temps pris pour récupérer le dernier décalage des sources. |
| Le temps nécessaire pour générer le plan d'exécution. |
| Le temps nécessaire pour planifier et exécuter le micro-lot. |
| Le temps nécessaire pour commit les nouveaux décalages disponibles. |
| Le temps nécessaire pour commit les données écrites dans le sink pendant |
| Le temps nécessaire pour commit le batch dans le log de commit. |
objet eventTime
Type d'objet : ju.Map[String, String]
Informations sur la valeur temporelle de l'événement observée dans les données traitées dans le micro-batch. Ces données sont utilisées par le filigrane pour déterminer comment élaguer l'état en vue du traitement des agrégations avec état définies dans le Job Structured Streaming.
Champs | Description |
|---|---|
| Le temps d'événement moyen observé dans ce Trigger. |
| Le temps d'événement maximal observé dans ce Trigger. |
| Le temps minimal d'événement observé dans ce Trigger. |
| La valeur du filigrane utilisé dans ce trigger. |
objet stateOperators
Type d'objet : Array[StateOperatorProgress]
L'objet stateOperators contient des informations sur les opérations avec état définies dans le Job Structured Streaming et les agrégations qui en découlent.
Pour plus de détails sur les opérateurs d'état de flux, voir Qu'est-ce que le streaming avec état ?.
Champs | Description |
|---|---|
| Le nom de l'opérateur avec état auquel les métriques sont liées, tels que |
| Le nombre total de lignes en état à la suite d'un opérateur avec état ou d'une agrégation. |
| Le nombre total de lignes mises à jour dans l'état à la suite d'un opérateur avec état ou d'une agrégation. |
| Cette métrique n'est actuellement pas mesurable par Spark et il est prévu de la supprimer lors de futures mises à jour. |
| Le nombre total de lignes supprimées de l'état suite à un opérateur avec état ou à une agrégation. |
| Cette métrique n'est actuellement pas mesurable par Spark et il est prévu de la supprimer lors de futures mises à jour. |
| Le temps nécessaire pour commit toutes les mises à jour (ajouts et suppressions) et renvoyer une nouvelle version. |
| Mémoire utilisée par le magasin d'état. |
| Le nombre de lignes considérées comme trop tardives pour être incluses dans une agrégation avec état. Agrégations en streaming uniquement : Le nombre de lignes supprimées après agrégation (pas les lignes d'entrée brutes). Ce nombre n'est pas précis, mais fournit une indication qu'il y a des données tardives qui sont perdues. |
| Le nombre de partitions de brassage pour cet opérateur avec état. |
| L'instance de magasin d'état réelle que l'opérateur a initialisée et maintenue. Pour de nombreux opérateurs avec état, cela correspond au nombre de partitions. Cependant, les jointures Stream-Stream initialisent quatre instances de magasin d'état par partition. |
| Consultez stateOperators.customMetrics dans cette rubrique pour en savoir plus. |
Objet StateOperatorProgress.customMetrics
Type d'objet : ju.Map[String, JLong]
StateOperatorProgress contient un champ, customMetrics, qui comprend les métriques spécifiques à la fonctionnalité que vous utilisez lors de la collecte de ces métriques.
Fonctionnalité | Description |
|---|---|
Magasin d'état RocksDB | |
Magasin d’état HDFS | |
Déduplication de Stream | |
Agrégation de Stream | |
Opérateur de jointure de Stream | |
|
Métriques personnalisées du magasin d'état RocksDB
Informations recueillies auprès de RocksDB capturant des métriques sur ses performances et ses opérations concernant les valeurs avec état qu'il maintient pour le Job Structured Streaming. Pour plus d'informations, consultez Configurer le magasin d'état RocksDB sur Databricks.
Champs | Description |
|---|---|
| Le nombre d'octets copiés tel que suivi par le gestionnaire de fichiers RocksDB. |
| Le temps nécessaire, en millisecondes, pour prendre un instantané de RocksDB natif et l’écrire dans un répertoire local. |
| La durée en millisecondes de la compression (facultatif) pendant le commit du point de contrôle. |
| Le temps de compactage pendant le commit, en millisecondes. |
| Temps, en millisecondes, nécessaire à la synchronisation de l’instantané natif RocksDB vers le stockage externe (l’emplacement de point de contrôle). |
| Le temps en millisecondes pour vider les modifications en mémoire de RocksDB sur le disque local. |
| Le temps en millisecondes d'arrêt des threads worker en arrière-plan dans le cadre du commit de checkpoint, par exemple pour la compaction. |
| Le temps en millisecondes d'application des écritures intermédiaires dans une structure en mémoire ( |
| Le nombre de fichiers copiés, tels que suivis par le gestionnaire de fichiers RocksDB. |
| Le nombre de fichiers réutilisés et suivis par le gestionnaire de fichiers RocksDB. |
| Le nombre d'appels |
| Le temps moyen en nanosecondes pour l'appel natif |
| Le nombre d'accès au cache à partir du cache de blocs dans RocksDB. |
| Le nombre d'échecs du cache de blocs dans RocksDB. |
| La taille de tous les fichiers Static Sorted Table (SST) dans l'instance RocksDB. |
| Le nombre d'octets décompressés lus par les opérations |
| Le nombre total d’octets non compressés écrits par |
| Le nombre total d'octets de données non compressées lus à l'aide d'un itérateur. Certaines opérations avec état (par exemple, le traitement des délais d'attente dans |
| Le nombre d'octets que le processus de compactage lit sur le disque. |
| Le nombre total d'octets que le processus de compaction écrit sur le disque. |
| Le temps en millisecondes pour les compactages RocksDB, y compris les compactages en arrière-plan et le compactage facultatif lancé pendant le commit. |
| Le temps de vidage total, y compris le vidage en arrière-plan. Les opérations de vidage sont des processus par lesquels le |
| La taille en octets des fichiers zip décompressés telle que signalée par le Gestionnaire de fichiers. Le Gestionnaire de fichiers gère l'utilisation et la suppression de l'espace disque des fichiers physiques SST. |
| La version la plus récente de l'instantané RocksDB enregistrée à l'emplacement du point de contrôle. Une valeur de « -1 » indique qu'aucun instantané n'a jamais été enregistré. Étant donné que les instantanés sont spécifiques à chaque instance de magasin d'état, cette métrique s'applique à un ID de partition et à un nom de magasin d'état particuliers. |
| La latence totale des options d'achat et de vente. |
| Le nombre d'appels de mise. |
| Le temps d'attente du processus d'écriture pour que la compaction ou le vidage se termine. |
| Le nombre total d'octets écrits par vidage. |
| L'utilisation de la mémoire pour les blocs épinglés |
| Le nombre de clés internes pour les familles de colonnes internes |
| Le nombre de familles de colonnes externes |
| Le nombre de familles de colonnes internes |
Métriques personnalisées du magasin d'état HDFS
Informations collectées sur les comportements et les opérations du fournisseur de magasin d'état HDFS.
Champs | Description |
|---|---|
| La taille estimée de l'état uniquement sur la version actuelle. |
| Le nombre d'accès réussis au cache sur les états mis en cache dans le fournisseur. |
| Le nombre de manques de cache sur les états mis en cache chez le fournisseur. |
| La dernière version upload de l'instantané pour une instance de magasin d'état spécifique. |
Déduplication des métriques personnalisées
Informations recueillies concernant les comportements et opérations de déduplication.
Champs | Description |
|---|---|
| Le nombre de lignes en double supprimées. |
| Le nombre de lignes d'état lues pendant l'éviction d'état. |
Métriques personnalisées d'agrégation
Informations recueillies concernant les comportements d'agrégation et les opérations.
Champs | Description |
|---|---|
| Le nombre de lignes d'état lues pendant l'éviction d'état. |
Métriques personnalisées de jonction de Stream
Informations collectées concernant les comportements et les opérations de jointure de stream.
Champs | Description |
|---|---|
| Le nombre de valeurs |
métriques personnalisées transformWithState
Informations collectées sur les comportements et Opérations de transformWithState (TWS). Pour plus de détails sur transformWithState, consultez Créer une application avec état personnalisée.
Champs | Description |
|---|---|
| Nombre de millisecondes nécessaires pour traiter tout l'état initial. |
| Nombre de variables d'état de valeur. Également présent pour |
| Nombre de variables d'état de liste. Également présent pour |
| Nombre de variables d'état de carte. Également présent pour |
| Nombre de variables d'état supprimées. Également présent pour |
| Nombre de millisecondes nécessaires au traitement de tous les compteurs. |
| Nombre de minuteurs enregistrés. Également présent pour |
| Nombre de minuteurs supprimés. Également présent pour |
| Nombre de minuteurs expirés. Également présent pour |
| Nombre de variables d'état de la valeur dotées d'un TTL. Également présent pour |
| Nombre de variables d'état de liste avec TTL. Également présent pour |
| Nombre de variables d'état de carte avec TTL. Également présent pour |
| Nombre de valeurs supprimées en raison de l'expiration du TTL. Également présent pour |
| Nombre de valeurs supprimées de manière incrémentielle en raison de l'expiration de la TTL. |
Objets sources
Type d'objet : Array[SourceProgress]
L'objet sources contient des informations et des métriques pour les sources de données en streaming.
Champs | Description |
|---|---|
| Description détaillée de la table de source de données de streaming. |
| Le numéro de décalage de start dans la table de la source de données à partir de laquelle le Job de streaming a démarré. |
| Le dernier décalage traité par le micro-batch. |
| Le décalage le plus récent traité par le micro-batch. |
| Le nombre de lignes d'entrée traitées à partir de cette source. |
| Le taux, en secondes, auquel les données arrivent pour traitement depuis cette source. |
| Le débit auquel Spark traite les données de cette source. |
| Type : |
Databricks fournit l'implémentation d'objets sources suivante :
Pour les champs définis sous la forme sources.<startOffset / endOffset / latestOffset>.* (ou une autre variante), interprétez-les comme l’un des 3 champs possibles (jusqu’à ces limites), contenant tous le champ enfant indiqué :
sources.startOffset.<child-field>sources.endOffset.<child-field>sources.latestOffset.<child-field>
Objet sources Delta Lake
Définitions des mesures personnalisées utilisées pour les sources de données de streaming de table Delta Lake.
Champs | Description |
|---|---|
| La description de la source à partir de laquelle la query de streaming lit. Par exemple : |
| La version de sérialisation avec laquelle ce décalage est encodé. |
| L'ID de la table en cours de lecture. Ceci est utilisé pour détecter une mauvaise configuration lors du redémarrage d'une query. Consultez Mapper les identifiants de table de métriques Unity Catalog, Delta Lake et Structured Streaming. |
| La version de la table qui est en cours de traitement. |
| L'index dans la séquence de |
| Identifie si l'offset actuel marque le start d'une nouvelle query de streaming plutôt que le traitement des changements survenus après le traitement initial des données. Lorsque vous start une nouvelle query, toutes les données présentes dans la table au start sont traitées en premier, puis toutes les nouvelles données qui arrivent. |
| Heure de l'événement enregistrée pour l'ordonnancement par heure de l'événement. L'heure de l'événement des données d'instantané initial qui sont en attente de traitement. Utilisé lors du traitement d’un instantané initial avec un ordre temporel des événements. |
| L'offset le plus récent traité par la query micro-batch. |
| Le nombre de lignes d'entrée traitées à partir de cette source. |
| Le rythme auquel les données arrivent pour traitement depuis cette source. |
| Le débit auquel Spark traite les données de cette source. |
| La taille combinée des fichiers en attente (fichiers suivis par RocksDB). Il s'agit de la métrique de backlog pour Delta et Auto Loader en tant que source de streaming. |
| Le nombre de fichiers en attente à traiter. Il s'agit de la métrique de backlog pour Delta et Auto Loader en tant que source de streaming. |
Objet sources Apache Kafka
Définitions des métriques personnalisées utilisées pour les sources de données de streaming Apache Kafka.
Champs | Description |
|---|---|
| Description détaillée de la source Kafka, spécifiant le sujet Kafka exact à partir duquel la lecture est effectuée. Par exemple : |
| Le numéro de décalage de début dans le sujet Kafka auquel le Job de streaming a démarré. |
| Le dernier offset traité par le micro-batch. Cela pourrait être égal à |
| Le dernier décalage calculé par le micro-batch. Le processus de micro-batching pourrait ne pas traiter tous les offsets en cas de limitation, ce qui entraînerait une différence entre |
| Le nombre de lignes d'entrée traitées à partir de cette source. |
| Le rythme auquel les données arrivent pour traitement depuis cette source. |
| Le débit auquel Spark traite les données de cette source. |
| Le nombre moyen de décalages de la query en streaming par rapport au dernier décalage disponible parmi tous les sujets abonnés. |
| Le nombre estimé d'octets que le processus de query n'a pas consommés à partir des sujets abonnés. |
| Le nombre maximal de décalages que la query en streaming est en retard par rapport au dernier décalage disponible parmi tous les sujets abonnés. |
| Le nombre minimum de décalages que la query de streaming accuse par rapport au dernier décalage disponible parmi tous les sujets abonnés. |
Dans Databricks Runtime 17.1 et versions ultérieures, les derniers offsets Kafka sont récupérés après l'achèvement de chaque micro-batch. Sur les rubriques qui reçoivent continuellement des données, les métriques de backlog peuvent afficher des valeurs non nulles, faibles et persistantes. Il s'agit d'un comportement attendu et n'indique pas que le Stream prend du retard.
Dans Databricks Runtime 17.0 et versions antérieures, les derniers offsets Kafka sont récupérés à l'heure de start du micro-batch. Les métriques de backlog peuvent renvoyer 0 lorsque les queries en streaming consomment systématiquement tous les enregistrements disponibles au start du micro-batch.
Objet source AWS Kinesis
Définitions des métriques personnalisées utilisées pour les sources de données de streaming AWS Kinesis.
Champs | Description |
|---|---|
| La description de la source Kinesis, spécifiant le Kinesis Stream exact à partir duquel la query de streaming est lue. Par exemple : |
| Le nombre moyen de millisecondes de retard d'un consommateur par rapport au début d'un stream. |
| Le nombre maximal de millisecondes pendant lesquelles un consommateur a pris du retard par rapport au début d'un stream. |
| Le nombre minimal de millisecondes qu'un consommateur a pris de retard par rapport au début d'un stream. |
| Le nombre d’octets restants à traiter. Ceci est la métrique de backlog pour Kinesis en tant que source. |
| Nom du Stream Kinesis. |
| ID du shard Kinesis Stream. |
| Le premier numéro de séquence des enregistrements dans une partition Kinesis qui ont été consommés dans un batch donné. |
| Le dernier numéro de séquence des enregistrements consommés à partir d'un shard Kinesis dans un batch donné. |
| Si le fragment Kinesis a été fermé par le Stream Kinesis. |
| Le temps approximatif du retard de la query de streaming par rapport aux dernières données dans le Kinesis stream. |
| Le numéro de séquence du dernier enregistrement consommé, et est utilisé pour la vérification des pertes de données. Notez que |
| Le mode consommateur utilisé pour l'exécution de la query de streaming. Peut être |
| Le nombre de Kinesis Stream traités dans ce micro-batch. |
| Le nombre total de fragments traités dans ce micro-batch. |
| Le nombre de shards fermés traités dans ce micro-batch. |
| Le nombre d'octets traités dans ce micro-batch. |
| Le nombre d'enregistrements traités dans ce micro-batch. |
| Le nombre de consommateurs enregistrés utilisés en mode EFO. |
Pour plus d'information, voir Surveiller les métriques Kinesis.
Métriques source Auto Loader
Définitions des métriques personnalisées utilisées pour les sources de données de streaming Auto Loader.
Champs | Description |
|---|---|
| La position actuelle dans la séquence des fichiers traités, dans l'ordre de leur découverte. |
| La version d'implémentation de la source cloudFiles. |
| L'heure de start de l'opération de remplissage la plus récente. |
| L'heure de fin de la plus récente opération de backfill. |
| Le dernier chemin d'entrée fourni par l'utilisateur du Stream avant que le Stream ne soit redémarré. |
| Le nombre de fichiers en attente |
| La taille (en octets) des fichiers en attente |
| La taille approximative de la file d'attente de messages. Uniquement lorsque l'option cloudFiles.useNotifications est activée. |
| Le nombre de lignes d'entrée traitées à partir de cette source. Pour le format source |
Indicateurs des sources PubSub
Définitions des métriques personnalisées utilisées pour les sources de données de streaming PubSub. Pour plus de détails sur le monitoring des sources de streaming PubSub, consultez Monitor Pub/Sub streaming metrics.
Champs | Description |
|---|---|
| La version d'implémentation avec laquelle ce décalage est encodé. |
| Le numéro de séquence persistant qui est en cours de traitement. |
| La plus grande époque de récupération en cours de traitement. |
| Le nombre d’enregistrements disponibles pour le traitement dans le backlog actuel. |
| La taille totale en octets des données non traitées dans l'arriéré actuel. |
| Le nombre total d'enregistrements dupliqués traités par le Stream depuis son start. |
Mesures des sources Pulsar
Définitions pour les métriques personnalisées utilisées pour les sources de données de streaming Pulsar.
Champs | Description |
|---|---|
| Le nombre de lignes traitées dans le micro-batch actuel. |
| Le nombre total d'octets traités dans le micro-batch actuel. |
objet sink
Type d'objet : SinkProgress
Champs | Description |
|---|---|
| La description du puits, détaillant l'implémentation spécifique du puits utilisée. |
| Le nombre de lignes de sortie. Différents types de récepteurs peuvent avoir des comportements ou des restrictions différents pour les valeurs. Voir les types pris en charge spécifiques. |
|
|
Actuellement, Databricks propose deux implémentations d'objets sink spécifiques :
Type de puits | Détails |
|---|---|
table Delta Lake | Consultez l’ objet Delta sink. |
Sujet Apache Kafka | Consultez objet récepteur Kafka. |
Le champ sink.metrics se comporte de la même manière pour les deux variantes de l'objet sink.
Objet récepteur Delta Lake
Champs | Description |
|---|---|
| La description du récepteur Delta, détaillant l'implémentation spécifique du récepteur Delta utilisée. Par exemple : |
| Le nombre de lignes est toujours |
Objet récepteur Apache Kafka
Champs | Description |
|---|---|
| La description du récepteur Kafka vers lequel la query de streaming écrit, détaillant l'implémentation spécifique du récepteur Kafka utilisée. Par exemple : |
| Le nombre de lignes qui ont été écrites dans la table de sortie ou le puits de données dans le cadre du micro-batch. Dans certaines situations, cette valeur peut être « -1 » et peut généralement être interprétée comme « inconnue ». |
Exemples
Exemple d'événement StreamingQueryListener Kafka-à-Kafka
{
"id" : "3574feba-646d-4735-83c4-66f657e52517",
"runId" : "38a78903-9e55-4440-ad81-50b591e4746c",
"name" : "STREAMING_QUERY_NAME_UNIQUE",
"timestamp" : "2022-10-31T20:09:30.455Z",
"batchId" : 1377,
"numInputRows" : 687,
"inputRowsPerSecond" : 32.13433743393049,
"processedRowsPerSecond" : 34.067241892293964,
"durationMs" : {
"addBatch" : 18352,
"getBatch" : 0,
"latestOffset" : 31,
"queryPlanning" : 977,
"triggerExecution" : 20165,
"walCommit" : 342
},
"eventTime" : {
"avg" : "2022-10-31T20:09:18.070Z",
"max" : "2022-10-31T20:09:30.125Z",
"min" : "2022-10-31T20:09:09.793Z",
"watermark" : "2022-10-31T20:08:46.355Z"
},
"stateOperators" : [ {
"operatorName" : "stateStoreSave",
"numRowsTotal" : 208,
"numRowsUpdated" : 73,
"allUpdatesTimeMs" : 434,
"numRowsRemoved" : 76,
"allRemovalsTimeMs" : 515,
"commitTimeMs" : 0,
"memoryUsedBytes" : 167069743,
"numRowsDroppedByWatermark" : 0,
"numShufflePartitions" : 20,
"numStateStoreInstances" : 20,
"customMetrics" : {
"SnapshotLastUploaded.partition_0_default" : 1370,
"SnapshotLastUploaded.partition_1_default" : 1370,
"SnapshotLastUploaded.partition_2_default" : 1362,
"SnapshotLastUploaded.partition_3_default" : 1370,
"SnapshotLastUploaded.partition_4_default" : 1356,
"rocksdbBytesCopied" : 0,
"rocksdbCommitCheckpointLatency" : 0,
"rocksdbCommitCompactLatency" : 0,
"rocksdbCommitFileSyncLatencyMs" : 0,
"rocksdbCommitFlushLatency" : 0,
"rocksdbCommitPauseLatency" : 0,
"rocksdbCommitWriteBatchLatency" : 0,
"rocksdbFilesCopied" : 0,
"rocksdbFilesReused" : 0,
"rocksdbGetCount" : 222,
"rocksdbGetLatency" : 0,
"rocksdbPutCount" : 0,
"rocksdbPutLatency" : 0,
"rocksdbReadBlockCacheHitCount" : 165,
"rocksdbReadBlockCacheMissCount" : 41,
"rocksdbSstFileSize" : 232729,
"rocksdbTotalBytesRead" : 12844,
"rocksdbTotalBytesReadByCompaction" : 0,
"rocksdbTotalBytesReadThroughIterator" : 161238,
"rocksdbTotalBytesWritten" : 0,
"rocksdbTotalBytesWrittenByCompaction" : 0,
"rocksdbTotalCompactionLatencyMs" : 0,
"rocksdbTotalFlushLatencyMs" : 0,
"rocksdbWriterStallLatencyMs" : 0,
"rocksdbZipFileBytesUncompressed" : 0
}
}, {
"operatorName" : "dedupe",
"numRowsTotal" : 2454744,
"numRowsUpdated" : 73,
"allUpdatesTimeMs" : 4155,
"numRowsRemoved" : 0,
"allRemovalsTimeMs" : 0,
"commitTimeMs" : 0,
"memoryUsedBytes" : 137765341,
"numRowsDroppedByWatermark" : 34,
"numShufflePartitions" : 20,
"numStateStoreInstances" : 20,
"customMetrics" : {
"SnapshotLastUploaded.partition_0_default" : 1360,
"SnapshotLastUploaded.partition_1_default" : 1360,
"SnapshotLastUploaded.partition_2_default" : 1352,
"SnapshotLastUploaded.partition_3_default" : 1360,
"SnapshotLastUploaded.partition_4_default" : 1346,
"numDroppedDuplicateRows" : 193,
"rocksdbBytesCopied" : 0,
"rocksdbCommitCheckpointLatency" : 0,
"rocksdbCommitCompactLatency" : 0,
"rocksdbCommitFileSyncLatencyMs" : 0,
"rocksdbCommitFlushLatency" : 0,
"rocksdbCommitPauseLatency" : 0,
"rocksdbCommitWriteBatchLatency" : 0,
"rocksdbFilesCopied" : 0,
"rocksdbFilesReused" : 0,
"rocksdbGetCount" : 146,
"rocksdbGetLatency" : 0,
"rocksdbPutCount" : 0,
"rocksdbPutLatency" : 0,
"rocksdbReadBlockCacheHitCount" : 3,
"rocksdbReadBlockCacheMissCount" : 3,
"rocksdbSstFileSize" : 78959140,
"rocksdbTotalBytesRead" : 0,
"rocksdbTotalBytesReadByCompaction" : 0,
"rocksdbTotalBytesReadThroughIterator" : 0,
"rocksdbTotalBytesWritten" : 0,
"rocksdbTotalBytesWrittenByCompaction" : 0,
"rocksdbTotalCompactionLatencyMs" : 0,
"rocksdbTotalFlushLatencyMs" : 0,
"rocksdbWriterStallLatencyMs" : 0,
"rocksdbZipFileBytesUncompressed" : 0
}
}, {
"operatorName" : "symmetricHashJoin",
"numRowsTotal" : 2583,
"numRowsUpdated" : 682,
"allUpdatesTimeMs" : 9645,
"numRowsRemoved" : 508,
"allRemovalsTimeMs" : 46,
"commitTimeMs" : 21,
"memoryUsedBytes" : 668544484,
"numRowsDroppedByWatermark" : 0,
"numShufflePartitions" : 20,
"numStateStoreInstances" : 80,
"customMetrics" : {
"SnapshotLastUploaded.partition_0_left-keyToNumValues" : 1310,
"SnapshotLastUploaded.partition_1_left-keyWithIndexToValue" : 1318,
"SnapshotLastUploaded.partition_2_left-keyToNumValues" : 1305,
"SnapshotLastUploaded.partition_2_right-keyWithIndexToValue" : 1306,
"SnapshotLastUploaded.partition_4_left-keyWithIndexToValue" : 1310,
"rocksdbBytesCopied" : 0,
"rocksdbCommitCheckpointLatency" : 0,
"rocksdbCommitCompactLatency" : 0,
"rocksdbCommitFileSyncLatencyMs" : 0,
"rocksdbCommitFlushLatency" : 0,
"rocksdbCommitPauseLatency" : 0,
"rocksdbCommitWriteBatchLatency" : 0,
"rocksdbFilesCopied" : 0,
"rocksdbFilesReused" : 0,
"rocksdbGetCount" : 4218,
"rocksdbGetLatency" : 3,
"rocksdbPutCount" : 0,
"rocksdbPutLatency" : 0,
"rocksdbReadBlockCacheHitCount" : 3425,
"rocksdbReadBlockCacheMissCount" : 149,
"rocksdbSstFileSize" : 742827,
"rocksdbTotalBytesRead" : 866864,
"rocksdbTotalBytesReadByCompaction" : 0,
"rocksdbTotalBytesReadThroughIterator" : 0,
"rocksdbTotalBytesWritten" : 0,
"rocksdbTotalBytesWrittenByCompaction" : 0,
"rocksdbTotalCompactionLatencyMs" : 0,
"rocksdbTotalFlushLatencyMs" : 0,
"rocksdbWriterStallLatencyMs" : 0,
"rocksdbZipFileBytesUncompressed" : 0
}
} ],
"sources" : [ {
"description" : "KafkaV2[Subscribe[KAFKA_TOPIC_NAME_INPUT_A]]",
"startOffset" : {
"KAFKA_TOPIC_NAME_INPUT_A" : {
"0" : 349706380
}
},
"endOffset" : {
"KAFKA_TOPIC_NAME_INPUT_A" : {
"0" : 349706672
}
},
"latestOffset" : {
"KAFKA_TOPIC_NAME_INPUT_A" : {
"0" : 349706672
}
},
"numInputRows" : 292,
"inputRowsPerSecond" : 13.65826278123392,
"processedRowsPerSecond" : 14.479817514628582,
"metrics" : {
"avgOffsetsBehindLatest" : "0.0",
"estimatedTotalBytesBehindLatest" : "0.0",
"maxOffsetsBehindLatest" : "0",
"minOffsetsBehindLatest" : "0"
}
}, {
"description" : "KafkaV2[Subscribe[KAFKA_TOPIC_NAME_INPUT_B]]",
"startOffset" : {
KAFKA_TOPIC_NAME_INPUT_B" : {
"2" : 143147812,
"1" : 129288266,
"0" : 138102966
}
},
"endOffset" : {
"KAFKA_TOPIC_NAME_INPUT_B" : {
"2" : 143147812,
"1" : 129288266,
"0" : 138102966
}
},
"latestOffset" : {
"KAFKA_TOPIC_NAME_INPUT_B" : {
"2" : 143147812,
"1" : 129288266,
"0" : 138102966
}
},
"numInputRows" : 0,
"inputRowsPerSecond" : 0.0,
"processedRowsPerSecond" : 0.0,
"metrics" : {
"avgOffsetsBehindLatest" : "0.0",
"maxOffsetsBehindLatest" : "0",
"minOffsetsBehindLatest" : "0"
}
} ],
"sink" : {
"description" : "org.apache.spark.sql.kafka010.KafkaSourceProvider$KafkaTable@e04b100",
"numOutputRows" : 76
}
}
Exemple d'événement StreamingQueryListener de Delta Lake à Delta Lake
{
"id" : "aeb6bc0f-3f7d-4928-a078-ba2b304e2eaf",
"runId" : "35d751d9-2d7c-4338-b3de-6c6ae9ebcfc2",
"name" : "silverTransformFromBronze",
"timestamp" : "2022-11-01T18:21:29.500Z",
"batchId" : 4,
"numInputRows" : 0,
"inputRowsPerSecond" : 0.0,
"processedRowsPerSecond" : 0.0,
"durationMs" : {
"latestOffset" : 62,
"triggerExecution" : 62
},
"stateOperators" : [ ],
"sources" : [ {
"description" : "DeltaSource[dbfs:/FileStore/<user>/stateful-trade-analysis-demo/table]",
"startOffset" : {
"sourceVersion" : 1,
"reservoirId" : "84590dac-da51-4e0f-8eda-6620198651a9",
"reservoirVersion" : 3216,
"index" : 3214,
"isStartingVersion" : true
},
"endOffset" : {
"sourceVersion" : 1,
"reservoirId" : "84590dac-da51-4e0f-8eda-6620198651a9",
"reservoirVersion" : 3216,
"index" : 3214,
"isStartingVersion" : true
},
"latestOffset" : null,
"numInputRows" : 0,
"inputRowsPerSecond" : 0.0,
"processedRowsPerSecond" : 0.0,
"metrics" : {
"numBytesOutstanding" : "0",
"numFilesOutstanding" : "0"
}
} ],
"sink" : {
"description" : "DeltaSink[dbfs:/user/hive/warehouse/<user>.db/trade_history_silver_delta_demo2]",
"numOutputRows" : -1
}
}
Exemple d'événement StreamingQueryListener de Kinesis à Delta Lake
{
"id" : "3ce9bd93-da16-4cb3-a3b6-e97a592783b5",
"runId" : "fe4a6bda-dda2-4067-805d-51260d93260b",
"name" : null,
"timestamp" : "2024-05-14T02:09:20.846Z",
"batchId" : 0,
"batchDuration" : 59322,
"numInputRows" : 20,
"inputRowsPerSecond" : 0.0,
"processedRowsPerSecond" : 0.33714304979602844,
"durationMs" : {
"addBatch" : 5397,
"commitBatch" : 4429,
"commitOffsets" : 211,
"getBatch" : 5,
"latestOffset" : 21998,
"queryPlanning" : 12128,
"triggerExecution" : 59313,
"walCommit" : 220
},
"stateOperators" : [ ],
"sources" : [ {
"description" : "KinesisV2[KinesisTestUtils-7199466178786508570-at-1715652545256]",
"startOffset" : null,
"endOffset" : [ {
"shard" : {
"stream" : "KinesisTestUtils-7199466178786508570-at-1715652545256",
"shardId" : "shardId-000000000000"
},
"firstSeqNum" : "49652022592149344892294981243280420130985816456924495874",
"lastSeqNum" : "49652022592149344892294981243290091537542733559041622018",
"closed" : false,
"msBehindLatest" : "0",
"lastRecordSeqNum" : "49652022592149344892294981243290091537542733559041622018"
}, {
"shard" : {
"stream" : "KinesisTestUtils-7199466178786508570-at-1715652545256",
"shardId" : "shardId-000000000001"
},
"firstSeqNum" : "49652022592171645637493511866421955849258464818430476306",
"lastSeqNum" : "49652022592171645637493511866434045107454611178897014802",
"closed" : false,
"msBehindLatest" : "0",
"lastRecordSeqNum" : "49652022592171645637493511866434045107454611178897014802"
} ],
"latestOffset" : null,
"numInputRows" : 20,
"inputRowsPerSecond" : 0.0,
"processedRowsPerSecond" : 0.33714304979602844,
"metrics" : {
"avgMsBehindLatest" : "0.0",
"maxMsBehindLatest" : "0",
"minMsBehindLatest" : "0",
"mode" : "efo",
"numClosedShards" : "0",
"numProcessedBytes" : "30",
"numProcessedRecords" : "18",
"numRegisteredConsumers" : "1",
"numStreams" : "1",
"numTotalShards" : "2",
"totalPrefetchedBytes" : "0"
}
} ],
"sink" : {
"description" : "DeltaSink[dbfs:/streaming/test/KinesisToDeltaServerlessLiteSuite/<run-id>/deltaTable]",
"numOutputRows" : -1
}
}
Exemple d'événement StreamingQueryListener Kafka+Delta Lake vers Delta Lake
{
"id" : "210f4746-7caa-4a51-bd08-87cabb45bdbe",
"runId" : "42a2f990-c463-4a9c-9aae-95d6990e63f4",
"name" : null,
"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]]",
"startOffset" : null,
"endOffset" : {
"topic-1" : {
"1" : 5,
"0" : 5
}
},
"latestOffset" : {
"topic-1" : {
"1" : 5,
"0" : 5
}
},
"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]",
"startOffset" : null,
"endOffset" : {
"sourceVersion" : 1,
"reservoirId" : "b207a1cd-0fbe-4652-9c8f-e5cc467ae84f",
"reservoirVersion" : 1,
"index" : -1,
"isStartingVersion" : false
},
"latestOffset" : null,
"numInputRows" : 10,
"inputRowsPerSecond" : 0.0,
"processedRowsPerSecond" : 2.77700638711469,
"metrics" : {
"numBytesOutstanding" : "0",
"numFilesOutstanding" : "0"
}
} ],
"sink" : {
"description" : "DeltaSink[/tmp/spark-d445c92a-4640-4827-a9bd-47246a30bb04]",
"numOutputRows" : -1
}
}
Source de taux d'exemple vers l'événement StreamingQueryListener de Delta Lake
{
"id" : "912ebdc1-edf2-48ec-b9fb-1a9b67dd2d9e",
"runId" : "85de73a5-92cc-4b7f-9350-f8635b0cf66e",
"name" : "dataGen",
"timestamp" : "2022-11-01T18:28:20.332Z",
"batchId" : 279,
"numInputRows" : 300,
"inputRowsPerSecond" : 114.15525114155251,
"processedRowsPerSecond" : 158.9825119236884,
"durationMs" : {
"addBatch" : 1771,
"commitOffsets" : 54,
"getBatch" : 0,
"latestOffset" : 0,
"queryPlanning" : 4,
"triggerExecution" : 1887,
"walCommit" : 58
},
"stateOperators" : [ ],
"sources" : [ {
"description" : "RateStreamV2[rowsPerSecond=100, rampUpTimeSeconds=0, numPartitions=default",
"startOffset" : 560,
"endOffset" : 563,
"latestOffset" : 563,
"numInputRows" : 300,
"inputRowsPerSecond" : 114.15525114155251,
"processedRowsPerSecond" : 158.9825119236884
} ],
"sink" : {
"description" : "DeltaSink[dbfs:/user/hive/warehouse/<user>.db/trade_history_bronze_delta_demo]",
"numOutputRows" : -1
}
}