Apache ORC 是针对大规模分析工作负荷优化的列式文件格式。 它使用内置索引和统计信息在读取期间跳过不相关的数据。 Azure Databricks支持 ORC 通过 Apache Spark 进行读取和写入,包括架构规范、分区和写入压缩。
先决条件
Azure Databricks不需要其他配置即可使用 ORC 文件。 但是,若要流式传输 ORC 文件,需要 自动加载程序。
选项
使用 .option() 和 .options() 的 DataFrameReader 和 DataFrameWriter 方法来配置 ORC 数据源。 有关支持选项的完整列表,请参阅 DataFrameReader ORC 选项 和 DataFrameWriter ORC 选项。
Usage
以下示例使用 Wanderbricks 示例数据集演示如何使用 Spark 数据帧 API 和 SQL 读取和写入 ORC 文件。
读取和写入 ORC 文件
Python
# Write wanderbricks reviews to ORC format
df = spark.read.table("samples.wanderbricks.reviews")
df.write.format("orc").save("/Volumes/<catalog>/<schema>/<volume>/reviews_orc")
# Read an ORC file into a DataFrame
df = spark.read.format("orc").load("/Volumes/<catalog>/<schema>/<volume>/reviews_orc")
display(df)
# Write with overwrite mode
df.write.format("orc").mode("overwrite").save("/Volumes/<catalog>/<schema>/<volume>/reviews_orc")
Scala
// Write wanderbricks reviews to ORC format
val reviews = spark.read.table("samples.wanderbricks.reviews")
reviews.write.format("orc").save("/Volumes/<catalog>/<schema>/<volume>/reviews_orc")
// Read an ORC file into a DataFrame
val df = spark.read.format("orc").load("/Volumes/<catalog>/<schema>/<volume>/reviews_orc")
df.show()
// Write with overwrite mode
df.write.format("orc").mode("overwrite").save("/Volumes/<catalog>/<schema>/<volume>/reviews_orc")
SQL
-- Write wanderbricks reviews to ORC format
CREATE TABLE reviews_orc
USING ORC
AS SELECT * FROM samples.wanderbricks.reviews;
SELECT * FROM reviews_orc;
使用 SQL 读取 ORC 文件
使用 read_files 用 SQL 直接从云存储查询 ORC 文件,而无需创建表。
SELECT * FROM read_files(
'/Volumes/<catalog>/<schema>/<volume>/reviews_orc',
format => 'orc'
)
指定架构
读取 ORC 文件时指定架构,以避免架构推理开销。 例如,定义一个包含 review_id、rating 和 comment 字段的架构,并将 reviews_orc 读入 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("orc").schema(schema).load("/Volumes/<catalog>/<schema>/<volume>/reviews_orc")
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("orc").schema(schema).load("/Volumes/<catalog>/<schema>/<volume>/reviews_orc")
df.printSchema()
df.show()
SQL
-- Create a table with an explicit schema from ORC files
CREATE TABLE reviews_orc (
review_id STRING,
rating INT,
comment STRING
)
USING ORC
OPTIONS (path "/Volumes/<catalog>/<schema>/<volume>/reviews_orc");
SELECT * FROM reviews_orc;
写入分区的 ORC 文件
编写分区的 ORC 文件,以优化大型数据集的查询性能。 例如,读取samples.wanderbricks.bookings并将其写入bookings_orc_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("orc").partitionBy("year", "month").save("/Volumes/<catalog>/<schema>/<volume>/bookings_orc_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("orc").partitionBy("year", "month").save("/Volumes/<catalog>/<schema>/<volume>/bookings_orc_partitioned")
SQL
-- Write partitioned ORC files by year and month
CREATE TABLE bookings_orc_partitioned
USING ORC
PARTITIONED BY (year, month)
AS SELECT *, year(check_in) AS year, month(check_in) AS month
FROM samples.wanderbricks.bookings;
其他资源
- 什么是 Azure Databricks 中的 Delta Lake?:如果你正从使用 ORC 的 Hive 或 Hadoop 环境迁移,Delta Lake 是 Databricks 推荐的原生格式。 它在基于 Parquet 的存储之上增加了 ACID 事务、架构约束、时间旅行和优化的读取性能。
- 读取和写入 Parquet 文件:如果工作负荷需要 Databricks 以外的最广泛的生态系统兼容性,Parquet 是查询引擎和云存储工具中支持的最广泛的列式格式。