write (DataSourceStreamArrowWriter)

Akış havuzuna PyArrow RecordBatch nesnelerinin yineleyicisini yazar.

Bu yöntem, yürütücülerde her mikrobatch'teki akış veri havuzuna veri yazmak için çağrılır. PyArrow RecordBatch nesnelerinin yineleyicisini kabul eder ve bir işleme iletisini temsil eden veya None işleme iletisi olmayan tek bir satır döndürür.

Sürücü varsa, tüm yürütücülerden işleme iletilerini toplar ve tüm görevler başarıyla çalıştırılırsa bunları yöntemine commit() geçirir. Herhangi bir görev başarısız olursa, abort() yöntemi toplanan işleme iletileriyle çağrılır.

Sözdizimi

write(iterator: Iterator[RecordBatch])

Parametreler

Parametre Türü Açıklama
iterator Yineleyici[RecordBatch] Giriş verilerini temsil eden PyArrow RecordBatch nesnelerinin yineleyicisi.

İadeler

WriterCommitMessage

Seri hale getirilebilir bir işleme iletisi.

Örnekler

from dataclasses import dataclass

@dataclass
class MyCommitMessage(WriterCommitMessage):
    num_rows: int
    batch_id: int

def write(self, iterator: Iterator["RecordBatch"]) -> "WriterCommitMessage":
    total_rows = 0
    for batch in iterator:
        total_rows += len(batch)
    return MyCommitMessage(num_rows=total_rows, batch_id=self.current_batch_id)