Aller au contenu principal

write (DataSourceArrowWriter)

Écrit un itérateur d'objets PyArrow RecordBatch vers le récepteur.

Cette méthode est appelée une fois sur chaque exécuteur pour écrire des données dans la source de données. Il accepte un itérateur d'objets PyArrow RecordBatch et renvoie une seule ligne représentant un message de commit, ou None s'il n'y a pas de message de commit.

Le Driver recueille les messages de commit, le cas échéant, de tous les exécuteurs et les transmet à la méthode commit() si toutes les tâches s'exécutent correctement. Si une tâche échoue, la méthode abort() sera appelée avec les messages de commit recueillis.

Syntaxe

write(iterator: Iterator[RecordBatch])

parameter

parameter

Type

Description

iterator

Itérateur[RecordBatch]

Un itérateur d'objets RecordBatch PyArrow représentant les données d'entrée.

parameter

Type

Description

iterator

Itérateur[RecordBatch]

Un itérateur d'objets RecordBatch PyArrow représentant les données d'entrée.

Renvoie

WriterCommitMessage

Un message de commit sérialisable.

Exemples

Python
from dataclasses import dataclass

@dataclass
class MyCommitMessage(WriterCommitMessage):
num_rows: int

def write(self, iterator: Iterator["RecordBatch"]) -> "WriterCommitMessage":
total_rows = 0
for batch in iterator:
total_rows += len(batch)
return MyCommitMessage(num_rows=total_rows)