Aller au contenu principal

DataSourceArrowWriter

Une classe de base pour les rédacteurs de sources de données qui traitent les données à l'aide de RecordBatch de PyArrow.

Contrairement à DataSourceWriter, qui fonctionne avec un itérateur d'objets Spark Row, cette classe est optimisée pour le format Arrow lors de l'écriture des données. Il peut offrir de meilleures performances lors de l'interfaçage avec des systèmes ou des bibliothèques qui supportent Arrow en mode natif. Implémentez cette classe et renvoyez une instance de DataSource.writer() pour rendre une source de données inscriptible à l'aide d'Arrow.

Syntaxe

Python
from pyspark.sql.datasource import DataSourceArrowWriter

class MyDataSourceArrowWriter(DataSourceArrowWriter):
def write(self, iterator):
...

Méthodes

Méthode

Description

write(iterator)

Écrit un itérateur d'objets PyArrow RecordBatch dans le récepteur. Appelé une fois sur chaque exécuteur. 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. Hérité de DataSourceWriter.

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é. Hérité de DataSourceWriter.

Méthode

Description

write(iterator)

Écrit un itérateur d'objets PyArrow RecordBatch dans le récepteur. Appelé une fois sur chaque exécuteur. 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. Hérité de DataSourceWriter.

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é. Hérité de DataSourceWriter.

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 writer basé sur Arrow qui compte les lignes dans tous les batchs :

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

@dataclass
class MyCommitMessage(WriterCommitMessage):
num_rows: int

class MyDataSourceArrowWriter(DataSourceArrowWriter):
def write(self, iterator):
total_rows = 0
for batch in iterator:
total_rows += len(batch)
return MyCommitMessage(num_rows=total_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")