DataSourceStreamReader

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ın read().
  • 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 dict veya özyinelemeli dict olarak 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")