Aller au contenu principal

StreamingQuery

Un handle vers une query qui s'exécute en continu en arrière-plan à mesure que de nouvelles données arrivent. Toutes les méthodes sont thread-safe.

Syntaxe

Python
# Returned by DataStreamWriter.start() or DataStreamWriter.toTable()
q = df.writeStream.format("console").start()

Propriétés

Propriété

Description

id

Renvoie l'identifiant unique de cette query qui persiste après les redémarrages à partir des données de point de contrôle.

runId

Renvoie l'ID unique de cette query qui ne persiste pas après les redémarrages.

name

Renvoie le nom de la query spécifié par l'utilisateur, ou None s'il n'est pas spécifié.

isActive

Indique si cette query de streaming est actuellement active.

status

Renvoie l'état actuel de la query sous forme de dictionnaire.

recentProgress

Renvoie un tableau des StreamingQueryProgress mises à jour les plus récentes pour cette query.

lastProgress

Renvoie la mise à jour StreamingQueryProgress la plus récente, ou None s'il n'y a pas eu de mises à jour.

Propriété

Description

id

Renvoie l'identifiant unique de cette query qui persiste après les redémarrages à partir des données de point de contrôle.

runId

Renvoie l'ID unique de cette query qui ne persiste pas après les redémarrages.

name

Renvoie le nom de la query spécifié par l'utilisateur, ou None s'il n'est pas spécifié.

isActive

Indique si cette query de streaming est actuellement active.

status

Renvoie l'état actuel de la query sous forme de dictionnaire.

recentProgress

Renvoie un tableau des StreamingQueryProgress mises à jour les plus récentes pour cette query.

lastProgress

Renvoie la mise à jour StreamingQueryProgress la plus récente, ou None s'il n'y a pas eu de mises à jour.

Méthodes

Méthode

Description

awaitTermination(timeout)

Attend la fin de cette query, soit par stop(), soit par une exception.

processAllAvailable()

Bloque jusqu'à ce que toutes les données disponibles dans la source aient été traitées et validées vers la destination. Destiné aux tests.

stop()

Arrête cette query de streaming.

explain(extended)

Affiche les plans (logiques et physiques) dans la console pour le debugging.

exception()

Renvoie StreamingQueryException si la query s’est terminée par une exception, ou None.

Méthode

Description

awaitTermination(timeout)

Attend la fin de cette query, soit par stop(), soit par une exception.

processAllAvailable()

Bloque jusqu'à ce que toutes les données disponibles dans la source aient été traitées et validées vers la destination. Destiné aux tests.

stop()

Arrête cette query de streaming.

explain(extended)

Affiche les plans (logiques et physiques) dans la console pour le debugging.

exception()

Renvoie StreamingQueryException si la query s’est terminée par une exception, ou None.

Exemples

Python
sdf = spark.readStream.format("rate").load()
sq = sdf.writeStream.format('memory').queryName('this_query').start()
sq.isActive
# True
sq.name
# 'this_query'
sq.awaitTermination(5)
# False
sq.stop()
sq.isActive
# False