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ış veri kaynağı okuyucuları için temel sınıf.
Veri kaynağı akış okuyucuları, akış veri kaynağından veri çıkışından sorumludur. Bu sınıfı uygulayın ve veri kaynağını akış kaynağı olarak okunabilir hale getirmek için öğesinden DataSource.streamReader() bir örnek döndür.
Databricks Runtime 15.2'ye eklendi
Sözdizimi
from pyspark.sql.datasource import DataSourceStreamReader
class MyDataSourceStreamReader(DataSourceStreamReader):
def initialOffset(self):
...
def partitions(self, start, end):
...
def read(self, partition):
...
Methods
| Yöntem | Açıklama |
|---|---|
initialOffset() |
Akış veri kaynağının ilk uzaklığını olarak dictdöndürür. Yeni bir akış sorgusu bu uzaklıktan okumaya başlar. JSON veya dict biçimde ilkel türlerin uzaklık anahtar-değer çiftlerini döndürmelidir.
PySparkNotImplementedError Uygulanmazsa yükseltir. |
latestOffset(start, limit) |
Başlangıç uzaklığı ve okuma sınırı göz önünde bulundurulduğunda, olarak dictkullanılabilen en son uzaklığı döndürür. Kaynak, yeni veri yokmuş gibi start aynı uzaklığı döndürebilir. Kaynak her zaman verilen limitöğesine saygı duymalıdır. JSON veya dict biçimde ilkel türlerin uzaklık anahtar-değer çiftlerini döndürmelidir.
PySparkNotImplementedError Uygulanmazsa yükseltir. |
partitions(start, end) |
ve InputPartition uzaklıkları arasındaki start verileri temsil eden bir nesne dizisi end döndürür. değerine eşitse startendboş bir sıra döndürür. Her InputPartition biri, bir Spark görevi tarafından işlenebilen bir veri bölmeyi temsil eder. |
read(partition) |
Belirli bir bölüm için veri oluşturur ve tanımlama kümeleri, satırlar veya PyArrow RecordBatch nesnelerinin yineleyicisini döndürür. Her tanımlama grubu veya satır, son DataFrame'deki bir satıra dönüştürülür. Bu yöntem soyut ve uygulanması gerekir. |
commit(end) |
Spark'ın değerinden küçük veya buna eşit uzaklıklar için tüm verileri işlemeyi endtamamladığını kaynağa bildirir. Spark yalnızca gelecektekinden daha büyük end uzaklıklar isteyecektir. |
stop() |
Kaynağı durdurur ve ayırdığı tüm kaynakları serbest bırakır. Akış sorgusu sonlandırıldığında çağrılır. |
Notlar
-
read()statiktir ve durum bilgisi yoktur. Değiştirilebilir sınıf üyelerine erişmeyin veya farklı çağrıları arasında bellek içi durumu tutmayınread(). - tarafından
partitions()döndürülen tüm bölüm değerleri seçilebilir nesneler olmalıdır. - Uzaklıklar, anahtarları ve değerleri ilkel türler olan bir
dictveya özyinelemelidictolarak temsil edilir: tamsayı, dize veya boole.
Örnekler
Dizine alınan kayıt dizisinden okuyan bir akış okuyucu uygulayın:
from pyspark.sql.datasource import (
DataSource,
DataSourceStreamReader,
InputPartition,
)
class MyDataSourceStreamReader(DataSourceStreamReader):
def initialOffset(self):
return {"index": 0}
def latestOffset(self, start, limit):
return {"index": start["index"] + 10}
def partitions(self, start, end):
return [
InputPartition(i)
for i in range(start["index"], end["index"])
]
def read(self, partition):
yield (partition.value, f"record-{partition.value}")
def commit(self, end):
print(f"Committed up to offset {end}")
def stop(self):
print("Stopping stream reader")