Aller au contenu principal

StreamingQueryManager

Gère toutes les instances StreamingQuery actives associées à un SparkSession. Utilisez spark.streams pour y accéder.

Syntaxe​

Python
# Access through SparkSession
spark.streams

Propriétés​

Propriété

Description

active

Retourne une liste de toutes les requêtes de streaming actives associées à ce SparkSession.

Propriété

Description

active

Retourne une liste de toutes les requêtes de streaming actives associées à ce SparkSession.

Méthodes​

Méthode

Description

get(id)

Retourne une query active par son ID unique.

awaitAnyTermination(timeout)

Attend qu'une query active se termine, ou que le délai d'expiration soit atteint.

resetTerminated()

Oublie les requêtes terminées passées afin que awaitAnyTermination() puisse être utilisé à nouveau pour attendre de nouvelles terminaisons.

addListener(listener)

Enregistre un StreamingQueryListener pour recevoir les rappels d'événements de cycle de vie.

removeListener(listener)

Désenregistre un StreamingQueryListener.

Méthode

Description

get(id)

Retourne une query active par son ID unique.

awaitAnyTermination(timeout)

Attend qu'une query active se termine, ou que le délai d'expiration soit atteint.

resetTerminated()

Oublie les requêtes terminées passées afin que awaitAnyTermination() puisse être utilisé à nouveau pour attendre de nouvelles terminaisons.

addListener(listener)

Enregistre un StreamingQueryListener pour recevoir les rappels d'événements de cycle de vie.

removeListener(listener)

Désenregistre un StreamingQueryListener.

Exemples​

Python
sdf = spark.readStream.format("rate").load()
sq = sdf.writeStream.format('memory').queryName('this_query').start()
sqm = spark.streams
[q.name for q in sqm.active]
# ['this_query']
sqm.awaitAnyTermination(5)
# True
sq.stop()
sqm.resetTerminated()