Aller au contenu principal

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.

important

Les limitations suivantes s'appliquent aux workloads utilisant les modes d'accès compute compatibles Unity Catalog :

  • StreamingQueryListener né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é.
  • StreamingQueryListener né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é).
remarque

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
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 = {}
}

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.

    QueryExecutionListener est appelée lorsque la query est terminée. Accédez aux métriques à l'aide de la carte QueryExecution.observedMetrics.

  • Streaming ou micro-batch : StreamingQueryListenerutilisez.

    StreamingQueryListener est appelée lorsque la query de streaming termine une époque. Accédez aux métriques à l'aide de la carte StreamingQueryProgress.observedMetrics. Databricks ne prend pas en charge le mode continuous Trigger pour le streaming.

Par exemple :

Scala
// 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
}
}
}
})

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 :

SQL
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

id

Un ID de query unique qui persiste après les redémarrages.

runId

Un ID de query unique pour chaque start/redémarrage. Consultez StreamingQuery.runId().

name

Le nom de la query spécifié par l'utilisateur. Le nom est nul si aucun nom n'est spécifié.

timestamp

Le Timestamp de l'exécution du micro-batch.

batchId

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

batchDuration

La durée de traitement d'une opération de batch, en millisecondes.

numInputRows

Nombre total (toutes sources confondues) d'enregistrements traités dans un trigger.

inputRowsPerSecond

Le taux agrégé (toutes sources confondues) des données entrantes.

processedRowsPerSecond

Le taux agrégé (toutes sources confondues) auquel Spark traite les données.

Champs

Description

id

Un ID de query unique qui persiste après les redémarrages.

runId

Un ID de query unique pour chaque start/redémarrage. Consultez StreamingQuery.runId().

name

Le nom de la query spécifié par l'utilisateur. Le nom est nul si aucun nom n'est spécifié.

timestamp

Le Timestamp de l'exécution du micro-batch.

batchId

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

batchDuration

La durée de traitement d'une opération de batch, en millisecondes.

numInputRows

Nombre total (toutes sources confondues) d'enregistrements traités dans un trigger.

inputRowsPerSecond

Le taux agrégé (toutes sources confondues) des données entrantes.

processedRowsPerSecond

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

durationMs

Type : ju.Map[String, JLong]. Voir objet durationMs.

eventTime

Type : ju.Map[String, String]. Voir objet eventTime.

stateOperators

Type : Array[StateOperatorProgress]. Consultez l'objet stateOperators.

sources

Type : Array[SourceProgress]. Consultez objet sources.

sink

Type : SinkProgress. Consultez l’ objet puits.

observedMetrics

Type : ju.Map[String, Row]. Fonctions d'agrégation arbitraires nommées pouvant être définies sur un DataFrame/une query (tel que df.observe).

Champs

Description

durationMs

Type : ju.Map[String, JLong]. Voir objet durationMs.

eventTime

Type : ju.Map[String, String]. Voir objet eventTime.

stateOperators

Type : Array[StateOperatorProgress]. Consultez l'objet stateOperators.

sources

Type : Array[SourceProgress]. Consultez objet sources.

sink

Type : SinkProgress. Consultez l’ objet puits.

observedMetrics

Type : ju.Map[String, Row]. Fonctions d'agrégation arbitraires nommées pouvant être définies sur un DataFrame/une query (tel que df.observe).

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

durationMs.addBatch

La durée d'exécution du micro-batch. Cela exclut le temps que Spark prend pour planifier le micro-batch.

durationMs.getBatch

Le temps nécessaire pour récupérer les métadonnées concernant les offsets de la source.

durationMs.latestOffset

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.

durationMs.queryPlanning

Le temps nécessaire pour générer le plan d'exécution.

durationMs.triggerExecution

Le temps nécessaire pour planifier et exécuter le micro-lot.

durationMs.walCommit

Le temps nécessaire pour commit les nouveaux décalages disponibles.

durationMs.commitBatch

Le temps nécessaire pour commit les données écrites dans le sink pendant addBatch. Uniquement présent pour les destinations qui prennent en charge le commit.

durationMs.commitOffsets

Le temps nécessaire pour commit le batch dans le log de commit.

Champs

Description

durationMs.addBatch

La durée d'exécution du micro-batch. Cela exclut le temps que Spark prend pour planifier le micro-batch.

durationMs.getBatch

Le temps nécessaire pour récupérer les métadonnées concernant les offsets de la source.

durationMs.latestOffset

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.

durationMs.queryPlanning

Le temps nécessaire pour générer le plan d'exécution.

durationMs.triggerExecution

Le temps nécessaire pour planifier et exécuter le micro-lot.

durationMs.walCommit

Le temps nécessaire pour commit les nouveaux décalages disponibles.

durationMs.commitBatch

Le temps nécessaire pour commit les données écrites dans le sink pendant addBatch. Uniquement présent pour les destinations qui prennent en charge le commit.

durationMs.commitOffsets

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

eventTime.avg

Le temps d'événement moyen observé dans ce Trigger.

eventTime.max

Le temps d'événement maximal observé dans ce Trigger.

eventTime.min

Le temps minimal d'événement observé dans ce Trigger.

eventTime.watermark

La valeur du filigrane utilisé dans ce trigger.

Champs

Description

eventTime.avg

Le temps d'événement moyen observé dans ce Trigger.

eventTime.max

Le temps d'événement maximal observé dans ce Trigger.

eventTime.min

Le temps minimal d'événement observé dans ce Trigger.

eventTime.watermark

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

stateOperators.operatorName

Le nom de l'opérateur avec état auquel les métriques sont liées, tels que symmetricHashJoin, dedupe ou stateStoreSave.

stateOperators.numRowsTotal

Le nombre total de lignes en état à la suite d'un opérateur avec état ou d'une agrégation.

stateOperators.numRowsUpdated

Le nombre total de lignes mises à jour dans l'état à la suite d'un opérateur avec état ou d'une agrégation.

stateOperators.allUpdatesTimeMs

Cette métrique n'est actuellement pas mesurable par Spark et il est prévu de la supprimer lors de futures mises à jour.

stateOperators.numRowsRemoved

Le nombre total de lignes supprimées de l'état suite à un opérateur avec état ou à une agrégation.

stateOperators.allRemovalsTimeMs

Cette métrique n'est actuellement pas mesurable par Spark et il est prévu de la supprimer lors de futures mises à jour.

stateOperators.commitTimeMs

Le temps nécessaire pour commit toutes les mises à jour (ajouts et suppressions) et renvoyer une nouvelle version.

stateOperators.memoryUsedBytes

Mémoire utilisée par le magasin d'état.

stateOperators.numRowsDroppedByWatermark

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.

stateOperators.numShufflePartitions

Le nombre de partitions de brassage pour cet opérateur avec état.

stateOperators.numStateStoreInstances

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.

stateOperators.customMetrics

Consultez stateOperators.customMetrics dans cette rubrique pour en savoir plus.

Champs

Description

stateOperators.operatorName

Le nom de l'opérateur avec état auquel les métriques sont liées, tels que symmetricHashJoin, dedupe ou stateStoreSave.

stateOperators.numRowsTotal

Le nombre total de lignes en état à la suite d'un opérateur avec état ou d'une agrégation.

stateOperators.numRowsUpdated

Le nombre total de lignes mises à jour dans l'état à la suite d'un opérateur avec état ou d'une agrégation.

stateOperators.allUpdatesTimeMs

Cette métrique n'est actuellement pas mesurable par Spark et il est prévu de la supprimer lors de futures mises à jour.

stateOperators.numRowsRemoved

Le nombre total de lignes supprimées de l'état suite à un opérateur avec état ou à une agrégation.

stateOperators.allRemovalsTimeMs

Cette métrique n'est actuellement pas mesurable par Spark et il est prévu de la supprimer lors de futures mises à jour.

stateOperators.commitTimeMs

Le temps nécessaire pour commit toutes les mises à jour (ajouts et suppressions) et renvoyer une nouvelle version.

stateOperators.memoryUsedBytes

Mémoire utilisée par le magasin d'état.

stateOperators.numRowsDroppedByWatermark

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.

stateOperators.numShufflePartitions

Le nombre de partitions de brassage pour cet opérateur avec état.

stateOperators.numStateStoreInstances

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.

stateOperators.customMetrics

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.

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

customMetrics.rocksdbBytesCopied

Le nombre d'octets copiés tel que suivi par le gestionnaire de fichiers RocksDB.

customMetrics.rocksdbCommitCheckpointLatency

Le temps nécessaire, en millisecondes, pour prendre un instantané de RocksDB natif et l’écrire dans un répertoire local.

customMetrics.rocksdbCompactLatency

La durée en millisecondes de la compression (facultatif) pendant le commit du point de contrôle.

customMetrics.rocksdbCommitCompactLatency

Le temps de compactage pendant le commit, en millisecondes.

customMetrics.rocksdbCommitFileSyncLatencyMs

Temps, en millisecondes, nécessaire à la synchronisation de l’instantané natif RocksDB vers le stockage externe (l’emplacement de point de contrôle).

customMetrics.rocksdbCommitFlushLatency

Le temps en millisecondes pour vider les modifications en mémoire de RocksDB sur le disque local.

customMetrics.rocksdbCommitPauseLatency

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.

customMetrics.rocksdbCommitWriteBatchLatency

Le temps en millisecondes d'application des écritures intermédiaires dans une structure en mémoire (WriteBatch) à RocksDB natif.

customMetrics.rocksdbFilesCopied

Le nombre de fichiers copiés, tels que suivis par le gestionnaire de fichiers RocksDB.

customMetrics.rocksdbFilesReused

Le nombre de fichiers réutilisés et suivis par le gestionnaire de fichiers RocksDB.

customMetrics.rocksdbGetCount

Le nombre d'appels get (n'inclut pas gets de WriteBatch – batch en mémoire utilisé pour les écritures intermédiaires).

customMetrics.rocksdbGetLatency

Le temps moyen en nanosecondes pour l'appel natif RocksDB::Get sous-jacent.

customMetrics.rocksdbReadBlockCacheHitCount

Le nombre d'accès au cache à partir du cache de blocs dans RocksDB.

customMetrics.rocksdbReadBlockCacheMissCount

Le nombre d'échecs du cache de blocs dans RocksDB.

customMetrics.rocksdbSstFileSize

La taille de tous les fichiers Static Sorted Table (SST) dans l'instance RocksDB.

customMetrics.rocksdbTotalBytesRead

Le nombre d'octets décompressés lus par les opérations get.

customMetrics.rocksdbTotalBytesWritten

Le nombre total d’octets non compressés écrits par put Opérations.

customMetrics.rocksdbTotalBytesReadThroughIterator

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 FlatMapGroupsWithState et le filigranage) nécessitent la lecture de données dans Databricks via un itérateur.

customMetrics.rocksdbTotalBytesReadByCompaction

Le nombre d'octets que le processus de compactage lit sur le disque.

customMetrics.rocksdbTotalBytesWrittenByCompaction

Le nombre total d'octets que le processus de compaction écrit sur le disque.

customMetrics.rocksdbTotalCompactionLatencyMs

Le temps en millisecondes pour les compactages RocksDB, y compris les compactages en arrière-plan et le compactage facultatif lancé pendant le commit.

customMetrics.rocksdbTotalFlushLatencyMs

Le temps de vidage total, y compris le vidage en arrière-plan. Les opérations de vidage sont des processus par lesquels le MemTable est vidé vers le stockage une fois qu'il est plein. MemTables sont le premier niveau où les données sont stockées dans RocksDB.

customMetrics.rocksdbZipFileBytesUncompressed

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.

customMetrics.SnapshotLastUploaded.partition_<partition-id>_<state-store-name>

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.

customMetrics.rocksdbPutLatency

La latence totale des options d'achat et de vente.

customMetrics.rocksdbPutCount

Le nombre d'appels de mise.

customMetrics.rocksdbWriterStallLatencyMs

Le temps d'attente du processus d'écriture pour que la compaction ou le vidage se termine.

customMetrics.rocksdbTotalBytesWrittenByFlush

Le nombre total d'octets écrits par vidage.

customMetrics.rocksdbPinnedBlocksMemoryUsage

L'utilisation de la mémoire pour les blocs épinglés

customMetrics.rocksdbNumInternalColFamiliesKeys

Le nombre de clés internes pour les familles de colonnes internes

customMetrics.rocksdbNumExternalColumnFamilies

Le nombre de familles de colonnes externes

customMetrics.rocksdbNumInternalColumnFamilies

Le nombre de familles de colonnes internes

Champs

Description

customMetrics.rocksdbBytesCopied

Le nombre d'octets copiés tel que suivi par le gestionnaire de fichiers RocksDB.

customMetrics.rocksdbCommitCheckpointLatency

Le temps nécessaire, en millisecondes, pour prendre un instantané de RocksDB natif et l’écrire dans un répertoire local.

customMetrics.rocksdbCompactLatency

La durée en millisecondes de la compression (facultatif) pendant le commit du point de contrôle.

customMetrics.rocksdbCommitCompactLatency

Le temps de compactage pendant le commit, en millisecondes.

customMetrics.rocksdbCommitFileSyncLatencyMs

Temps, en millisecondes, nécessaire à la synchronisation de l’instantané natif RocksDB vers le stockage externe (l’emplacement de point de contrôle).

customMetrics.rocksdbCommitFlushLatency

Le temps en millisecondes pour vider les modifications en mémoire de RocksDB sur le disque local.

customMetrics.rocksdbCommitPauseLatency

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.

customMetrics.rocksdbCommitWriteBatchLatency

Le temps en millisecondes d'application des écritures intermédiaires dans une structure en mémoire (WriteBatch) à RocksDB natif.

customMetrics.rocksdbFilesCopied

Le nombre de fichiers copiés, tels que suivis par le gestionnaire de fichiers RocksDB.

customMetrics.rocksdbFilesReused

Le nombre de fichiers réutilisés et suivis par le gestionnaire de fichiers RocksDB.

customMetrics.rocksdbGetCount

Le nombre d'appels get (n'inclut pas gets de WriteBatch – batch en mémoire utilisé pour les écritures intermédiaires).

customMetrics.rocksdbGetLatency

Le temps moyen en nanosecondes pour l'appel natif RocksDB::Get sous-jacent.

customMetrics.rocksdbReadBlockCacheHitCount

Le nombre d'accès au cache à partir du cache de blocs dans RocksDB.

customMetrics.rocksdbReadBlockCacheMissCount

Le nombre d'échecs du cache de blocs dans RocksDB.

customMetrics.rocksdbSstFileSize

La taille de tous les fichiers Static Sorted Table (SST) dans l'instance RocksDB.

customMetrics.rocksdbTotalBytesRead

Le nombre d'octets décompressés lus par les opérations get.

customMetrics.rocksdbTotalBytesWritten

Le nombre total d’octets non compressés écrits par put Opérations.

customMetrics.rocksdbTotalBytesReadThroughIterator

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 FlatMapGroupsWithState et le filigranage) nécessitent la lecture de données dans Databricks via un itérateur.

customMetrics.rocksdbTotalBytesReadByCompaction

Le nombre d'octets que le processus de compactage lit sur le disque.

customMetrics.rocksdbTotalBytesWrittenByCompaction

Le nombre total d'octets que le processus de compaction écrit sur le disque.

customMetrics.rocksdbTotalCompactionLatencyMs

Le temps en millisecondes pour les compactages RocksDB, y compris les compactages en arrière-plan et le compactage facultatif lancé pendant le commit.

customMetrics.rocksdbTotalFlushLatencyMs

Le temps de vidage total, y compris le vidage en arrière-plan. Les opérations de vidage sont des processus par lesquels le MemTable est vidé vers le stockage une fois qu'il est plein. MemTables sont le premier niveau où les données sont stockées dans RocksDB.

customMetrics.rocksdbZipFileBytesUncompressed

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.

customMetrics.SnapshotLastUploaded.partition_<partition-id>_<state-store-name>

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.

customMetrics.rocksdbPutLatency

La latence totale des options d'achat et de vente.

customMetrics.rocksdbPutCount

Le nombre d'appels de mise.

customMetrics.rocksdbWriterStallLatencyMs

Le temps d'attente du processus d'écriture pour que la compaction ou le vidage se termine.

customMetrics.rocksdbTotalBytesWrittenByFlush

Le nombre total d'octets écrits par vidage.

customMetrics.rocksdbPinnedBlocksMemoryUsage

L'utilisation de la mémoire pour les blocs épinglés

customMetrics.rocksdbNumInternalColFamiliesKeys

Le nombre de clés internes pour les familles de colonnes internes

customMetrics.rocksdbNumExternalColumnFamilies

Le nombre de familles de colonnes externes

customMetrics.rocksdbNumInternalColumnFamilies

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

customMetrics.stateOnCurrentVersionSizeBytes

La taille estimée de l'état uniquement sur la version actuelle.

customMetrics.loadedMapCacheHitCount

Le nombre d'accès réussis au cache sur les états mis en cache dans le fournisseur.

customMetrics.loadedMapCacheMissCount

Le nombre de manques de cache sur les états mis en cache chez le fournisseur.

customMetrics.SnapshotLastUploaded.partition_<partition-id>_<state-store-name>

La dernière version upload de l'instantané pour une instance de magasin d'état spécifique.

Champs

Description

customMetrics.stateOnCurrentVersionSizeBytes

La taille estimée de l'état uniquement sur la version actuelle.

customMetrics.loadedMapCacheHitCount

Le nombre d'accès réussis au cache sur les états mis en cache dans le fournisseur.

customMetrics.loadedMapCacheMissCount

Le nombre de manques de cache sur les états mis en cache chez le fournisseur.

customMetrics.SnapshotLastUploaded.partition_<partition-id>_<state-store-name>

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

customMetrics.numDroppedDuplicateRows

Le nombre de lignes en double supprimées.

customMetrics.numRowsReadDuringEviction

Le nombre de lignes d'état lues pendant l'éviction d'état.

Champs

Description

customMetrics.numDroppedDuplicateRows

Le nombre de lignes en double supprimées.

customMetrics.numRowsReadDuringEviction

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

customMetrics.numRowsReadDuringEviction

Le nombre de lignes d'état lues pendant l'éviction d'état.

Champs

Description

customMetrics.numRowsReadDuringEviction

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

customMetrics.skippedNullValueCount

Le nombre de valeurs null ignorées, lorsque spark.sql.streaming.stateStore.skipNullsForStreamStreamJoins.enabled est défini sur true.

Champs

Description

customMetrics.skippedNullValueCount

Le nombre de valeurs null ignorées, lorsque spark.sql.streaming.stateStore.skipNullsForStreamStreamJoins.enabled est défini sur true.

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

customMetrics.initialStateProcessingTimeMs

Nombre de millisecondes nécessaires pour traiter tout l'état initial.

customMetrics.numValueStateVars

Nombre de variables d'état de valeur. Également présent pour transformWithStateInPandas.

customMetrics.numListStateVars

Nombre de variables d'état de liste. Également présent pour transformWithStateInPandas.

customMetrics.numMapStateVars

Nombre de variables d'état de carte. Également présent pour transformWithStateInPandas.

customMetrics.numDeletedStateVars

Nombre de variables d'état supprimées. Également présent pour transformWithStateInPandas.

customMetrics.timerProcessingTimeMs

Nombre de millisecondes nécessaires au traitement de tous les compteurs.

customMetrics.numRegisteredTimers

Nombre de minuteurs enregistrés. Également présent pour transformWithStateInPandas.

customMetrics.numDeletedTimers

Nombre de minuteurs supprimés. Également présent pour transformWithStateInPandas.

customMetrics.numExpiredTimers

Nombre de minuteurs expirés. Également présent pour transformWithStateInPandas.

customMetrics.numValueStateWithTTLVars

Nombre de variables d'état de la valeur dotées d'un TTL. Également présent pour transformWithStateInPandas.

customMetrics.numListStateWithTTLVars

Nombre de variables d'état de liste avec TTL. Également présent pour transformWithStateInPandas.

customMetrics.numMapStateWithTTLVars

Nombre de variables d'état de carte avec TTL. Également présent pour transformWithStateInPandas.

customMetrics.numValuesRemovedDueToTTLExpiry

Nombre de valeurs supprimées en raison de l'expiration du TTL. Également présent pour transformWithStateInPandas.

customMetrics.numValuesIncrementallyRemovedDueToTTLExpiry

Nombre de valeurs supprimées de manière incrémentielle en raison de l'expiration de la TTL.

Champs

Description

customMetrics.initialStateProcessingTimeMs

Nombre de millisecondes nécessaires pour traiter tout l'état initial.

customMetrics.numValueStateVars

Nombre de variables d'état de valeur. Également présent pour transformWithStateInPandas.

customMetrics.numListStateVars

Nombre de variables d'état de liste. Également présent pour transformWithStateInPandas.

customMetrics.numMapStateVars

Nombre de variables d'état de carte. Également présent pour transformWithStateInPandas.

customMetrics.numDeletedStateVars

Nombre de variables d'état supprimées. Également présent pour transformWithStateInPandas.

customMetrics.timerProcessingTimeMs

Nombre de millisecondes nécessaires au traitement de tous les compteurs.

customMetrics.numRegisteredTimers

Nombre de minuteurs enregistrés. Également présent pour transformWithStateInPandas.

customMetrics.numDeletedTimers

Nombre de minuteurs supprimés. Également présent pour transformWithStateInPandas.

customMetrics.numExpiredTimers

Nombre de minuteurs expirés. Également présent pour transformWithStateInPandas.

customMetrics.numValueStateWithTTLVars

Nombre de variables d'état de la valeur dotées d'un TTL. Également présent pour transformWithStateInPandas.

customMetrics.numListStateWithTTLVars

Nombre de variables d'état de liste avec TTL. Également présent pour transformWithStateInPandas.

customMetrics.numMapStateWithTTLVars

Nombre de variables d'état de carte avec TTL. Également présent pour transformWithStateInPandas.

customMetrics.numValuesRemovedDueToTTLExpiry

Nombre de valeurs supprimées en raison de l'expiration du TTL. Également présent pour transformWithStateInPandas.

customMetrics.numValuesIncrementallyRemovedDueToTTLExpiry

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

Description détaillée de la table de source de données de streaming.

startOffset

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

endOffset

Le dernier décalage traité par le micro-batch.

latestOffset

Le décalage le plus récent traité par le micro-batch.

numInputRows

Le nombre de lignes d'entrée traitées à partir de cette source.

inputRowsPerSecond

Le taux, en secondes, auquel les données arrivent pour traitement depuis cette source.

processedRowsPerSecond

Le débit auquel Spark traite les données de cette source.

metrics

Type : ju.Map[String, String]. Contient des métriques personnalisées pour une source de données spécifique.

Champs

Description

description

Description détaillée de la table de source de données de streaming.

startOffset

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

endOffset

Le dernier décalage traité par le micro-batch.

latestOffset

Le décalage le plus récent traité par le micro-batch.

numInputRows

Le nombre de lignes d'entrée traitées à partir de cette source.

inputRowsPerSecond

Le taux, en secondes, auquel les données arrivent pour traitement depuis cette source.

processedRowsPerSecond

Le débit auquel Spark traite les données de cette source.

metrics

Type : ju.Map[String, String]. Contient des métriques personnalisées pour une source de données spécifique.

Databricks fournit l'implémentation d'objets sources suivante :

remarque

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

sources.description

La description de la source à partir de laquelle la query de streaming lit. Par exemple : ”DeltaSource[table]”.

sources.<startOffset / endOffset>.sourceVersion

La version de sérialisation avec laquelle ce décalage est encodé.

sources.<startOffset / endOffset>.reservoirId

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.

sources.<startOffset / endOffset>.reservoirVersion

La version de la table qui est en cours de traitement.

sources.<startOffset / endOffset>.index

L'index dans la séquence de AddFiles dans cette version. Ceci est utilisé pour diviser les commits volumineux en plusieurs batches. Cet index est créé en triant sur modificationTimestamp et path.

sources.<startOffset / endOffset>.isStartingVersion

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.

sources.<startOffset / endOffset / latestOffset>.eventTimeMillis

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.

sources.latestOffset

L'offset le plus récent traité par la query micro-batch.

sources.numInputRows

Le nombre de lignes d'entrée traitées à partir de cette source.

sources.inputRowsPerSecond

Le rythme auquel les données arrivent pour traitement depuis cette source.

sources.processedRowsPerSecond

Le débit auquel Spark traite les données de cette source.

sources.metrics.numBytesOutstanding

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.

sources.metrics.numFilesOutstanding

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.

Champs

Description

sources.description

La description de la source à partir de laquelle la query de streaming lit. Par exemple : ”DeltaSource[table]”.

sources.<startOffset / endOffset>.sourceVersion

La version de sérialisation avec laquelle ce décalage est encodé.

sources.<startOffset / endOffset>.reservoirId

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.

sources.<startOffset / endOffset>.reservoirVersion

La version de la table qui est en cours de traitement.

sources.<startOffset / endOffset>.index

L'index dans la séquence de AddFiles dans cette version. Ceci est utilisé pour diviser les commits volumineux en plusieurs batches. Cet index est créé en triant sur modificationTimestamp et path.

sources.<startOffset / endOffset>.isStartingVersion

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.

sources.<startOffset / endOffset / latestOffset>.eventTimeMillis

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.

sources.latestOffset

L'offset le plus récent traité par la query micro-batch.

sources.numInputRows

Le nombre de lignes d'entrée traitées à partir de cette source.

sources.inputRowsPerSecond

Le rythme auquel les données arrivent pour traitement depuis cette source.

sources.processedRowsPerSecond

Le débit auquel Spark traite les données de cette source.

sources.metrics.numBytesOutstanding

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.

sources.metrics.numFilesOutstanding

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

sources.description

Description détaillée de la source Kafka, spécifiant le sujet Kafka exact à partir duquel la lecture est effectuée. Par exemple : ”KafkaV2[Subscribe[KAFKA_TOPIC_NAME_INPUT_A]]”.

sources.startOffset

Le numéro de décalage de début dans le sujet Kafka auquel le Job de streaming a démarré.

sources.endOffset

Le dernier offset traité par le micro-batch. Cela pourrait être égal à latestOffset pour une exécution micro-batch en cours.

sources.latestOffset

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 endOffset et latestOffset.

sources.numInputRows

Le nombre de lignes d'entrée traitées à partir de cette source.

sources.inputRowsPerSecond

Le rythme auquel les données arrivent pour traitement depuis cette source.

sources.processedRowsPerSecond

Le débit auquel Spark traite les données de cette source.

sources.metrics.avgOffsetsBehindLatest

Le nombre moyen de décalages de la query en streaming par rapport au dernier décalage disponible parmi tous les sujets abonnés.

sources.metrics.estimatedTotalBytesBehindLatest

Le nombre estimé d'octets que le processus de query n'a pas consommés à partir des sujets abonnés.

sources.metrics.maxOffsetsBehindLatest

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.

sources.metrics.minOffsetsBehindLatest

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.

Champs

Description

sources.description

Description détaillée de la source Kafka, spécifiant le sujet Kafka exact à partir duquel la lecture est effectuée. Par exemple : ”KafkaV2[Subscribe[KAFKA_TOPIC_NAME_INPUT_A]]”.

sources.startOffset

Le numéro de décalage de début dans le sujet Kafka auquel le Job de streaming a démarré.

sources.endOffset

Le dernier offset traité par le micro-batch. Cela pourrait être égal à latestOffset pour une exécution micro-batch en cours.

sources.latestOffset

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 endOffset et latestOffset.

sources.numInputRows

Le nombre de lignes d'entrée traitées à partir de cette source.

sources.inputRowsPerSecond

Le rythme auquel les données arrivent pour traitement depuis cette source.

sources.processedRowsPerSecond

Le débit auquel Spark traite les données de cette source.

sources.metrics.avgOffsetsBehindLatest

Le nombre moyen de décalages de la query en streaming par rapport au dernier décalage disponible parmi tous les sujets abonnés.

sources.metrics.estimatedTotalBytesBehindLatest

Le nombre estimé d'octets que le processus de query n'a pas consommés à partir des sujets abonnés.

sources.metrics.maxOffsetsBehindLatest

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.

sources.metrics.minOffsetsBehindLatest

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

sources.description

La description de la source Kinesis, spécifiant le Kinesis Stream exact à partir duquel la query de streaming est lue. Par exemple : ”KinesisV2[stream]”.

sources.metrics.avgMsBehindLatest

Le nombre moyen de millisecondes de retard d'un consommateur par rapport au début d'un stream.

sources.metrics.maxMsBehindLatest

Le nombre maximal de millisecondes pendant lesquelles un consommateur a pris du retard par rapport au début d'un stream.

sources.metrics.minMsBehindLatest

Le nombre minimal de millisecondes qu'un consommateur a pris de retard par rapport au début d'un stream.

sources.metrics.totalPrefetchedBytes

Le nombre d’octets restants à traiter. Ceci est la métrique de backlog pour Kinesis en tant que source.

sources.<startOffset / endOffset / latestOffset>(index).shard.stream

Nom du Stream Kinesis.

sources.<startOffset / endOffset / latestOffset>(index).shard.shardId

ID du shard Kinesis Stream.

sources.<startOffset / endOffset / latestOffset>(index).firstSeqNum

Le premier numéro de séquence des enregistrements dans une partition Kinesis qui ont été consommés dans un batch donné.

sources.<startOffset / endOffset / latestOffset>(index).lastSeqNum

Le dernier numéro de séquence des enregistrements consommés à partir d'un shard Kinesis dans un batch donné.

sources.<startOffset / endOffset / latestOffset>(index).closed

Si le fragment Kinesis a été fermé par le Stream Kinesis.

sources.<startOffset / endOffset / latestOffset>(index).msBehindLatest

Le temps approximatif du retard de la query de streaming par rapport aux dernières données dans le Kinesis stream.

sources.<startOffset / endOffset / latestOffset>(index).lastRecordSeqNum

Le numéro de séquence du dernier enregistrement consommé, et est utilisé pour la vérification des pertes de données. Notez que lastRecordSeqNum peut être différent d'endSeqNum pour les lectures EFO.

sources.metrics.mode

Le mode consommateur utilisé pour l'exécution de la query de streaming. Peut être Polling ou EFO. Mode.

sources.metrics.numStreams

Le nombre de Kinesis Stream traités dans ce micro-batch.

sources.metrics.numTotalShards

Le nombre total de fragments traités dans ce micro-batch.

sources.metrics.numClosedShards

Le nombre de shards fermés traités dans ce micro-batch.

sources.metrics.numProcessedBytes

Le nombre d'octets traités dans ce micro-batch.

sources.metrics.numProcessedRecords

Le nombre d'enregistrements traités dans ce micro-batch.

sources.metrics.numRegisteredConsumers

Le nombre de consommateurs enregistrés utilisés en mode EFO.

Champs

Description

sources.description

La description de la source Kinesis, spécifiant le Kinesis Stream exact à partir duquel la query de streaming est lue. Par exemple : ”KinesisV2[stream]”.

sources.metrics.avgMsBehindLatest

Le nombre moyen de millisecondes de retard d'un consommateur par rapport au début d'un stream.

sources.metrics.maxMsBehindLatest

Le nombre maximal de millisecondes pendant lesquelles un consommateur a pris du retard par rapport au début d'un stream.

sources.metrics.minMsBehindLatest

Le nombre minimal de millisecondes qu'un consommateur a pris de retard par rapport au début d'un stream.

sources.metrics.totalPrefetchedBytes

Le nombre d’octets restants à traiter. Ceci est la métrique de backlog pour Kinesis en tant que source.

sources.<startOffset / endOffset / latestOffset>(index).shard.stream

Nom du Stream Kinesis.

sources.<startOffset / endOffset / latestOffset>(index).shard.shardId

ID du shard Kinesis Stream.

sources.<startOffset / endOffset / latestOffset>(index).firstSeqNum

Le premier numéro de séquence des enregistrements dans une partition Kinesis qui ont été consommés dans un batch donné.

sources.<startOffset / endOffset / latestOffset>(index).lastSeqNum

Le dernier numéro de séquence des enregistrements consommés à partir d'un shard Kinesis dans un batch donné.

sources.<startOffset / endOffset / latestOffset>(index).closed

Si le fragment Kinesis a été fermé par le Stream Kinesis.

sources.<startOffset / endOffset / latestOffset>(index).msBehindLatest

Le temps approximatif du retard de la query de streaming par rapport aux dernières données dans le Kinesis stream.

sources.<startOffset / endOffset / latestOffset>(index).lastRecordSeqNum

Le numéro de séquence du dernier enregistrement consommé, et est utilisé pour la vérification des pertes de données. Notez que lastRecordSeqNum peut être différent d'endSeqNum pour les lectures EFO.

sources.metrics.mode

Le mode consommateur utilisé pour l'exécution de la query de streaming. Peut être Polling ou EFO. Mode.

sources.metrics.numStreams

Le nombre de Kinesis Stream traités dans ce micro-batch.

sources.metrics.numTotalShards

Le nombre total de fragments traités dans ce micro-batch.

sources.metrics.numClosedShards

Le nombre de shards fermés traités dans ce micro-batch.

sources.metrics.numProcessedBytes

Le nombre d'octets traités dans ce micro-batch.

sources.metrics.numProcessedRecords

Le nombre d'enregistrements traités dans ce micro-batch.

sources.metrics.numRegisteredConsumers

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

sources.<startOffset / endOffset / latestOffset>.seqNum

La position actuelle dans la séquence des fichiers traités, dans l'ordre de leur découverte.

sources.<startOffset / endOffset / latestOffset>.sourceVersion

La version d'implémentation de la source cloudFiles.

sources.<startOffset / endOffset / latestOffset>.lastBackfillStartTimeMs

L'heure de start de l'opération de remplissage la plus récente.

sources.<startOffset / endOffset / latestOffset>.lastBackfillFinishTimeMs

L'heure de fin de la plus récente opération de backfill.

sources.<startOffset / endOffset / latestOffset>.lastInputPath

Le dernier chemin d'entrée fourni par l'utilisateur du Stream avant que le Stream ne soit redémarré.

sources.metrics.numFilesOutstanding

Le nombre de fichiers en attente

sources.metrics.numBytesOutstanding

La taille (en octets) des fichiers en attente

sources.metrics.approximateQueueSize

La taille approximative de la file d'attente de messages. Uniquement lorsque l'option cloudFiles.useNotifications est activée.

sources.numInputRows

Le nombre de lignes d'entrée traitées à partir de cette source. Pour le format source binaryFile, numInputRows est égal au nombre de fichiers.

Champs

Description

sources.<startOffset / endOffset / latestOffset>.seqNum

La position actuelle dans la séquence des fichiers traités, dans l'ordre de leur découverte.

sources.<startOffset / endOffset / latestOffset>.sourceVersion

La version d'implémentation de la source cloudFiles.

sources.<startOffset / endOffset / latestOffset>.lastBackfillStartTimeMs

L'heure de start de l'opération de remplissage la plus récente.

sources.<startOffset / endOffset / latestOffset>.lastBackfillFinishTimeMs

L'heure de fin de la plus récente opération de backfill.

sources.<startOffset / endOffset / latestOffset>.lastInputPath

Le dernier chemin d'entrée fourni par l'utilisateur du Stream avant que le Stream ne soit redémarré.

sources.metrics.numFilesOutstanding

Le nombre de fichiers en attente

sources.metrics.numBytesOutstanding

La taille (en octets) des fichiers en attente

sources.metrics.approximateQueueSize

La taille approximative de la file d'attente de messages. Uniquement lorsque l'option cloudFiles.useNotifications est activée.

sources.numInputRows

Le nombre de lignes d'entrée traitées à partir de cette source. Pour le format source binaryFile, numInputRows est égal au nombre de fichiers.

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

sources.<startOffset / endOffset / latestOffset>.sourceVersion

La version d'implémentation avec laquelle ce décalage est encodé.

sources.<startOffset / endOffset / latestOffset>.seqNum

Le numéro de séquence persistant qui est en cours de traitement.

sources.<startOffset / endOffset / latestOffset>.fetchEpoch

La plus grande époque de récupération en cours de traitement.

sources.metrics.numRecordsReadyToProcess

Le nombre d’enregistrements disponibles pour le traitement dans le backlog actuel.

sources.metrics.sizeOfRecordsReadyToProcess

La taille totale en octets des données non traitées dans l'arriéré actuel.

sources.metrics.numDuplicatesSinceStreamStart

Le nombre total d'enregistrements dupliqués traités par le Stream depuis son start.

Champs

Description

sources.<startOffset / endOffset / latestOffset>.sourceVersion

La version d'implémentation avec laquelle ce décalage est encodé.

sources.<startOffset / endOffset / latestOffset>.seqNum

Le numéro de séquence persistant qui est en cours de traitement.

sources.<startOffset / endOffset / latestOffset>.fetchEpoch

La plus grande époque de récupération en cours de traitement.

sources.metrics.numRecordsReadyToProcess

Le nombre d’enregistrements disponibles pour le traitement dans le backlog actuel.

sources.metrics.sizeOfRecordsReadyToProcess

La taille totale en octets des données non traitées dans l'arriéré actuel.

sources.metrics.numDuplicatesSinceStreamStart

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

sources.metrics.numInputRows

Le nombre de lignes traitées dans le micro-batch actuel.

sources.metrics.numInputBytes

Le nombre total d'octets traités dans le micro-batch actuel.

Champs

Description

sources.metrics.numInputRows

Le nombre de lignes traitées dans le micro-batch actuel.

sources.metrics.numInputBytes

Le nombre total d'octets traités dans le micro-batch actuel.

objet sink

Type d'objet : SinkProgress

Champs

Description

sink.description

La description du puits, détaillant l'implémentation spécifique du puits utilisée.

sink.numOutputRows

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.

sink.metrics

ju.Map[String, String] des métriques du récepteur.

Champs

Description

sink.description

La description du puits, détaillant l'implémentation spécifique du puits utilisée.

sink.numOutputRows

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.

sink.metrics

ju.Map[String, String] des métriques du récepteur.

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.

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

sink.description

La description du récepteur Delta, détaillant l'implémentation spécifique du récepteur Delta utilisée. Par exemple : ”DeltaSink[table]”.

sink.numOutputRows

Le nombre de lignes est toujours -1 car Spark ne peut pas inférer les lignes de sortie pour les sinks DSv1, qui est la classification pour le sink Delta Lake.

Champs

Description

sink.description

La description du récepteur Delta, détaillant l'implémentation spécifique du récepteur Delta utilisée. Par exemple : ”DeltaSink[table]”.

sink.numOutputRows

Le nombre de lignes est toujours -1 car Spark ne peut pas inférer les lignes de sortie pour les sinks DSv1, qui est la classification pour le sink Delta Lake.

Objet récepteur Apache Kafka

Champs

Description

sink.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 : ”org.apache.spark.sql.kafka010.KafkaSourceProvider$KafkaTable@e04b100”.

sink.numOutputRows

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

Champs

Description

sink.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 : ”org.apache.spark.sql.kafka010.KafkaSourceProvider$KafkaTable@e04b100”.

sink.numOutputRows

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

Python
{
"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

Python
{
"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

Python
{
"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

Python
{
"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

Python
{
"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
}
}