Aller au contenu principal

Source de données

Classe de base pour les source de données.

Cette classe représente une source de données personnalisée qui permet d’y lire des données et/ou d’y écrire des données. La source de données fournit des méthodes pour créer des lecteurs et des rédacteurs pour la lecture et l’écriture de données, respectivement. Au moins l'une des méthodes reader() ou writer() doit être implémentée par toute sous-classe afin de rendre la source de données lisible ou modifiable (voire les deux).

Après avoir implémenté cette interface, vous pouvez charger votre source de données à l'aide de spark.read.format(...).load() et enregistrer les données à l'aide de df.write.format(...).save().

Pour plus d'informations, consultez sources de données personnalisées PySpark.

Syntaxe

Python
from pyspark.sql.datasource import DataSource

class MyDataSource(DataSource):
@classmethod
def name(cls):
return "my_data_source"

parameter

parameter

Type

Description

options

dict

Un dictionnaire insensible à la casse représentant les options pour cette source de données.

parameter

Type

Description

options

dict

Un dictionnaire insensible à la casse représentant les options pour cette source de données.

Méthodes

Méthode

Description

name()

Renvoie une chaîne représentant le nom de format de cette source de données. By default, returns the class name. Remplacer pour fournir un nom abrégé personnalisé.

schema()

Renvoie le schéma de la source de données en tant que StructType ou chaîne DDL. Si aucun schéma n'est implémenté et qu'aucun schéma n'est fourni par l'utilisateur, une exception est levée.

reader(schema)

Renvoie une instance DataSourceReader pour la lecture des données. Requis pour les sources de données lisibles.

writer(schema, overwrite)

Renvoie une instance DataSourceWriter pour l'écriture de données. Requis pour les sources de données inscriptibles.

streamWriter(schema, overwrite)

Renvoie une instance DataSourceStreamWriter pour l'écriture de données dans un récepteur de streaming. Requis pour les sources de données de streaming inscriptibles.

simpleStreamReader(schema)

Renvoie une instance SimpleDataSourceStreamReader pour la lecture des données en streaming. Utilisé uniquement lorsque streamReader() n'est pas implémenté.

streamReader(schema)

Renvoie une instance DataSourceStreamReader pour la lecture des données en streaming. Prend la priorité sur simpleStreamReader().

Méthode

Description

name()

Renvoie une chaîne représentant le nom de format de cette source de données. By default, returns the class name. Remplacer pour fournir un nom abrégé personnalisé.

schema()

Renvoie le schéma de la source de données en tant que StructType ou chaîne DDL. Si aucun schéma n'est implémenté et qu'aucun schéma n'est fourni par l'utilisateur, une exception est levée.

reader(schema)

Renvoie une instance DataSourceReader pour la lecture des données. Requis pour les sources de données lisibles.

writer(schema, overwrite)

Renvoie une instance DataSourceWriter pour l'écriture de données. Requis pour les sources de données inscriptibles.

streamWriter(schema, overwrite)

Renvoie une instance DataSourceStreamWriter pour l'écriture de données dans un récepteur de streaming. Requis pour les sources de données de streaming inscriptibles.

simpleStreamReader(schema)

Renvoie une instance SimpleDataSourceStreamReader pour la lecture des données en streaming. Utilisé uniquement lorsque streamReader() n'est pas implémenté.

streamReader(schema)

Renvoie une instance DataSourceStreamReader pour la lecture des données en streaming. Prend la priorité sur simpleStreamReader().

Exemples

Définissez et enregistrez une source de données lisible personnalisée :

Python
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()

Définissez une source de données avec un schéma StructType :

Python
from pyspark.sql.types import StructType, StructField, IntegerType, StringType

class MyDataSource(DataSource):
def schema(self):
return StructType().add("a", "int").add("b", "string")