Odczytywanie i zapisywanie plików Parquet

Apache Parquet to format pliku kolumnowego zoptymalizowany pod kątem obciążeń analitycznych. Umożliwia on aparatom zapytań odczytywanie tylko potrzebnych kolumn i pomijanie nieistotnych grup wierszy. Parquet jest bazowym formatem przechowywania dla Delta Lake(/delta/index.md), dzięki czemu jest najczęściej spotykanym formatem danych przechowywanych w usłudze Azure Databricks. Azure Databricks obsługuje język Parquet na potrzeby odczytu i zapisu na platformie Apache Spark, w tym specyfikacji schematu, partycjonowania i kompresji zapisu.

Prerequisites

Azure Databricks nie wymaga dodatkowej konfiguracji do korzystania z plików Parquet. Jednak aby przesyłać strumieniowo pliki Parquet, potrzebujesz Auto Loader.

Opcje

.option() Użyj metod .options() i DataFrameReader , DataFrameWriter aby skonfigurować źródła danych Parquet. Aby uzyskać pełną listę obsługiwanych opcji, zobacz DataFrameReader Opcje Parquet i DataFrameWriter Opcje Parquet.

Usage

W poniższych przykładach użyto przykładowego zestawu danych usługi Wanderbricks , aby zademonstrować odczytywanie i zapisywanie plików Parquet przy użyciu interfejsu API ramki danych platformy Spark i języka SQL.

Odczytywanie plików Parquet przy użyciu języka SQL

Użyj read_files polecenia , aby wysyłać zapytania do plików Parquet bezpośrednio z magazynu w chmurze przy użyciu języka SQL bez tworzenia tabeli.

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

Odczytywanie i zapisywanie plików Parquet

Poniższe przykłady zapisują recenzje Wanderbricks do formatu Parquet, odczytują je z powrotem do obiektu DataFrame i demonstrują tryb nadpisywania.

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;

Określanie schematu

Określ schemat podczas odczytu plików Parquet, aby uniknąć narzutu wynikającego z automatycznego wykrywania schematu. Na przykład zdefiniuj schemat z polami review_id, rating i comment, a następnie wczytaj reviews_parquet do ramki danych.

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;

Zapisywanie partycjonowanych plików Parquet

Zapisywanie partycjonowanych plików Parquet pod kątem zoptymalizowanej wydajności zapytań w dużych zestawach danych. Na przykład odczytaj samples.wanderbricks.bookings i zapisz je w bookings_parquet_partitioned, z podziałem na partycje według year i month, wyprowadzonych z kolumny 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.{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;

Dodatkowe zasoby

  • Czym jest Delta Lake w Azure Databricks?: Jeśli potrzebujesz transakcji ACID, wymuszania schematu lub funkcji podróży w czasie, przy jednoczesnym zachowaniu kolumnowej wydajności formatu Parquet, Delta Lake jest zalecanym formatem przechowywania danych w usłudze Azure Databricks.