Aller au contenu principal

DataSourceStreamWriter

Une classe de base pour les enregistreurs de données Stream.

Les graveurs de flux de données sont chargés d'écrire des données dans un récepteur de streaming. Implémentez cette classe et renvoyez une instance de DataSource.streamWriter() pour rendre une source de données inscriptible en tant que récepteur de streaming. write() est appelé sur les exécuteurs pour chaque microbatch, et commit() ou abort() est appelé sur le Driver une fois que toutes les tâches du microbatch sont terminées.

Syntaxe

Python
from pyspark.sql.datasource import DataSourceStreamWriter

class MyDataSourceStreamWriter(DataSourceStreamWriter):
def write(self, iterator):
...

Méthodes

Méthode

Description

write(iterator)

Écrit des données dans le puits de streaming. Appelé sur les exécuteurs une fois par microlot. Accepte un itérateur de Row objets et renvoie un WriterCommitMessage, ou None s'il n'y a pas de message de commit. Cette méthode est abstraite et doit être mise en œuvre.

commit(messages, batchId)

Commits le micro-lot à l'aide d'une liste de messages de commit collectés auprès de tous les exécuteurs. Invoqué sur le Driver lorsque toutes les tâches du micro-lot s'exécutent avec succès.

abort(messages, batchId)

Annule le microbatch à l'aide d'une liste de messages de commit collectés auprès de tous les exécuteurs. Appelé sur le driver lorsqu'une ou plusieurs tâches du microbatch ont échoué.

Méthode

Description

write(iterator)

Écrit des données dans le puits de streaming. Appelé sur les exécuteurs une fois par microlot. Accepte un itérateur de Row objets et renvoie un WriterCommitMessage, ou None s'il n'y a pas de message de commit. Cette méthode est abstraite et doit être mise en œuvre.

commit(messages, batchId)

Commits le micro-lot à l'aide d'une liste de messages de commit collectés auprès de tous les exécuteurs. Invoqué sur le Driver lorsque toutes les tâches du micro-lot s'exécutent avec succès.

abort(messages, batchId)

Annule le microbatch à l'aide d'une liste de messages de commit collectés auprès de tous les exécuteurs. Appelé sur le driver lorsqu'une ou plusieurs tâches du microbatch ont échoué.

Notes

  • Le Driver collecte les messages de commit de tous les exécuteurs et les transmet à commit() si toutes les tâches réussissent, ou à abort() si une tâche échoue.
  • Si une tâche d'écriture échoue, son message de commit sera None dans la liste passée à commit() ou abort().
  • batchId identifie de manière unique chaque microbatch et s'incrémente de 1 à chaque microbatch traité.

Exemples

Implémentez un Stream Writer qui ajoute des lignes à un fichier :

Python
from dataclasses import dataclass
from pyspark.sql.datasource import DataSource, DataSourceStreamWriter, WriterCommitMessage

@dataclass
class MyCommitMessage(WriterCommitMessage):
num_rows: int

class MyDataSourceStreamWriter(DataSourceStreamWriter):
def __init__(self, options):
self.path = options.get("path")

def write(self, iterator):
rows = list(iterator)
with open(self.path, "a") as f:
for row in rows:
f.write(str(row) + "\n")
return MyCommitMessage(num_rows=len(rows))

def commit(self, messages, batchId):
total = sum(m.num_rows for m in messages if m is not None)
print(f"Committed batch {batchId} with {total} rows")

def abort(self, messages, batchId):
print(f"Batch {batchId} failed, performing cleanup")