Aller au contenu principal

foreach (DataStreamWriter)

Définit le résultat de la query de streaming à traiter à l’aide du writer fourni. La logique de traitement peut être spécifiée comme une fonction qui prend une ligne en entrée, ou comme un objet avec process(row) et des méthodes open(partition_id, epoch_id) et close(error) facultatives.

Syntaxe

foreach(f)

parameter

parameter

Type

Description

f

appelable ou objet

Une fonction qui prend une Row en entrée, ou un objet avec une méthode process(row) et des méthodes open et close facultatives.

parameter

Type

Description

f

appelable ou objet

Une fonction qui prend une Row en entrée, ou un objet avec une méthode process(row) et des méthodes open et close facultatives.

Renvoie

DataStreamWriter

Notes

L’objet fourni doit être sérialisable. Toute initialisation pour l'écriture de données (par exemple, l'ouverture d'une connexion) doit être effectuée à l'intérieur de open(), et non au moment de la construction.

Exemples

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

Traiter chaque ligne à l'aide d'une fonction :

Python
def print_row(row):
print(row)

q = df.writeStream.foreach(print_row).start()
time.sleep(3)
q.stop()

Traitez chaque ligne à l'aide d'un objet avec les méthodes open, process et close :

Python
class RowPrinter:
def open(self, partition_id, epoch_id):
print("Opened %d, %d" % (partition_id, epoch_id))
return True

def process(self, row):
print(row)

def close(self, error):
print("Closed with error: %s" % str(error))

q = df.writeStream.foreach(RowPrinter()).start()
time.sleep(3)
q.stop()