Lire et écrire des fichiers Parquet
Apache Parquet est un format de fichier en colonnes optimisé pour les charges de travail analytiques. Cela permet aux moteurs de query de lire uniquement les colonnes nécessaires et d'ignorer les groupes de lignes non pertinents. Parquet est le format de stockage sous-jacent pour Delta Lake (/delta/index.md), ce qui en fait le format le plus courant pour les données stockées dans Databricks. Databricks prend en charge Parquet pour la lecture et l'écriture avec Apache Spark, y compris la spécification de schéma, le partitionnement et la compression en écriture.
Prérequis
Databricks ne nécessite pas de configuration supplémentaire pour utiliser les fichiers Parquet. Cependant, pour Stream des fichiers Parquet, vous avez besoin de Auto Loader.
Options
Utilisez les méthodes .option() et .options() de DataFrameReader et DataFrameWriter pour configurer les sources de données Parquet. Pour une liste complète des options prises en charge, consultez les options ParquetDataFrameReader et les options ParquetDataFrameWriter.
Utilisation
Les exemples suivants utilisent l'exemple de dataset Wanderbricks pour démontrer la lecture et l'écriture de fichiers Parquet à l'aide de l'API Spark DataFrame et de SQL.
Lire les fichiers Parquet à l'aide de SQL
Utilisez read_files pour interroger les fichiers Parquet directement depuis le stockage cloud à l'aide de SQL sans créer de table.
SELECT * FROM read_files(
'/Volumes/<catalog>/<schema>/<volume>/reviews_parquet',
format => 'parquet'
)
Lire et écrire des fichiers Parquet
Les exemples suivants écrivent les revues Wanderbricks au format Parquet, les relisent dans un DataFrame et démontrent le mode d'écrasement.
- Python
- Scala
- SQL
# Write wanderbricks reviews to Parquet format
df = spark.read.table("samples.wanderbricks.reviews")
df.write.format("parquet").save("/Volumes/<catalog>/<schema>/<volume>/reviews_parquet")
# Read a Parquet file into a DataFrame
df = spark.read.format("parquet").load("/Volumes/<catalog>/<schema>/<volume>/reviews_parquet")
display(df)
# Write with overwrite mode
df.write.format("parquet").mode("overwrite").save("/Volumes/<catalog>/<schema>/<volume>/reviews_parquet")
// Write wanderbricks reviews to Parquet format
val reviews = spark.read.table("samples.wanderbricks.reviews")
reviews.write.format("parquet").save("/Volumes/<catalog>/<schema>/<volume>/reviews_parquet")
// Read a Parquet file into a DataFrame
val df = spark.read.format("parquet").load("/Volumes/<catalog>/<schema>/<volume>/reviews_parquet")
df.show()
// Write with overwrite mode
df.write.format("parquet").mode("overwrite").save("/Volumes/<catalog>/<schema>/<volume>/reviews_parquet")
-- Write wanderbricks reviews to Parquet format
CREATE TABLE reviews_parquet
USING PARQUET
AS SELECT * FROM samples.wanderbricks.reviews;
SELECT * FROM reviews_parquet;
Spécifiez un schéma
Spécifiez un schéma lors de la lecture des fichiers Parquet afin d’éviter la surcharge de l’inférence de schéma. Par exemple, définissez un schéma avec les champs review_id, rating et comment et lisez reviews_parquet dans un DataFrame.
- Python
- Scala
- SQL
from pyspark.sql.types import StructType, StructField, StringType, IntegerType
schema = StructType([
StructField("review_id", StringType(), True),
StructField("rating", IntegerType(), True),
StructField("comment", StringType(), True)
])
df = spark.read.format("parquet").schema(schema).load("/Volumes/<catalog>/<schema>/<volume>/reviews_parquet")
df.printSchema()
df.show()
import org.apache.spark.sql.types.{StructType, StructField, StringType, IntegerType}
val schema = StructType(Array(
StructField("review_id", StringType, nullable = true),
StructField("rating", IntegerType, nullable = true),
StructField("comment", StringType, nullable = true)
))
val df = spark.read.format("parquet").schema(schema).load("/Volumes/<catalog>/<schema>/<volume>/reviews_parquet")
df.printSchema()
df.show()
-- Create a table with an explicit schema from Parquet files
CREATE TABLE reviews_parquet (
review_id STRING,
rating INT,
comment STRING
)
USING PARQUET
OPTIONS (path "/Volumes/<catalog>/<schema>/<volume>/reviews_parquet");
SELECT * FROM reviews_parquet;
Écrire des fichiers Parquet partitionnés
Écrire des fichiers Parquet partitionnés pour des performances de query optimisées sur de grands datasets. Par exemple, lisez samples.wanderbricks.bookings et écrivez-le dans bookings_parquet_partitioned partitionné par year et month dérivé de la colonne check_in.
- Python
- Scala
- SQL
from pyspark.sql.functions import year, month
df = spark.read.table("samples.wanderbricks.bookings")
df_with_parts = df.withColumn("year", year("check_in")).withColumn("month", month("check_in"))
df_with_parts.write.format("parquet").partitionBy("year", "month").save("/Volumes/<catalog>/<schema>/<volume>/bookings_parquet_partitioned")
import org.apache.spark.sql.functions.{year, month}
val bookings = spark.read.table("samples.wanderbricks.bookings")
val bookingsWithParts = bookings.withColumn("year", year(col("check_in"))).withColumn("month", month(col("check_in")))
bookingsWithParts.write.format("parquet").partitionBy("year", "month").save("/Volumes/<catalog>/<schema>/<volume>/bookings_parquet_partitioned")
-- Write partitioned Parquet files by year and month
CREATE TABLE bookings_parquet_partitioned
USING PARQUET
PARTITIONED BY (year, month)
AS SELECT *, year(check_in) AS year, month(check_in) AS month
FROM samples.wanderbricks.bookings;
Ressources supplémentaires
- Qu'est-ce que Delta Lake dans Databricks ?: Si vous avez besoin de transactions ACID, d'application des schémas ou de time travel, outre les performances en colonne de Parquet, Delta Lake est le format recommandé pour les données stockées dans Databricks.