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.
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()