读取和写入 Parquet 文件

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

Prerequisites

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

选项

使用 .option().options()DataFrameReaderDataFrameWriter 方法来配置 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_idratingcomment 字段的架构,并将 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列派生出的monthcheck_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;

其他资源

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