Ler e escrever ficheiros Parquet

Apache Parquet é um formato de ficheiro colunar otimizado para cargas de trabalho analíticas. Permite que os motores de consulta leiam apenas as colunas necessárias e saltem grupos de linhas irrelevantes. Parquet é o formato de armazenamento subjacente para Delta Lake(/delta/index.md), tornando-se o formato mais comum para dados armazenados no Azure Databricks. O Azure Databricks suporta o Parquet tanto para leitura como para escrita com o Apache Spark, incluindo especificação de esquema, particionamento e compressão de escrita.

Pré-requisitos

O Azure Databricks não requer configuração adicional para utilizar ficheiros Parquet. No entanto, para transmitir ficheiros Parquet, precisas do Auto Loader.

Opções

Utilize os métodos .option() e .options() de DataFrameReader e DataFrameWriter para configurar origens de dados Parquet. Para uma lista completa de opções suportadas, veja DataFrameReader Opções de Parquet e DataFrameWriter Opções de Parquet.

Usage

Os exemplos seguintes utilizam o conjunto de dados de exemplo Wanderbricks para demonstrar a leitura e escrita de ficheiros Parquet usando a API Spark DataFrame e SQL.

Leia ficheiros Parquet usando SQL

Uso read_files para consultar ficheiros Parquet diretamente a partir de armazenamento na cloud usando SQL sem criar uma tabela.

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

Ler e escrever ficheiros Parquet

Os exemplos seguintes escrevem as análises do Wanderbricks para o formato Parquet, lêem-nas novamente num DataFrame e demonstram o modo de sobrescrição.

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;

Especifique um esquema

Especifique um esquema ao ler ficheiros Parquet para evitar a sobrecarga da inferência de esquemas. Por exemplo, defina um esquema com os campos review_id, rating e comment e leia reviews_parquet para um 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;

Escrever ficheiros Parquet particionados

Escrever ficheiros Parquet particionados para otimizar o desempenho das consultas em grandes conjuntos de dados. Por exemplo, leia samples.wanderbricks.bookings e escreva em bookings_parquet_partitioned particionado por year e month derivado da check_in coluna.

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;

Recursos adicionais

  • O que é Delta Lake no Azure Databricks?: Se precisar de transações ACID, aplicação de esquemas ou viagens no tempo juntamente com o desempenho colunar do Parquet, o Delta Lake é o formato recomendado para dados armazenados no Azure Databricks.