Aller au contenu principal

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.

remarque

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
df = (spark.read
.format("statestore")
.load("/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

key

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.

value

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.

partition_id

Entier

La partition du point de contrôle d'état qui contient l'enregistrement d'opérateur avec état.

Colonne

Type

Description

key

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.

value

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.

partition_id

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

batch_id

Long

L'identifiant du batch auquel le changement d'état appartient.

change_type

Chaîne

Le type de modification appliqué par le batch : update pour les insertions et les mises à jour, delete pour les suppressions.

key

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.

value

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. null pour les enregistrements où change_type est delete.

partition_id

Entier

La partition du point de contrôle d'état qui contient l'enregistrement d'opérateur avec état.

Colonne

Type

Description

batch_id

Long

L'identifiant du batch auquel le changement d'état appartient.

change_type

Chaîne

Le type de modification appliqué par le batch : update pour les insertions et les mises à jour, delete pour les suppressions.

key

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.

value

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. null pour les enregistrements où change_type est delete.

partition_id

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
df = (spark.read
.format("statestore")
.option("readChangeFeed", True)
.option("changeStartBatchId", 2)
.load("<checkpointLocation>")
)

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
df = (spark.read
.format("state-metadata")
.load("<checkpointLocation>"))

Les données renvoyées ont le schema suivant :

Colonne

Type

Description

operatorId

Entier

L'identifiant entier de l'opérateur de streaming avec état.

operatorName

Chaîne

Nom de l'opérateur de streaming avec état.

stateStoreName

Chaîne

Nom du stockage d'état de l'opérateur.

numPartitions

Entier

Nombre de partitions du magasin d'état.

minBatchId

Long

L'ID de batch minimum disponible pour l'état de la query.

maxBatchId

Long

ID de batch maximal disponible pour l’interrogation de l’état.

Colonne

Type

Description

operatorId

Entier

L'identifiant entier de l'opérateur de streaming avec état.

operatorName

Chaîne

Nom de l'opérateur de streaming avec état.

stateStoreName

Chaîne

Nom du stockage d'état de l'opérateur.

numPartitions

Entier

Nombre de partitions du magasin d'état.

minBatchId

Long

L'ID de batch minimum disponible pour l'état de la query.

maxBatchId

Long

ID de batch maximal disponible pour l’interrogation de l’état.

remarque

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
left_df = (spark.read
.format("statestore")
.option("joinSide", "left")
.load("/checkpoint/path"))

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
df = (spark.read
.format("state-metadata")
.load("<checkpointLocation>"))

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

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
left_df = (spark.read
.format("statestore")
.option("operatorId", 1)
.load("/checkpoint/path"))