读取和写入 Parquet 文件

Apache Parquet 是针对分析工作负荷优化的列式文件格式。 它允许查询引擎仅读取所需的列,并跳过不相关的行组。 Parquet 是 Delta Lake(/delta/index.md)的基础存储格式,使其成为存储在 Azure Databricks 中的数据的最常见格式。 Azure Databricks 支持使用 Apache Spark 读取和写入 Parquet 文件,包括架构指定、分区和写入压缩。

Prerequisites

Azure Databricks不需要其他配置才能使用 Parquet 文件。 但是,若要流式传输 Parquet 文件,需要 自动加载程序。

选项

使用 .option() 和 .options() 的 DataFrameReader 和 DataFrameWriter 方法来配置 Parquet 数据源。 有关受支持选项的完整列表,请参阅 DataFrameReader Parquet 选项 和 DataFrameWriter Parquet 选项。

Usage

以下示例使用 Wanderbricks 示例数据集演示如何使用 Spark 数据帧 API 和 SQL 读取和写入 Parquet 文件。

使用 SQL 读取 Parquet 文件

使用 read_files 直接通过 SQL 从云存储中查询 Parquet 文件,而无需创建表。

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

读取和写入 Parquet 文件

以下示例将 Wanderbricks 评论写入 Parquet 格式,将其读回到 DataFrame 中,并演示覆盖模式。

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;

指定架构

读取 Parquet 文件时指定模式,以避免模式推断带来的开销。 例如,定义一个包含 review_id、rating 和 comment 字段的架构,并将 reviews_parquet 读入 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;

写入分区的 Parquet 文件

编写分区的 Parquet 文件,以优化大型数据集的查询性能。 例如,读取samples.wanderbricks.bookings并将其写入bookings_parquet_partitioned,按由year列派生出的month和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;

其他资源

  • Azure Databricks 中的 Delta Lake 是什么?:如果你既需要 ACID 事务、架构强制执行或时光回溯功能,又希望具备 Parquet 的列式性能,那么 Delta Lake 是 Azure Databricks 中数据存储的推荐格式。