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