Not
Bu sayfaya erişim yetkilendirme gerektiriyor. Oturum açmayı veya dizinleri değiştirmeyi deneyebilirsiniz.
Bu sayfaya erişim yetkilendirme gerektiriyor. Dizinleri değiştirmeyi deneyebilirsiniz.
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)