Aller au contenu principal

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

Python
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

initialOffset()

Renvoie le décalage initial de la source de données de streaming en tant que dict. Une nouvelle query de streaming start à lire à partir de ce décalage. Doit renvoyer des paires clé-valeur de décalage de types primitifs au format JSON ou dict. Déclenche PySparkNotImplementedError si non implémenté.

latestOffset(start, limit)

Renvoie le décalage le plus récent disponible en tant que dict, étant donné un décalage de start et une limite de lecture. La source peut renvoyer le même décalage que start s'il n'y a pas de nouvelles données. La source doit toujours respecter le/la limit donné(e). Doit renvoyer des paires clé-valeur de décalage de types primitifs au format JSON ou dict. Génère une erreur PySparkNotImplementedError si elle n'est pas implémentée.

partitions(start, end)

Retourne une séquence de InputPartition objets représentant les données entre les décalages start et end. Retourne une séquence vide si start est égal à end. Chaque InputPartition représente un fractionnement de données qui peut être traité par une tâche Spark.

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.

commit(end)

Informe la source que Spark a terminé le traitement de toutes les données pour les décalages inférieurs ou égaux à end. Spark ne demandera que les décalages supérieurs à end à l'avenir.

stop()

Arrête la source et libère toutes les ressources qu'elle a allouées. Invoqué lorsque le streaming query se termine.

Méthode

Description

initialOffset()

Renvoie le décalage initial de la source de données de streaming en tant que dict. Une nouvelle query de streaming start à lire à partir de ce décalage. Doit renvoyer des paires clé-valeur de décalage de types primitifs au format JSON ou dict. Déclenche PySparkNotImplementedError si non implémenté.

latestOffset(start, limit)

Renvoie le décalage le plus récent disponible en tant que dict, étant donné un décalage de start et une limite de lecture. La source peut renvoyer le même décalage que start s'il n'y a pas de nouvelles données. La source doit toujours respecter le/la limit donné(e). Doit renvoyer des paires clé-valeur de décalage de types primitifs au format JSON ou dict. Génère une erreur PySparkNotImplementedError si elle n'est pas implémentée.

partitions(start, end)

Retourne une séquence de InputPartition objets représentant les données entre les décalages start et end. Retourne une séquence vide si start est égal à end. Chaque InputPartition représente un fractionnement de données qui peut être traité par une tâche Spark.

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.

commit(end)

Informe la source que Spark a terminé le traitement de toutes les données pour les décalages inférieurs ou égaux à end. Spark ne demandera que les décalages supérieurs à end à l'avenir.

stop()

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 de read().
  • Toutes les valeurs de partition renvoyées par partitions() doivent être des objets picklables.
  • Les offsets sont représentés comme un dict ou un dict ré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 :

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