Notitie
Voor toegang tot deze pagina is autorisatie vereist. U kunt proberen u aan te melden of de directory te wijzigen.
Voor toegang tot deze pagina is autorisatie vereist. U kunt proberen de mappen te wijzigen.
Apache Parquet is een kolombestandsindeling die is geoptimaliseerd voor analytische workloads. Hiermee kunnen query-engines alleen de benodigde kolommen lezen en irrelevante rijgroepen overslaan. Parquet is de onderliggende opslagindeling voor Delta Lake(/delta/index.md), waardoor het de meest voorkomende indeling is voor gegevens die zijn opgeslagen in Azure Databricks. Azure Databricks ondersteunt Parquet voor zowel lezen als schrijven met Apache Spark, waaronder schemaspecificatie, partitionering en schrijfcompressie.
Vereiste voorwaarden
Azure Databricks vereist geen aanvullende configuratie voor het gebruik van Parquet-bestanden. Om Parquet-bestanden te streamen, hebt u echter Auto Loader nodig.
Opties
Gebruik de .option() en .options() methoden van DataFrameReader en DataFrameWriter om Parquet-gegevensbronnen te configureren. Zie DataFrameReader Parquet-opties en DataFrameWriter Parquet-opties voor een volledige lijst met ondersteunde opties.
Usage
In de volgende voorbeelden wordt de Wanderbricks-voorbeeldgegevensset gebruikt om het lezen en schrijven van Parquet-bestanden te demonstreren met behulp van de Spark DataFrame-API en SQL.
Parquet-bestanden lezen met SQL
Gebruik read_files dit om query's uit te voeren op Parquet-bestanden rechtstreeks vanuit cloudopslag met behulp van SQL zonder een tabel te maken.
SELECT * FROM read_files(
'/Volumes/<catalog>/<schema>/<volume>/reviews_parquet',
format => 'parquet'
)
Parquet-bestanden lezen en schrijven
In de volgende voorbeelden worden de Wanderbricks-reviews in Parquet-indeling geschreven, vervolgens weer ingelezen in een DataFrame en wordt de overschrijfmodus gedemonstreerd.
Python
# 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")
Scala
// 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")
SQL
-- Write wanderbricks reviews to Parquet format
CREATE TABLE reviews_parquet
USING PARQUET
AS SELECT * FROM samples.wanderbricks.reviews;
SELECT * FROM reviews_parquet;
Een schema opgeven
Geef een schema op bij het lezen van Parquet-bestanden om de overhead van schemadeductie te voorkomen. Definieer bijvoorbeeld een schema met de velden review_id, rating en comment en lees reviews_parquet in in een DataFrame.
Python
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()
Scala
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()
SQL
-- 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;
Gepartitioneerde Parquet-bestanden schrijven
Schrijf gepartitioneerde Parquet-bestanden voor geoptimaliseerde queryprestaties voor grote gegevenssets. Lees samples.wanderbricks.bookings en schrijf deze bijvoorbeeld naar bookings_parquet_partitioned gepartitioneerd door year en month afgeleid van de check_in kolom.
Python
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")
Scala
import org.apache.spark.sql.functions.{col, month, year}
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")
SQL
-- 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;
Aanvullende bronnen
- Wat is Delta Lake in Azure Databricks?: Als u ACID-transacties, schemaafdwinging of tijdreizen naast de kolomprestaties van Parquet nodig hebt, is Delta Lake de aanbevolen indeling voor gegevens die zijn opgeslagen in Azure Databricks.