Aller au contenu principal

DataSourceWriter

Une classe de base pour les rédacteurs de sources de données.

Les rédacteurs de source de données sont responsables de l’enregistrement des données dans une source de données. Implémentez cette classe et renvoyez une instance à partir de DataSource.writer() pour rendre une source de données inscriptible.

Ajouté dans Databricks Runtime 14.3 LTS

Syntaxe

Python
from pyspark.sql.datasource import DataSourceWriter

class MyDataSourceWriter(DataSourceWriter):
def write(self, iterator):
...

Méthodes

Méthode

Description

write(iterator)

Écrit des données dans la source de données. Appelé une fois sur chaque exécuteur. 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)

Valide le job d'écriture à 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 s'exécutent avec succès.

abort(messages)

Interrompt le job d'écriture à l'aide d'une liste de messages de commit collectés auprès de tous les exécuteurs. Invoqué sur le Driver lorsqu'une ou plusieurs tâches ont échoué.

Méthode

Description

write(iterator)

Écrit des données dans la source de données. Appelé une fois sur chaque exécuteur. 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)

Valide le job d'écriture à 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 s'exécutent avec succès.

abort(messages)

Interrompt le job d'écriture à l'aide d'une liste de messages de commit collectés auprès de tous les exécuteurs. Invoqué sur le Driver lorsqu'une ou plusieurs tâches 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().

Exemples

Implémenter un enregistreur de base qui sauvegarde les lignes dans un fichier :

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

@dataclass
class MyCommitMessage(WriterCommitMessage):
num_rows: int

class MyDataSourceWriter(DataSourceWriter):
def __init__(self, options):
self.path = options.get("path")

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

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

def abort(self, messages):
print("Write job failed, performing cleanup")