Aller au contenu principal

DataSourceStreamArrowWriter

Une classe de base pour les enregistreurs de Stream de données qui traitent les données à l'aide du RecordBatch de PyArrow.

Contrairement à DataSourceStreamWriter, qui fonctionne avec un itérateur d'objets Spark Row, cette classe est optimisée pour le format Arrow lors de l'écriture de données en streaming. Il peut offrir de meilleures performances lors de l’interfaçage avec des systèmes ou des bibliothèques qui prennent en charge en mode natif Arrow pour les cas d’utilisation en streaming. Implémentez cette classe et renvoyez une instance de DataSource.streamWriter() pour rendre une source de données inscriptible en tant que *streaming sink* en utilisant Arrow.

Syntaxe

Python
from pyspark.sql.datasource import DataSourceStreamArrowWriter

class MyDataSourceStreamArrowWriter(DataSourceStreamArrowWriter):
def write(self, iterator):
...

Méthodes

Méthode

Description

write(iterator)

Écrit un itérateur d'objets PyArrow RecordBatch vers le récepteur de streaming. Appelé sur les exécuteurs une fois par microlot. 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. Hérité de DataSourceStreamWriter.

abort(messages, batchId)

Annule le microbatch à 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 du micro-lot ont échoué. Hérité de DataSourceStreamWriter.

Méthode

Description

write(iterator)

Écrit un itérateur d'objets PyArrow RecordBatch vers le récepteur de streaming. Appelé sur les exécuteurs une fois par microlot. 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. Hérité de DataSourceStreamWriter.

abort(messages, batchId)

Annule le microbatch à 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 du micro-lot ont échoué. Hérité de DataSourceStreamWriter.

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 graveur de Stream basé sur Arrow qui compte les lignes par micro-lot :

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

@dataclass
class MyCommitMessage(WriterCommitMessage):
num_rows: int

class MyDataSourceStreamArrowWriter(DataSourceStreamArrowWriter):
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, 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")