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
from pyspark.sql.datasource import DataSourceReader
class MyDataSourceReader(DataSourceReader):
def read(self, partition):
...
Méthodes
Méthode | Description |
|---|---|
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. | |
Renvoie une séquence de | |
Génère des données pour une partition donnée et renvoie un itérateur de tuples, de lignes ou d'objets PyArrow |
Exemples
Implémentez un lecteur de base qui renvoie les lignes d’une liste de partitions :
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:
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 :
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):
...