DataSourceStreamArrowWriter

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

Spark DataSourceStreamWriter nesnelerinin yineleyicisiyle çalışan öğesinden farklı Rowolarak, bu sınıf akış verileri yazılırken Ok biçimi için iyileştirilmiştir. Akış kullanım örnekleri için Yerel olarak Ok'un 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 akış havuzu olarak yazılabilir hale getirmek için öğesinden DataSource.streamWriter() bir örnek döndürür.

Sözdizimi

from pyspark.sql.datasource import DataSourceStreamArrowWriter

class MyDataSourceStreamArrowWriter(DataSourceStreamArrowWriter):
    def write(self, iterator):
        ...

Methods

Yöntem Açıklama
write(iterator) Akış havuzuna PyArrow RecordBatch nesnelerinin yineleyicisini yazar. Her mikrobatch için yürütücülere 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, 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. DataSourceStreamWriter öğesinden devralındı.
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. DataSourceStreamWriter öğ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.
  • batchId her mikrobatch'i benzersiz olarak tanımlar ve her mikrobatch işlenirken 1 artırır.

Örnekler

Mikrobatch başına satırları sayan Ok tabanlı bir akış yazıcısı uygulayın:

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

@dataclass
class MyCommitMessage(WriterCommitMessage):
    num_rows: int

class MyDataSourceStreamArrowWriter(DataSourceStreamArrowWriter):
    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, 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")