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()