Classe DataFrameWriter
Interface utilisée pour écrire un DataFrame dans des systèmes de stockage externes (par exemple, systèmes de fichiers, magasins clé-valeur, etc).
Prend en charge Spark Connect
Syntaxe
Utilisez DataFrame.write pour accéder à cette interface.
Méthodes
Méthode | Description |
|---|---|
Spécifie le comportement lorsque les données ou la table existent déjà. | |
Spécifie la source de données de sortie sous-jacente. | |
Ajoute une option de sortie pour la source de données sous-jacente. | |
Ajoute des options de sortie pour la source de données sous-jacente. | |
Partitionne la sortie par les colonnes données sur le système de fichiers. | |
Regroupe la sortie par les colonnes données. | |
Trie la sortie de chaque compartiment par les colonnes données sur le système de fichiers. | |
Clusters les données par les colonnes données pour optimiser les performances des queries. | |
Enregistre le contenu du DataFrame dans une source de données. | |
Insère le contenu du DataFrame dans la table spécifiée. | |
Enregistre le contenu du DataFrame en tant que table spécifiée. | |
Enregistre le contenu du DataFrame au format JSON au chemin spécifié. | |
Enregistre le contenu du DataFrame au format Parquet au chemin spécifié. | |
Enregistre le contenu du DataFrame dans un fichier texte au chemin d'accès spécifié. | |
Enregistre le contenu du DataFrame au format CSV au chemin spécifié. | |
Enregistre le contenu du DataFrame au format XML dans le chemin spécifié. | |
Enregistre le contenu du DataFrame au format ORC au chemin spécifié. | |
Enregistre le contenu du **DataFrame** au format Excel au chemin spécifié. | |
Enregistre le contenu du DataFrame dans une table de base de données externe via JDBC. |
Modes d'enregistrement
La méthode mode() prend en charge les options suivantes :
- ajouter : Ajoute le contenu de ce DataFrame aux données existantes.
- Remplacer : remplacer les données existantes.
- **error** ou **errorifexists** : Lever une exception si les données existent déjà (default).
- ignorer : Ignorer cette Opération en silence si les données existent déjà.
Exemples
Écriture vers différentes sources de données
# Access DataFrameWriter through DataFrame
df = spark.createDataFrame([{"name": "Alice", "age": 30}])
df.write
# Write to JSON file
df.write.json("path/to/output.json")
# Write to CSV file with options
df.write.option("header", "true").csv("path/to/output.csv")
# Write to Parquet file
df.write.parquet("path/to/output.parquet")
# Write to a table
df.write.saveAsTable("table_name")
Utilisation du format et de la sauvegarde
# Specify format explicitly
df.write.format("json").save("path/to/output.json")
# With options
df.write.format("csv") \
.option("header", "true") \
.option("compression", "gzip") \
.save("path/to/output.csv")
Spécification du mode d'enregistrement
# Overwrite existing data
df.write.mode("overwrite").parquet("path/to/output.parquet")
# Append to existing data
df.write.mode("append").parquet("path/to/output.parquet")
# Ignore if data exists
df.write.mode("ignore").json("path/to/output.json")
# Error if data exists (default)
df.write.mode("error").csv("path/to/output.csv")
Partitionnement des données
# Partition by single column
df.write.partitionBy("year").parquet("path/to/output.parquet")
# Partition by multiple columns
df.write.partitionBy("year", "month").parquet("path/to/output.parquet")
# Partition with bucketing
df.write \
.bucketBy(10, "id") \
.sortBy("age") \
.saveAsTable("bucketed_table")
Écriture dans JDBC
# Write to database table
df.write.jdbc(
url="jdbc:postgresql://localhost:5432/mydb",
table="users",
mode="overwrite",
properties={"user": "myuser", "password": "mypassword"}
)
Chaînage de méthodes
# Chain multiple configuration methods
df.write \
.format("parquet") \
.mode("overwrite") \
.option("compression", "snappy") \
.partitionBy("year", "month") \
.save("path/to/output")
Écriture dans les tables
# Save as managed table
df.write.saveAsTable("my_table")
# Save as managed table with options
df.write \
.mode("overwrite") \
.format("parquet") \
.partitionBy("year") \
.saveAsTable("partitioned_table")
# Insert into existing table
df.write.insertInto("existing_table")
# Insert into existing table with overwrite
df.write.insertInto("existing_table", overwrite=True)