Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Базовый класс для упрощенных средств чтения источников данных потоковой передачи.
По сравнению с DataSourceStreamReader, SimpleDataSourceStreamReader не требует планирования секций данных. Метод read() позволяет считывать данные и планировать последнее смещение одновременно.
Так как SimpleDataSourceStreamReader считывает записи в драйвере Spark для определения конечного смещения каждого пакета без секционирования, он подходит только для упрощенных вариантов использования, когда скорость ввода и размер пакета невелики. Используется DataSourceStreamReader при высокой пропускной способности чтения и не может обрабатываться одним процессом.
Добавлено в Databricks Runtime 15.3
Синтаксис
from pyspark.sql.datasource import SimpleDataSourceStreamReader
class MyStreamReader(SimpleDataSourceStreamReader):
def initialOffset(self):
return {"offset": 0}
def read(self, start):
...
def readBetweenOffsets(self, start, end):
...
Методы
| Метод | Описание |
|---|---|
initialOffset() |
Возвращает начальное смещение источника данных потоковой передачи. Новый запрос потоковой передачи начинает чтение из этого смещения. |
read(start) |
Считывает все доступные данные из смещения начала и возвращает кортеж итератора записей и конец смещения для следующей попытки чтения. |
readBetweenOffsets(start, end) |
Считывает все доступные данные между определенными смещениями начала и окончания. Вызывается во время восстановления сбоя для повторного чтения пакетной детерминированной. |
commit(end) |
Сообщает источнику, что Spark завершил обработку всех данных для смещения меньше или равно end. |
Примеры
Определите пользовательское упрощенное средство чтения источников данных потоковой передачи:
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()