SimpleDataSourceStreamReader

Базовый класс для упрощенных средств чтения источников данных потоковой передачи.

По сравнению с 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()