SimpleDataSourceStreamReader

Basitleştirilmiş akış veri kaynağı okuyucuları için temel sınıf.

ile DataSourceStreamReaderSimpleDataSourceStreamReader karşılaştırıldığında, veri bölümlerinin planlanması gerekmez. yöntemi, read() verilerin okunmasına ve en son uzaklığı aynı anda planlamaya olanak tanır.

SimpleDataSourceStreamReader Bölümleme olmadan her toplu iş için bitiş uzaklığını belirlemek için Spark sürücüsündeki kayıtları okuduğundan, yalnızca giriş hızının ve toplu iş boyutunun küçük olduğu basit kullanım örnekleri için uygundur. Okuma aktarım hızı yüksek olduğunda kullanın DataSourceStreamReader ve tek bir işlem tarafından işlenemez.

Databricks Runtime 15.3'e eklendi

Sözdizimi

from pyspark.sql.datasource import SimpleDataSourceStreamReader

class MyStreamReader(SimpleDataSourceStreamReader):
    def initialOffset(self):
        return {"offset": 0}

    def read(self, start):
        ...

    def readBetweenOffsets(self, start, end):
        ...

Methods

Yöntem Açıklama
initialOffset() Akış veri kaynağının ilk uzaklığını döndürür. Yeni bir akış sorgusu bu uzaklıktan okumaya başlar.
read(start) Başlangıç uzaklığından tüm kullanılabilir verileri okur ve bir kayıt yineleyicisi demetini ve sonraki okuma denemesi için bitiş uzaklığını döndürür.
readBetweenOffsets(start, end) Belirli başlangıç ve bitiş uzaklıkları arasındaki tüm kullanılabilir verileri okur. Bir toplu işlemi belirleyici bir şekilde yeniden okumak için hata kurtarma sırasında çağrılır.
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.

Örnekler

Özel bir basitleştirilmiş akış veri kaynağı okuyucusu tanımlayın:

from pyspark.sql.datasource import DataSource, SimpleDataSourceStreamReader

class MyStreamingDataSource(DataSource):
    @classmethod
    def name(cls):
        return "my_streaming_source"

    def schema(self):
        return "value STRING"

    def simpleStreamReader(self, schema):
        return MySimpleStreamReader()

class MySimpleStreamReader(SimpleDataSourceStreamReader):
    def initialOffset(self):
        return {"partition-1": {"index": 0}}

    def read(self, start):
        end = {"partition-1": {"index": start["partition-1"]["index"] + 1}}
        def records():
            yield ("hello",)
        return records(), end

    def readBetweenOffsets(self, start, end):
        def records():
            yield ("hello",)
        return records()

    def commit(self, end):
        pass

spark.dataSource.register(MyStreamingDataSource)
df = spark.readStream.format("my_streaming_source").load()