Aller au contenu principal

DataStreamReader

Interface utilisée pour charger un DataFrame en streaming à partir de systèmes de stockage externes (par exemple, les systèmes de fichiers et les magasins clé-valeur). Utilisez spark.readStream pour y accéder.

Syntaxe

Python
# Access through SparkSession
spark.readStream

Méthodes

Méthode

Description

format(source)

Spécifie le format de la source de données d'entrée.

schema(schema)

Spécifie le schéma du DataFrame de streaming.

option(key, value)

Ajouter une option d'entrée pour la source de données sous-jacente.

options(**options)

Ajoute plusieurs options d'entrée pour la source de données sous-jacente.

load(path)

Charge le DataFrame de streaming à partir du chemin d'accès donné et le renvoie.

json(path)

Charge un fichier JSON Stream et renvoie un DataFrame.

orc(path)

Charge un Stream de fichiers ORC et renvoie un DataFrame.

parquet(path)

Charge un Stream de fichiers Parquet et renvoie un DataFrame.

text(path)

Charge un Stream de fichier texte et renvoie un DataFrame.

csv(path)

Charge un Stream de fichier CSV et renvoie un DataFrame.

xml(path)

Charge un Stream de fichiers XML et renvoie un DataFrame.

table(tableName)

Charge une table Delta de streaming et renvoie un DataFrame.

name(source_name)

Attribue un nom à la source de streaming pour l'évolution des points de contrôle.

changes(tableName)

Retourne les changements au niveau des lignes (Change Data Capture) de la table spécifiée sous forme de DataFrame de streaming.

Méthode

Description

format(source)

Spécifie le format de la source de données d'entrée.

schema(schema)

Spécifie le schéma du DataFrame de streaming.

option(key, value)

Ajouter une option d'entrée pour la source de données sous-jacente.

options(**options)

Ajoute plusieurs options d'entrée pour la source de données sous-jacente.

load(path)

Charge le DataFrame de streaming à partir du chemin d'accès donné et le renvoie.

json(path)

Charge un fichier JSON Stream et renvoie un DataFrame.

orc(path)

Charge un Stream de fichiers ORC et renvoie un DataFrame.

parquet(path)

Charge un Stream de fichiers Parquet et renvoie un DataFrame.

text(path)

Charge un Stream de fichier texte et renvoie un DataFrame.

csv(path)

Charge un Stream de fichier CSV et renvoie un DataFrame.

xml(path)

Charge un Stream de fichiers XML et renvoie un DataFrame.

table(tableName)

Charge une table Delta de streaming et renvoie un DataFrame.

name(source_name)

Attribue un nom à la source de streaming pour l'évolution des points de contrôle.

changes(tableName)

Retourne les changements au niveau des lignes (Change Data Capture) de la table spécifiée sous forme de DataFrame de streaming.

Exemples

Python
spark.readStream
# <...streaming.readwriter.DataStreamReader object ...>

Chargez un stream de taux, appliquez une transformation, écrivez sur la console et arrêtez après 3 secondes.

Python
import time
df = spark.readStream.format("rate").load()
df = df.selectExpr("value % 3 as v")
q = df.writeStream.format("console").start()
time.sleep(3)
q.stop()