DataSourceStreamWriter

Veri akışı yazarları için temel sınıf.

Veri akışı yazarları bir akış havuzuna veri yazmakla sorumludur. Bu sınıfı uygulayın ve bir veri kaynağını akış havuzu olarak yazılabilir hale getirmek için öğesinden DataSource.streamWriter() bir örnek döndürür. write() her mikrobatch için yürütücülerde çağrılır ve commit() veya abort() mikrobatch içindeki tüm görevler tamamlandıktan sonra sürücüde çağrılır.

Sözdizimi

from pyspark.sql.datasource import DataSourceStreamWriter

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

Methods

Yöntem Açıklama
write(iterator) Akış havuzuna veri yazar. Her mikrobatch için yürütücülere bir kez çağrılır. Nesnelerin yineleyicisini Row kabul eder ve bir WriterCommitMessageveya None işleme iletisi yoksa döndürür. Bu yöntem soyut ve uygulanması gerekir.
commit(messages, batchId) Tüm yürütücülerden toplanan işleme iletilerinin listesini kullanarak mikrobatch'i işler. Mikrobatch'teki tüm görevler başarıyla çalıştırıldığında sürücüde çağrılır.
abort(messages, batchId) Tüm yürütücülerden toplanan işleme iletilerinin listesini kullanarak mikrobatch'i durdurur. Mikrobatch'teki bir veya daha fazla görev başarısız olduğunda sürücüde çağrılır.

Notlar

  • Sürücü, tüm yürütücülerden işleme iletilerini toplar ve tüm görevlerin başarılı olması veya herhangi bir görevin başarısız olması durumunda bu commit() iletileri iletir abort() .
  • Yazma görevi başarısız olursa, işleme iletisi veya Noneöğesine commit() geçirilen listede yer abort() alır.
  • batchId her mikrobatch'i benzersiz olarak tanımlar ve her mikrobatch işlenirken 1 artırır.

Örnekler

Bir dosyaya satır ekleyen bir akış yazıcısı uygulayın:

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")