Aller au contenu principal

DataSourceReader

Une classe de base pour les lecteurs de source de données.

Les lecteurs de sources de données sont responsables de la sortie de données depuis une source de données. Implémentez cette classe et renvoyez une instance de DataSource.reader() pour rendre une source de données lisible.

Syntaxe

Python
from pyspark.sql.datasource import DataSourceReader

class MyDataSourceReader(DataSourceReader):
def read(self, partition):
...

Méthodes

Méthode

Description

pushFilters(filters)

Appelée avec la liste des filtres qui peuvent être transmis à la source de données. Renvoie un itérable de filtres qui doivent encore être évalués par Spark. By default, renvoie tous les filtres, indiquant qu'aucun filtre n'est transmis. pushFilters() est autorisé à modifier self. L'objet doit rester sérialisable après modification. Les modifications apportées à self sont visibles par partitions() et read().

partitions()

Renvoie une séquence de InputPartition objets qui divisent la lecture des données en tâches parallèles. Par default, renvoie une seule partition. Remplacer pour de meilleures performances lors de la lecture de grands datasets. Toutes les valeurs de partition renvoyées par partitions() doivent être des objets sérialisables.

read(partition)

Génère des données pour une partition donnée et renvoie un itérateur de tuples, de lignes ou d'objets PyArrow RecordBatch. Chaque tuple ou ligne est converti(e) en ligne dans le DataFrame final. Cette méthode est abstraite et doit être implémentée.

Méthode

Description

pushFilters(filters)

Appelée avec la liste des filtres qui peuvent être transmis à la source de données. Renvoie un itérable de filtres qui doivent encore être évalués par Spark. By default, renvoie tous les filtres, indiquant qu'aucun filtre n'est transmis. pushFilters() est autorisé à modifier self. L'objet doit rester sérialisable après modification. Les modifications apportées à self sont visibles par partitions() et read().

partitions()

Renvoie une séquence de InputPartition objets qui divisent la lecture des données en tâches parallèles. Par default, renvoie une seule partition. Remplacer pour de meilleures performances lors de la lecture de grands datasets. Toutes les valeurs de partition renvoyées par partitions() doivent être des objets sérialisables.

read(partition)

Génère des données pour une partition donnée et renvoie un itérateur de tuples, de lignes ou d'objets PyArrow RecordBatch. Chaque tuple ou ligne est converti(e) en ligne dans le DataFrame final. Cette méthode est abstraite et doit être implémentée.

Exemples

Implémentez un lecteur de base qui renvoie les lignes d’une liste de partitions :

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

Retourner les lignes en utilisant PyArrow RecordBatch:

Python
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

Implémenter la descente de filtre pour prendre en charge les filtres EqualTo :

Python
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):
...