Poznámka:
Přístup k této stránce vyžaduje autorizaci. Můžete se zkusit přihlásit nebo změnit adresáře.
Přístup k této stránce vyžaduje autorizaci. Můžete zkusit změnit adresáře.
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 pokudabort()některý úkol selže. - Pokud úloha zápisu selže, zpráva potvrzení bude
Nonev seznamu předánacommit()neboabort(). -
batchIdjednoznač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")