DataSourceWriter

Veri kaynağı yazarları için temel sınıf.

Veri kaynağı yazarları, verileri bir veri kaynağına kaydetmekle sorumludur. Bu sınıfı uygulayın ve bir veri kaynağını yazılabilir hale getirmek için öğesinden DataSource.writer() bir örnek döndür.

Databricks Runtime 14.3 LTS'ye eklendi

Sözdizimi

from pyspark.sql.datasource import DataSourceWriter

class MyDataSourceWriter(DataSourceWriter):
    def write(self, iterator):
        ...

Methods

Yöntem Açıklama
write(iterator) Verileri veri kaynağına yazar. Her yürütücüde 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) Yazma işini, tüm yürütücülerden toplanan işleme iletilerinin listesini kullanarak işler. Tüm görevler başarıyla çalıştırıldığında sürücüde çağrılır.
abort(messages) Tüm yürütücülerden toplanan işleme iletilerinin listesini kullanarak yazma işini durdurur. 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.

Örnekler

Satırları bir dosyaya kaydeden temel bir yazıcı uygulayın:

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

@dataclass
class MyCommitMessage(WriterCommitMessage):
    num_rows: int

class MyDataSourceWriter(DataSourceWriter):
    def __init__(self, options):
        self.path = options.get("path")

    def write(self, iterator):
        rows = list(iterator)
        with open(self.path, "w") as f:
            for row in rows:
                f.write(str(row) + "\n")
        return MyCommitMessage(num_rows=len(rows))

    def commit(self, messages):
        total = sum(m.num_rows for m in messages if m is not None)
        print(f"Committed {total} rows")

    def abort(self, messages):
        print("Write job failed, performing cleanup")