DataSourceStreamReader
Une classe de base pour les lecteurs de source de données en streaming.
Les lecteurs de Stream de données sont responsables de la sortie des données d'une source de données de streaming. Implémentez cette classe et renvoyez une instance de DataSource.streamReader() pour rendre une source de données lisible en tant que source de streaming.
Ajouté dans Databricks Runtime 15.2
Syntaxe
from pyspark.sql.datasource import DataSourceStreamReader
class MyDataSourceStreamReader(DataSourceStreamReader):
def initialOffset(self):
...
def partitions(self, start, end):
...
def read(self, partition):
...
Méthodes
Méthode | Description |
|---|---|
Renvoie le décalage initial de la source de données de streaming en tant que | |
Renvoie le décalage le plus récent disponible en tant que | |
Retourne 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 | |
Informe la source que Spark a terminé le traitement de toutes les données pour les décalages inférieurs ou égaux à | |
Arrête la source et libère toutes les ressources qu'elle a allouées. Invoqué lorsque le streaming query se termine. |
Notes
read()est statique et sans état. N'accédez pas aux membres mutables de la classe et ne conservez pas d'état en mémoire entre les différentes invocations deread().- Toutes les valeurs de partition renvoyées par
partitions()doivent être des objets picklables. - Les offsets sont représentés comme un
dictou undictrécursif dont les clés et les valeurs sont des types primitifs : entier, chaîne de caractères ou booléen.
Exemples
Implémentez un lecteur de streaming qui lit à partir d'une séquence d'enregistrements indexés :
from pyspark.sql.datasource import (
DataSource,
DataSourceStreamReader,
InputPartition,
)
class MyDataSourceStreamReader(DataSourceStreamReader):
def initialOffset(self):
return {"index": 0}
def latestOffset(self, start, limit):
return {"index": start["index"] + 10}
def partitions(self, start, end):
return [
InputPartition(i)
for i in range(start["index"], end["index"])
]
def read(self, partition):
yield (partition.value, f"record-{partition.value}")
def commit(self, end):
print(f"Committed up to offset {end}")
def stop(self):
print("Stopping stream reader")