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
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 |
|---|---|
Appelé lorsqu'une query est start. | |
Appelée lorsqu'il y a une mise à jour du statut (taux d'ingestion mis à jour, etc.) | |
Appelé lorsque la query est inactive et attend de nouvelles données à traiter. | |
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
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())