讀寫 Avro 檔案

Apache Avro 是一種基於列的資料序列化格式,提供豐富的資料結構與緊湊且快速的二進位編碼。 Azure Databricks 使用者最常遇到這種情況是在從事件串流系統(如 Apache Kafka 和 Google Pub/Sub)匯入資料時,而 Avro 是主流的序列化格式。 Azure Databricks 支援使用 Apache Spark 讀取及寫入 Avro,包括 Avro 與 Spark SQL 類型之間的自動架構轉換、分割區、壓縮及自訂記錄名稱。

如果你是從 Apache Kafka 或其他訊息匯流排讀取 Avro 編碼的記錄,而不是從檔案讀取,請參閱 讀寫串流 Avro 資料,其中涵蓋了用於串流反序列化的 from_avro 和 to_avro 函式。

先決條件

Azure Databricks 使用 Avro 檔案不需要額外設定。 不過,要串流 Avro 檔案,你需要 Auto Loader。

選項

使用 DataFrameReader 和 DataFrameWriter 的 .option() 與 .options() 方法來設定 Avro 資料來源。 欲了解完整的支援選項清單,請參閱 DataFrameReader Avro 選項 及 DataFrameWriter Avro 選項。

Usage

以下範例使用 Wanderbricks 資料集 示範使用 Spark DataFrame API 與 SQL 讀寫 Avro 檔案。

使用 SQL 讀取 Avro 檔案

若要查詢 Avro 檔案而不註冊資料表,請使用 read_files。 Unity Catalog 對外部位置的權限會自動套用。

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

讀寫 Avro 檔案

當你需要讀寫下游系統的 Avro 檔案、在載入前套用轉換,或在寫入時控制分割與結構等選項時,請使用 Apache Spark DataFrame API。

以下範例使用 Wanderbricks 範例資料集。

Python

from pyspark.sql.functions import year, month

# Write wanderbricks reviews to Avro format
df = spark.read.table("samples.wanderbricks.reviews")
df.write.format("avro").save("/Volumes/<catalog>/<schema>/<volume>/reviews_avro")

# Read an Avro file into a DataFrame
df = spark.read.format("avro").load("/Volumes/<catalog>/<schema>/<volume>/reviews_avro")
display(df)

# Write with overwrite mode
df.write.format("avro").mode("overwrite").save("/Volumes/<catalog>/<schema>/<volume>/reviews_avro")

# Read using a custom Avro schema to select specific fields
avro_schema = """
{
  "type": "record",
  "name": "Review",
  "fields": [
    {"name": "review_id", "type": "string"},
    {"name": "rating", "type": "int"},
    {"name": "comment", "type": ["null", "string"]}
  ]
}
"""
df = spark.read.format("avro").option("avroSchema", avro_schema).load("/Volumes/<catalog>/<schema>/<volume>/reviews_avro")

# Write partitioned Avro files by year and 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("avro").partitionBy("year", "month").save("/Volumes/<catalog>/<schema>/<volume>/bookings_avro_partitioned")

# Write with a custom record name and namespace for Schema Registry compatibility
df = spark.read.table("samples.wanderbricks.reviews")
df.write.format("avro").options(
  recordName="Review",
  recordNamespace="com.wanderbricks"
).save("/Volumes/<catalog>/<schema>/<volume>/reviews_avro")

程式語言 Scala

import org.apache.spark.sql.functions.{col, month, year}

// Write wanderbricks reviews to Avro format
val reviews = spark.read.table("samples.wanderbricks.reviews")
reviews.write.format("avro").save("/Volumes/<catalog>/<schema>/<volume>/reviews_avro")

// Read an Avro file into a DataFrame
val df = spark.read.format("avro").load("/Volumes/<catalog>/<schema>/<volume>/reviews_avro")
df.show()

// Write with overwrite mode
df.write.format("avro").mode("overwrite").save("/Volumes/<catalog>/<schema>/<volume>/reviews_avro")

// Read using a custom Avro schema to select specific fields
val avroSchema = """
{
  "type": "record",
  "name": "Review",
  "fields": [
    {"name": "review_id", "type": "string"},
    {"name": "rating", "type": "int"},
    {"name": "comment", "type": ["null", "string"]}
  ]
}
"""
val filtered = spark.read.format("avro").option("avroSchema", avroSchema).load("/Volumes/<catalog>/<schema>/<volume>/reviews_avro")

// Write partitioned Avro files by year and 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("avro").partitionBy("year", "month").save("/Volumes/<catalog>/<schema>/<volume>/bookings_avro_partitioned")

// Write with a custom record name and namespace for Schema Registry compatibility
reviews.write.format("avro").options(Map(
  "recordName" -> "Review",
  "recordNamespace" -> "com.wanderbricks"
)).save("/Volumes/<catalog>/<schema>/<volume>/reviews_avro")

SQL

-- Write wanderbricks reviews to Avro format
CREATE TABLE reviews_avro
USING AVRO
AS SELECT * FROM samples.wanderbricks.reviews;

-- Write partitioned Avro files by year and month
CREATE TABLE bookings_avro_partitioned
USING AVRO
PARTITIONED BY (year, month)
AS SELECT *, year(check_in) AS year, month(check_in) AS month
FROM samples.wanderbricks.bookings;

SELECT * FROM bookings_avro_partitioned;

其他資源

  • 讀寫 Parquet 檔案:如果你的工作負載主要是分析和大量讀取,而非串流或寫入,Parquet 的欄式配置比 Avro 的列式儲存提供更有效率的查詢效能。