DataSourceStreamReader

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

Средства чтения потоков данных отвечают за вывод данных из источника потоковой передачи. Реализуйте этот класс и верните экземпляр из DataSource.streamReader() источника данных для чтения в качестве источника потоковой передачи.

Добавлено в Databricks Runtime 15.2

Синтаксис

from pyspark.sql.datasource import DataSourceStreamReader

class MyDataSourceStreamReader(DataSourceStreamReader):
    def initialOffset(self):
        ...

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

    def read(self, partition):
        ...

Методы

Метод Описание
initialOffset() Возвращает начальное смещение источника данных потоковой передачи в виде dict. Новый запрос потоковой передачи начинает чтение из этого смещения. Должен возвращать пары смещением ключа-значение для примитивных типов в формате JSON или dict формате. Вызывается PySparkNotImplementedError , если он не реализован.
latestOffset(start, limit) Возвращает последнее смещение, доступное как dictсмещение начала и ограничение чтения. Источник может вернуть то же смещение, что start и при отсутствии новых данных. Источник должен всегда уважать заданный limit. Должен возвращать пары смещением ключа-значение для примитивных типов в формате JSON или dict формате. Вызывается PySparkNotImplementedError , если он не реализован.
partitions(start, end) Возвращает последовательность InputPartition объектов, представляющих данные между start и end смещениями. Возвращает пустую последовательность, если start равно end. Каждый InputPartition представляет разделение данных, которое может обрабатываться одной задачей Spark.
read(partition) Создает данные для заданной секции и возвращает итератор кортежей, строк или объектов PyArrow RecordBatch . Каждый кортеж или строка преобразуется в строку в окончательном кадре данных. Этот метод является абстрактным и должен быть реализован.
commit(end) Сообщает источнику, что Spark завершил обработку всех данных для смещения меньше или равно end. Spark будет запрашивать смещение только больше, чем end в будущем.
stop() Останавливает источник и освобождает все выделенные ресурсы. Вызывается при завершении потокового запроса.

Примечания

  • read() является статическим и без отслеживания состояния. Не получите доступ к изменяемым членам класса или не сохраняйте состояние в памяти между различными вызовами read().
  • Все значения секций, возвращаемые partitions() объектами, должны быть выбранными.
  • Смещения представлены как dict рекурсивные dict , ключи и значения которых являются примитивными типами: целочисленное, строковое или логическое значение.

Примеры

Реализуйте средство чтения потоковой передачи, которое считывается из последовательности индексированных записей:

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")