DataSourceStreamWriter

Základní třída pro zapisovače datových proudů.

Za zápis dat do jímky streamování zodpovídají za zápis dat zapisovače datových proudů. Implementujte tuto třídu a vraťte instanci, DataSource.streamWriter() aby byl zdroj dat zapisovatelný jako jímka streamování. write() je volána na exekutory pro každý mikrobatch a commit() nebo abort() je volána na ovladači po dokončení všech úkolů v mikrobatchu.

Syntaxe

from pyspark.sql.datasource import DataSourceStreamWriter

class MyDataSourceStreamWriter(DataSourceStreamWriter):
    def write(self, iterator):
        ...

Methods

Metoda Description
write(iterator) Zapisuje data do jímky streamování. Volali jsme exekutory jednou za mikrobatch. Přijímá iterátor Row objektů a vrací WriterCommitMessagezprávu , nebo None pokud neexistuje žádná zpráva potvrzení. Tato metoda je abstraktní a musí být implementována.
commit(messages, batchId) Potvrdí mikrobatch pomocí seznamu zpráv potvrzení shromážděných ze všech exekutorů. Vyvolá se na ovladači, když se všechny úlohy v mikrobatchu úspěšně spustí.
abort(messages, batchId) Přeruší mikrobatch pomocí seznamu zpráv potvrzení shromážděných ze všech exekutorů. Vyvoláno na ovladači, když došlo k selhání jedné nebo více úloh v mikrobatchu.

Poznámky

  • Ovladač shromažďuje zprávy potvrzení ze všech exekutorů a předává je, commit() pokud jsou všechny úkoly úspěšné nebo pokud abort() některý úkol selže.
  • Pokud úloha zápisu selže, zpráva potvrzení bude None v seznamu předána commit() nebo abort().
  • batchId jednoznačně identifikuje jednotlivé mikrobatchy a přírůstky o 1 s každým zpracovaným mikrobatchem.

Příklady

Implementujte zapisovač streamu, který připojí řádky k souboru:

from dataclasses import dataclass
from pyspark.sql.datasource import DataSource, DataSourceStreamWriter, WriterCommitMessage

@dataclass
class MyCommitMessage(WriterCommitMessage):
    num_rows: int

class MyDataSourceStreamWriter(DataSourceStreamWriter):
    def __init__(self, options):
        self.path = options.get("path")

    def write(self, iterator):
        rows = list(iterator)
        with open(self.path, "a") as f:
            for row in rows:
                f.write(str(row) + "\n")
        return MyCommitMessage(num_rows=len(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")