DataSourceArrowWriter

PyArrow'un RecordBatchkullanarak verileri işleyen veri kaynağı yazarları için bir temel sınıf.

Spark DataSourceWriter nesnelerinin yineleyicisiyle çalışan öğesinden farklı olarakRow, bu sınıf veri yazarken Ok biçimi için iyileştirilmiştir. Arrow'un yerel olarak desteklenmesine neden olan sistem veya kitaplıklarla birlikte çalışırken daha iyi performans sunabilir. Bu sınıfı uygulayın ve bir veri kaynağını Ok kullanarak yazılabilir hale getirmek için öğesinden DataSource.writer() bir örnek döndür.

Sözdizimi

from pyspark.sql.datasource import DataSourceArrowWriter

class MyDataSourceArrowWriter(DataSourceArrowWriter):
    def write(self, iterator):
        ...

Methods

Yöntem Açıklama
write(iterator) Havuza PyArrow RecordBatch nesnelerinin yineleyicisini yazar. Her yürütücüde bir kez çağrılır. 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. DataSourceWriter öğesinden devralındı.
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. DataSourceWriter öğesinden devralındı.

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

Tüm toplu işlerde satırları sayan Ok tabanlı bir yazıcı uygulayın:

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

@dataclass
class MyCommitMessage(WriterCommitMessage):
    num_rows: int

class MyDataSourceArrowWriter(DataSourceArrowWriter):
    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):
        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")