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