Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Базовый класс для средств чтения источников данных.
Средства чтения источников данных отвечают за вывод данных из источника данных. Реализуйте этот класс и верните экземпляр, DataSource.reader() чтобы сделать источник данных читаемым.
Синтаксис
from pyspark.sql.datasource import DataSourceReader
class MyDataSourceReader(DataSourceReader):
def read(self, partition):
...
Методы
| Метод | Описание |
|---|---|
pushFilters(filters) |
Вызывается со списком фильтров, которые можно отправить в источник данных. Возвращает итератор фильтров, которые по-прежнему должны оцениваться Spark. По умолчанию возвращает все фильтры, указывающие, что фильтры не отправляются вниз.
pushFilters() разрешено изменять self. Объект должен оставаться выбранным после изменения.
self Изменения, которые будут видимы и partitions()read(). |
partitions() |
Возвращает последовательность InputPartition объектов, которые разделяют чтение данных на параллельные задачи. По умолчанию возвращается одна секция. Переопределите для повышения производительности при чтении больших наборов данных. Все значения секций, возвращаемые partitions() объектами, должны быть выбранными. |
read(partition) |
Создает данные для заданной секции и возвращает итератор кортежей, строк или объектов PyArrow RecordBatch . Каждый кортеж или строка преобразуется в строку в окончательном кадре данных. Этот метод является абстрактным и должен быть реализован. |
Примеры
Реализуйте базовое средство чтения, которое возвращает строки из списка разделов:
from pyspark.sql.datasource import DataSource, DataSourceReader, InputPartition
class MyDataSourceReader(DataSourceReader):
def partitions(self):
return [InputPartition(1), InputPartition(2), InputPartition(3)]
def read(self, partition):
yield (partition.value, 0)
yield (partition.value, 1)
Возвращать строки с помощью PyArrow RecordBatch:
class MyDataSourceReader(DataSourceReader):
def read(self, partition):
import pyarrow as pa
data = {
"partition": [partition.value] * 2,
"value": [0, 1]
}
table = pa.Table.from_pydict(data)
for batch in table.to_batches():
yield batch
Реализуйте pushdown фильтра для поддержки EqualTo фильтров:
from pyspark.sql.datasource import DataSourceReader, EqualTo
class MyDataSourceReader(DataSourceReader):
def __init__(self):
self.filters = []
def pushFilters(self, filters):
for f in filters:
if isinstance(f, EqualTo):
self.filters.append(f)
else:
yield f
def read(self, partition):
...