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 |
|---|---|---|
| appelable ou objet | Une fonction qui prend une Row en entrée, ou un objet avec une méthode |
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
import time
df = spark.readStream.format("rate").load()
Traiter chaque ligne à l'aide d'une fonction :
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 :
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()