Avro dosyalarını okuma ve yazma

Apache Avro , zengin veri yapıları ve kompakt, hızlı ikili kodlama sağlayan satır tabanlı bir veri serileştirme biçimidir. Azure Databricks kullanıcıları bununla en yaygın olarak, Avro’nun baskın serileştirme formatı olduğu Apache Kafka ve Google Pub/Sub gibi olay akışı sistemlerinden veri alımı sırasında karşılaşırlar. Azure Databricks, Avro ve Spark SQL türleri arasında otomatik şema dönüştürme, bölümleme, sıkıştırma ve özel kayıt adları dahil olmak üzere Apache Spark ile hem okuma hem de yazma için Avro'yı destekler.

Avro ile kodlanmış kayıtları dosyalardan değil de Apache Kafka'dan veya başka bir mesaj veri yolundan okuyorsanız, akışta seri durumdan çıkarma için kullanılan from_avro ve to_avro işlevlerini ele alan Akış Avro verilerini okuma ve yazma bölümüne bakın.

Prerequisites

Azure Databricks, Avro dosyalarını kullanmak için ek yapılandırma gerektirmez. Ancak Avro dosyalarının akışını yapmak için Otomatik Yükleyici gerekir.

Options

Avro veri kaynaklarını yapılandırmak için DataFrameReader ve DataFrameWriter öğelerinin .option() ve .options() yöntemlerini kullanın. Desteklenen seçeneklerin tam listesi için bkz DataFrameReader . Avro seçenekleri ve DataFrameWriter Avro seçenekleri.

Usage

Aşağıdaki örneklerde Spark DataFrame API'sini ve SQL'i kullanarak Avro dosyalarını okuma ve yazma işlemini göstermek için Wanderbricks veri kümesi kullanılır.

SQL kullanarak Avro dosyalarını okuma

Tablo kaydetmeden Avro dosyalarını sorgulamak için kullanın read_files. Dış konumdaki Unity Kataloğu izinleri otomatik olarak uygulanır.

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

Avro dosyalarını okuma ve yazma

Aşağı akış sistemi için Avro dosyalarını okumanız veya yazmanız, yüklemeden önce dönüşümleri uygulamanız veya yazma zamanında bölümleme ve şema gibi denetim seçeneklerini uygulamanız gerektiğinde Apache Spark DataFrame API'sini kullanın.

Aşağıdaki örneklerde Wanderbricks örnek veri kümesi kullanılmıştır.

Piton

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 programlama dili

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

// 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;

Ek kaynaklar

  • Parquet dosyalarını okuma ve yazma: İş yükünüz akış veya yoğun yazma yerine öncelikli olarak analitik ve okuma ağırlıklıysa Parquet'in sütun düzeni Avro'nun satır tabanlı depolamasından daha verimli sorgu performansı sunar.