write (DataSourceStreamArrowWriter)
Écrit un itérateur d’objets PyArrow RecordBatch dans le récepteur de streaming.
Cette méthode est appelée sur les exécuteurs pour écrire des données dans le récepteur de données de streaming dans chaque micro-lot. 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 |
|---|---|---|
| Itérateur[RecordBatch] | Un itérateur d'objets |
Renvoie
WriterCommitMessage
Un message de commit sérialisable.
Exemples
from dataclasses import dataclass
@dataclass
class MyCommitMessage(WriterCommitMessage):
num_rows: int
batch_id: 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, batch_id=self.current_batch_id)