Lire les informations d'état du Structured Streaming
Vous pouvez utiliser les opérations de DataFrame ou les fonctions de table à valeurs SQL pour interroger les données et métadonnées d'état de Structured Streaming. Utilisez ces fonctions pour observer les informations d'état des requêtes avec état de Structured Streaming, ce qui peut être utile pour le monitoring et le debugging.
Vous devez disposer d'un accès en lecture au chemin du point de contrôle pour une query en streaming afin d'interroger les données d'état ou les métadonnées. Les fonctions décrites dans cet article permettent un accès en lecture seule aux données d'état et aux métadonnées. Vous ne pouvez utiliser que la sémantique de lecture par batch pour interroger les informations d'état.
Vous ne pouvez pas interroger les informations d'état des LakeFlow Pipelines, des tables de streaming ou des vues matérialisées. Vous ne pouvez pas query les informations d'état à l'aide de compute Serverless ou de compute configuré avec le mode d'accès standard.
Exigences
-
Utilisez l'une des configurations de compute suivantes :
- Databricks Runtime 16.3 et versions ultérieures sur un compute configuré avec le mode d'accès standard.
- Databricks Runtime 14.3 LTS et versions ultérieures sur un compute configuré avec un mode d’accès dédié ou sans isolation.
-
Accès en lecture au chemin du point de contrôle utilisé par la query de streaming.
Lire le magasin d'état du Structured Streaming
Vous pouvez lire les informations de l'état du store pour les queries Structured Streaming exécutées dans n'importe quel Databricks Runtime pris en charge. Utilisez la syntaxe suivante :
- Python
- Scala
- SQL
df = (spark.read
.format("statestore")
.load("/checkpoint/path"))
val df = spark.read
.format("statestore")
.load("/checkpoint/path")
SELECT * FROM read_statestore('/checkpoint/path')
Options et schéma de l'API de lecture d'état
Pour une liste complète des options de format statestore, consultez State store.
Les données de sortie ont le schéma suivant :
Colonne | Type | Description |
|---|---|---|
| Structure (type dérivé de la clé d'état) | La clé d'un enregistrement d'opérateur avec état dans le point de contrôle de l'état. |
| Struct (type supplémentaire dérivé de la valeur d'état) | La valeur pour un enregistrement d'opérateur avec état dans le point de contrôle d'état. |
| Entier | La partition du point de contrôle d'état qui contient l'enregistrement d'opérateur avec état. |
Dans Databricks Runtime 16.4 LTS et versions supérieures, lorsque l'option readChangeFeed est définie sur true, les données de sortie ont le schéma suivant :
Colonne | Type | Description |
|---|---|---|
| Long | L'identifiant du batch auquel le changement d'état appartient. |
| Chaîne | Le type de modification appliqué par le batch : |
| Structure (type dérivé de la clé d'état) | La clé d'un enregistrement d'opérateur avec état dans le point de contrôle de l'état. |
| Struct (type supplémentaire dérivé de la valeur d'état) | La valeur d'un enregistrement d'opérateur avec état dans le point de contrôle d'état. |
| Entier | La partition du point de contrôle d'état qui contient l'enregistrement d'opérateur avec état. |
Consultez read_statestore fonction à valeur de table.
Lire les modifications de l'état de Structured Streaming
Disponible sur Databricks Runtime 16.4 LTS et versions ultérieures. Pour savoir comment l'état change entre les micro-lots au lieu de visualiser l'état complet à un seul micro-lot, définissez readChangeFeed sur true et spécifiez changeStartBatchId. Vous pouvez éventuellement spécifier changeEndBatchId. Pour une liste complète des options, reportez-vous à Magasin d'état.
Par exemple, pour lire les modifications d’état du batch 2 jusqu’au dernier batch validé :
- Python
- Scala
- SQL
df = (spark.read
.format("statestore")
.option("readChangeFeed", True)
.option("changeStartBatchId", 2)
.load("<checkpointLocation>")
)
val df = spark.read
.format("statestore")
.option("readChangeFeed", true)
.option("changeStartBatchId", 2)
.load("<checkpointLocation>")
SELECT * FROM read_statestore(
'<checkpointLocation>',
readChangeFeed => true,
changeStartBatchId => 2
);
Le schéma de sortie inclut des colonnes batch_id et change_type supplémentaires. Pour le schéma complet, voir Options et schéma de l'API de lecture d'état.
Lire les métadonnées d'état de Structured Streaming
Disponible sur Databricks Runtime 14.3 LTS ou version supérieure. Vous pouvez lire les informations de métadonnées d'état pour les requêtes Structured Streaming :
- Python
- Scala
- SQL
df = (spark.read
.format("state-metadata")
.load("<checkpointLocation>"))
val df = spark.read
.format("state-metadata")
.load("<checkpointLocation>")
SELECT * FROM read_state_metadata('/checkpoint/path')
Les données renvoyées ont le schema suivant :
Colonne | Type | Description |
|---|---|---|
| Entier | L'identifiant entier de l'opérateur de streaming avec état. |
| Chaîne | Nom de l'opérateur de streaming avec état. |
| Chaîne | Nom du stockage d'état de l'opérateur. |
| Entier | Nombre de partitions du magasin d'état. |
| Long | L'ID de batch minimum disponible pour l'état de la query. |
| Long | ID de batch maximal disponible pour l’interrogation de l’état. |
Les valeurs d'ID de batch fournies par minBatchId et maxBatchId reflètent l'état au moment où le point de contrôle a été écrit. Les anciens batchs sont automatiquement nettoyés avec l'exécution de micro-batchs, la valeur fournie ici ne peut donc pas être garantie comme étant toujours disponible.
Consultez read_state_metadata fonction à valeur de table.
Exemple : Query one side of a Stream-Stream join
Utilisez la syntaxe suivante pour interroger le côté gauche d'une jointure de Stream à Stream :
- Python
- Scala
- SQL
left_df = (spark.read
.format("statestore")
.option("joinSide", "left")
.load("/checkpoint/path"))
val leftDf = spark.read
.format("statestore")
.option("joinSide", "left")
.load("/checkpoint/path")
SELECT * FROM read_statestore(
'/checkpoint/path',
joinSide => 'left'
);
Exemple : Query state store pour Stream avec plusieurs opérateurs avec état
Cet exemple utilise le lecteur de métadonnées d'état pour recueillir les détails de métadonnées d'une query de streaming avec plusieurs opérateurs avec état, puis utilise les résultats de métadonnées comme options pour le lecteur d'état.
Le lecteur de métadonnées d'état prend le chemin du point de contrôle comme seule option, comme dans l'exemple de syntaxe suivant :
- Python
- Scala
- SQL
df = (spark.read
.format("state-metadata")
.load("<checkpointLocation>"))
val df = spark.read
.format("state-metadata")
.load("<checkpointLocation>")
SELECT * FROM read_state_metadata('/checkpoint/path')
Le tableau suivant représente un exemple de sortie des métadonnées du magasin d'état :
operatorId | Nom d'opérateur | stateStoreName | numPartitions | minBatchId | maxBatchId |
|---|---|---|---|---|---|
0 | stateStoreSave | par défaut | 200 | 0 | 13 |
1 | dedupeWithinWatermark | par défaut | 200 | 0 | 13 |
Pour obtenir les résultats de l'opérateur dedupeWithinWatermark, query le lecteur d'état avec l'option operatorId, comme dans l'exemple suivant :
- Python
- Scala
- SQL
left_df = (spark.read
.format("statestore")
.option("operatorId", 1)
.load("/checkpoint/path"))
val leftDf = spark.read
.format("statestore")
.option("operatorId", 1)
.load("/checkpoint/path")
SELECT * FROM read_statestore(
'/checkpoint/path',
operatorId => 1
);