Leggere e scrivere file Parquet

Apache Parquet è un formato di file a colonne ottimizzato per i carichi di lavoro analitici. Consente ai motori di query di leggere solo le colonne necessarie e ignorare i gruppi di righe irrilevanti. Parquet è il formato di archiviazione sottostante per Delta Lake(/delta/index.md), rendendolo il formato più comune per i dati archiviati in Azure Databricks. Azure Databricks supporta Parquet sia per la lettura che per la scrittura con Apache Spark, incluse le specifiche dello schema, il partizionamento e la compressione di scrittura.

Prerequisiti

Azure Databricks non richiede una configurazione aggiuntiva per l'uso di file Parquet. Tuttavia, per trasmettere i file Parquet, è necessario il caricatore automatico.

Opzioni

Utilizzare i metodi .option() e .options() di DataFrameReader e DataFrameWriter per configurare le origini dati Parquet. Per un elenco completo delle opzioni supportate, vedere DataFrameReader Opzioni Parquet e DataFrameWriter Opzioni Parquet.

Usage

Gli esempi seguenti usano il set di dati di esempio Wanderbricks per illustrare la lettura e la scrittura di file Parquet usando l'API DataFrame Spark e SQL.

Leggere i file Parquet con SQL

Usare read_files per eseguire query sui file Parquet direttamente dall'archiviazione cloud usando SQL senza creare una tabella.

SELECT * FROM read_files(
  '/Volumes/<catalog>/<schema>/<volume>/reviews_parquet',
  format => 'parquet'
)

Leggere e scrivere i file Parquet

Gli esempi seguenti scrivono le recensioni di Wanderbricks in formato Parquet, le rileggono in un DataFrame e dimostrano la modalità di sovrascrittura.

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;

Specificare uno schema

Specificare uno schema durante la lettura dei file Parquet per evitare il sovraccarico dell'inferenza dello schema. Ad esempio, definire uno schema con i campi review_id, rating e comment e leggere reviews_parquet in un 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;

Scrivere file Parquet partizionati

Scrivere file Parquet partizionati per ottimizzare le prestazioni delle query in set di dati di grandi dimensioni. Ad esempio, leggere samples.wanderbricks.bookings e scrivere in bookings_parquet_partitioned, partizionato per year e month derivati dalla colonna check_in.

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.{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")

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;

Risorse aggiuntive

  • Che cos'è Delta Lake in Azure Databricks?: se sono necessarie transazioni ACID, imposizione dello schema o tempo di spostamento insieme alle prestazioni a colonne di Parquet, Delta Lake è il formato consigliato per i dati archiviati in Azure Databricks.