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 kaynağı yazarları için temel sınıf.
Veri kaynağı yazarları, verileri bir veri kaynağına kaydetmekle sorumludur. Bu sınıfı uygulayın ve bir veri kaynağını yazılabilir hale getirmek için öğesinden DataSource.writer() bir örnek döndür.
Databricks Runtime 14.3 LTS'ye eklendi
Sözdizimi
from pyspark.sql.datasource import DataSourceWriter
class MyDataSourceWriter(DataSourceWriter):
def write(self, iterator):
...
Methods
| Yöntem | Açıklama |
|---|---|
write(iterator) |
Verileri veri kaynağına yazar. Her yürütücüde 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) |
Yazma işini, tüm yürütücülerden toplanan işleme iletilerinin listesini kullanarak işler. Tüm görevler başarıyla çalıştırıldığında sürücüde çağrılır. |
abort(messages) |
Tüm yürütücülerden toplanan işleme iletilerinin listesini kullanarak yazma işini durdurur. 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.
Örnekler
Satırları bir dosyaya kaydeden temel bir yazıcı uygulayın:
from dataclasses import dataclass
from pyspark.sql.datasource import DataSource, DataSourceWriter, WriterCommitMessage
@dataclass
class MyCommitMessage(WriterCommitMessage):
num_rows: int
class MyDataSourceWriter(DataSourceWriter):
def __init__(self, options):
self.path = options.get("path")
def write(self, iterator):
rows = list(iterator)
with open(self.path, "w") as f:
for row in rows:
f.write(str(row) + "\n")
return MyCommitMessage(num_rows=len(rows))
def commit(self, messages):
total = sum(m.num_rows for m in messages if m is not None)
print(f"Committed {total} rows")
def abort(self, messages):
print("Write job failed, performing cleanup")