Aller au contenu principal

StreamingQueryListener

Une classe abstraite pour écouter les événements liés à StreamingQuery. Héritez de cette classe et implémentez ses méthodes abstraites pour recevoir les rappels d'événements de cycle de vie pour les requêtes en streaming.

Syntaxe​

Python
from pyspark.sql.streaming import StreamingQueryListener

class MyListener(StreamingQueryListener):
def onQueryStarted(self, event):
pass

def onQueryProgress(self, event):
pass

def onQueryIdle(self, event):
pass

def onQueryTerminated(self, event):
pass

Méthodes​

Méthode

Description

onQueryStarted(event)

Appelé lorsqu'une query est start.

onQueryProgress(event)

Appelée lorsqu'il y a une mise à jour du statut (taux d'ingestion mis à jour, etc.)

onQueryIdle(event)

Appelé lorsque la query est inactive et attend de nouvelles données à traiter.

onQueryTerminated(event)

Appelée lorsqu'une query est arrêtée, avec ou sans erreur.

Méthode

Description

onQueryStarted(event)

Appelé lorsqu'une query est start.

onQueryProgress(event)

Appelée lorsqu'il y a une mise à jour du statut (taux d'ingestion mis à jour, etc.)

onQueryIdle(event)

Appelé lorsque la query est inactive et attend de nouvelles données à traiter.

onQueryTerminated(event)

Appelée lorsqu'une query est arrêtée, avec ou sans erreur.

Notes​

Les méthodes ne sont pas thread-safe car elles peuvent être appelées à partir de différents threads.

En mode Spark Connect, l'écouteur n'a pas accès aux variables définies en dehors de celui-ci. Utilisez self.spark au lieu de spark pour accéder à la session dans l'écouteur en mode Connect.

Exemples​

Python
from pyspark.sql.streaming import StreamingQueryListener

class MyListener(StreamingQueryListener):
def onQueryStarted(self, event):
# Do something with event.
pass

def onQueryProgress(self, event):
# Do something with event.
pass

def onQueryIdle(self, event):
# Do something with event.
pass

def onQueryTerminated(self, event):
# Do something with event.
pass

spark.streams.addListener(MyListener())