Aller au contenu principal

trigger (DataStreamWriter)

Définit le Trigger pour la query de streaming. Si non défini, la query s'exécute aussi vite que possible, équivalent à processingTime='0 seconds'. Un seul paramètre de trigger peut être défini à la fois.

Pour plus d'informations, consultez Configurer les intervalles de trigger de Structured Streaming.

Syntaxe

trigger(*, processingTime=None, once=None, continuous=None, availableNow=None, realTime=None)

parameter

parameter

Type

Description

processingTime

str, facultatif

Une chaîne d'intervalle de temps de traitement (par exemple, '5 seconds', '1 minute'). Exécute une query de micro-lots périodiquement en fonction du temps de traitement.

once

bool, facultatif

Si True, ne traite qu'un seul batch de données, puis termine la query.

continuous

str, facultatif

Une chaîne d'intervalle de temps (par exemple, '5 seconds'). Exécute une query continue avec un intervalle de point de contrôle donné.

availableNow

bool, facultatif

Si True, traite toutes les données disponibles par lots, puis met fin à la query.

realTime

str, facultatif

Une chaîne de durée de batch (par exemple, '5 seconds'). Exécute une query en temps réel avec des batchs à la durée spécifiée.

parameter

Type

Description

processingTime

str, facultatif

Une chaîne d'intervalle de temps de traitement (par exemple, '5 seconds', '1 minute'). Exécute une query de micro-lots périodiquement en fonction du temps de traitement.

once

bool, facultatif

Si True, ne traite qu'un seul batch de données, puis termine la query.

continuous

str, facultatif

Une chaîne d'intervalle de temps (par exemple, '5 seconds'). Exécute une query continue avec un intervalle de point de contrôle donné.

availableNow

bool, facultatif

Si True, traite toutes les données disponibles par lots, puis met fin à la query.

realTime

str, facultatif

Une chaîne de durée de batch (par exemple, '5 seconds'). Exécute une query en temps réel avec des batchs à la durée spécifiée.

Renvoie

DataStreamWriter

Exemples

Python
df = spark.readStream.format("rate").load()

Exécution du Trigger toutes les 5 secondes :

Python
df.writeStream.trigger(processingTime='5 seconds')
# <...streaming.readwriter.DataStreamWriter object ...>

Trigger continuous execution every 5 secondes :

:::note Compatibilité Serverless

trigger(continuous=) n'est pas pris en charge sur le compute Serverless de Databricks. Pour les pipelines continus sur serverless, utilisez le mode continu Lakeflow Pipelines plutôt.

:::

Python
df.writeStream.trigger(continuous='5 seconds')
# <...streaming.readwriter.DataStreamWriter object ...>

Traitez toutes les données disponibles par lots :

Python
df.writeStream.trigger(availableNow=True)
# <...streaming.readwriter.DataStreamWriter object ...>

Trigger l'exécution en temps réel toutes les 5 secondes :

Python
df.writeStream.trigger(realTime='5 seconds')
# <...streaming.readwriter.DataStreamWriter object ...>