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.
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 iletirabort(). - Yazma görevi başarısız olursa, işleme iletisi veya
Noneöğesinecommit()geçirilen listede yerabort()alır. -
batchIdher 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")