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.
Veri akışı yazarları için temel sınıf.
Veri akışı yazarları bir akış havuzuna veri yazmakla sorumludur. Bu sınıfı uygulayın ve bir veri kaynağını akış havuzu olarak yazılabilir hale getirmek için öğesinden DataSource.streamWriter() bir örnek döndürür.
write() her mikrobatch için yürütücülerde çağrılır ve commit() veya abort() mikrobatch içindeki tüm görevler tamamlandıktan sonra sürücüde çağrılır.
Sözdizimi
from pyspark.sql.datasource import DataSourceStreamWriter
class MyDataSourceStreamWriter(DataSourceStreamWriter):
def write(self, iterator):
...
Methods
| Yöntem | Açıklama |
|---|---|
write(iterator) |
Akış havuzuna veri yazar. Her mikrobatch için yürütücülere bir kez çağrılır. Nesnelerin yineleyicisini Row kabul eder ve 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. |
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. |
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
Bir dosyaya satır ekleyen bir akış yazıcısı uygulayın:
from dataclasses import dataclass
from pyspark.sql.datasource import DataSource, DataSourceStreamWriter, WriterCommitMessage
@dataclass
class MyCommitMessage(WriterCommitMessage):
num_rows: int
class MyDataSourceStreamWriter(DataSourceStreamWriter):
def __init__(self, options):
self.path = options.get("path")
def write(self, iterator):
rows = list(iterator)
with open(self.path, "a") as f:
for row in rows:
f.write(str(row) + "\n")
return MyCommitMessage(num_rows=len(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")