Aller au contenu principal

SimpleDataSourceStreamReader

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

Comparé à DataSourceStreamReader, SimpleDataSourceStreamReader ne nécessite pas de planifier les partitions de données. La méthode read() permet de lire les données et de planifier le décalage le plus récent en même temps.

Comme SimpleDataSourceStreamReader lit les enregistrements dans le driver Spark pour déterminer le décalage de fin de chaque batch sans partitionnement, il n'est adapté qu'aux cas d'utilisation légers où le débit d'entrée et la taille de batch sont faibles. Utilisez DataSourceStreamReader lorsque le throughput de lecture est élevé et ne peut pas être géré par un seul process.

Ajouté dans Databricks Runtime 15.3

Syntaxe

Python
from pyspark.sql.datasource import SimpleDataSourceStreamReader

class MyStreamReader(SimpleDataSourceStreamReader):
def initialOffset(self):
return {"offset": 0}

def read(self, start):
...

def readBetweenOffsets(self, start, end):
...

Méthodes

Méthode

Description

initialOffset()

Renvoie le décalage initial de la source de données de streaming. Une nouvelle query de streaming start à lire à partir de ce décalage.

read(start)

Lit toutes les données disponibles à partir du décalage de start et renvoie un tuple d'un itérateur d'enregistrements et le décalage de fin pour la prochaine tentative de lecture.

readBetweenOffsets(start, end)

Lit toutes les données disponibles entre des offsets start et fin spécifiques. Invoqué lors de la reprise après sinistre pour relire un batch de manière déterministe.

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.

Méthode

Description

initialOffset()

Renvoie le décalage initial de la source de données de streaming. Une nouvelle query de streaming start à lire à partir de ce décalage.

read(start)

Lit toutes les données disponibles à partir du décalage de start et renvoie un tuple d'un itérateur d'enregistrements et le décalage de fin pour la prochaine tentative de lecture.

readBetweenOffsets(start, end)

Lit toutes les données disponibles entre des offsets start et fin spécifiques. Invoqué lors de la reprise après sinistre pour relire un batch de manière déterministe.

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.

Exemples

Définissez un lecteur de source de données de streaming simplifié personnalisé :

Python
from pyspark.sql.datasource import DataSource, SimpleDataSourceStreamReader

class MyStreamingDataSource(DataSource):
@classmethod
def name(cls):
return "my_streaming_source"

def schema(self):
return "value STRING"

def simpleStreamReader(self, schema):
return MySimpleStreamReader()

class MySimpleStreamReader(SimpleDataSourceStreamReader):
def initialOffset(self):
return {"partition-1": {"index": 0}}

def read(self, start):
end = {"partition-1": {"index": start["partition-1"]["index"] + 1}}
def records():
yield ("hello",)
return records(), end

def readBetweenOffsets(self, start, end):
def records():
yield ("hello",)
return records()

def commit(self, end):
pass

spark.dataSource.register(MyStreamingDataSource)
df = spark.readStream.format("my_streaming_source").load()