Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Базовый класс для источников данных.
Этот класс представляет пользовательский источник данных, позволяющий читать данные из него и (или) записывать в него. Источник данных предоставляет методы для создания средств чтения и записи данных для чтения и записи данных соответственно. По крайней мере один из методов reader() или writer() должен быть реализован любым подклассом, чтобы сделать источник данных доступным для чтения или записи (или обоих).
После реализации этого интерфейса можно загрузить источник данных с помощью spark.read.format(...).load() и сохранить данные.df.write.format(...).save()
Дополнительные сведения см. в разделе "Пользовательские источники данных PySpark".
Синтаксис
from pyspark.sql.datasource import DataSource
class MyDataSource(DataSource):
@classmethod
def name(cls):
return "my_data_source"
Параметры
| Параметр | Тип | Описание |
|---|---|---|
options |
Дикт | Словарь без учета регистра, представляющий параметры этого источника данных. |
Методы
| Метод | Описание |
|---|---|
name() |
Возвращает строку, представляющую имя формата этого источника данных. По умолчанию возвращает имя класса. Переопределите, чтобы указать настраиваемое короткое имя. |
schema() |
Возвращает схему источника данных в виде StructType строки DDL. Если не реализована и схема не предоставляется пользователем, создается исключение. |
reader(schema) |
DataSourceReader Возвращает экземпляр для чтения данных. Требуется для доступных для чтения источников данных. |
writer(schema, overwrite) |
DataSourceWriter Возвращает экземпляр для записи данных. Требуется для записываемых источников данных. |
streamWriter(schema, overwrite) |
DataSourceStreamWriter Возвращает экземпляр для записи данных в приемник потоковой передачи. Требуется для источников данных потоковой передачи, доступных для записи. |
simpleStreamReader(schema) |
SimpleDataSourceStreamReader Возвращает экземпляр для чтения потоковых данных. Используется только в том случае, если streamReader() он не реализован. |
streamReader(schema) |
DataSourceStreamReader Возвращает экземпляр для чтения потоковых данных. Принимает приоритет simpleStreamReader()над . |
Примеры
Определите и зарегистрируйте пользовательский источник данных, доступный для чтения:
from pyspark.sql.datasource import DataSource, DataSourceReader, InputPartition
class MyDataSource(DataSource):
@classmethod
def name(cls):
return "my_data_source"
def schema(self):
return "a INT, b STRING"
def reader(self, schema):
return MyDataSourceReader(schema)
class MyDataSourceReader(DataSourceReader):
def read(self, partition):
yield (1, "hello")
yield (2, "world")
spark.dataSource.register(MyDataSource)
df = spark.read.format("my_data_source").load()
df.show()
Определите источник данных со схемой StructType :
from pyspark.sql.types import StructType, StructField, IntegerType, StringType
class MyDataSource(DataSource):
def schema(self):
return StructType().add("a", "int").add("b", "string")