Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Базовый класс для чтения источников данных потоковой передачи.
Средства чтения потоков данных отвечают за вывод данных из источника потоковой передачи. Реализуйте этот класс и верните экземпляр из 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")