Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Apache Avro — это формат сериализации данных на основе строк, предоставляющий широкие структуры данных и компактную, быструю двоичную кодировку. Azure Databricks пользователи чаще всего сталкиваются с ним при приеме данных из систем потоковой передачи событий, таких как Apache Kafka и Google Pub/Sub, где Avro является доминирующим форматом сериализации. Azure Databricks поддерживает Avro для чтения и записи с помощью Apache Spark, включая автоматическое преобразование схемы между типами Avro и Spark SQL, секционированием, сжатием и пользовательскими именами записей.
Если вы читаете записи в кодировке Avro из Apache Kafka или другой шины сообщений, а не из файлов, см. статью "Чтение и запись потоковых данных Avro", которая охватывает from_avro и to_avro функции, используемые для десериализации потоковой передачи.
Необходимые условия
Azure Databricks не требует дополнительной конфигурации для использования файлов Avro. Однако для потоковой передачи файлов Avro требуется автозагрузчик.
Options
Используйте методы .option() и .options() объектов DataFrameReader и DataFrameWriter для настройки источников данных Avro. Полный список поддерживаемых параметров см. в разделе DataFrameReader "Параметры Avro " и DataFrameWriter "Avro".
Usage
В следующих примерах используется набор данных Wanderbricks для демонстрации чтения и записи файлов Avro с помощью API Spark DataFrame и SQL.
Чтение файлов Avro с помощью SQL
Чтобы запросить файлы Avro без регистрации таблицы, используйте read_files. Разрешения на внешнее расположение в Unity Catalog применяются автоматически.
SELECT * FROM read_files(
'/Volumes/<catalog>/<schema>/<volume>/reviews_avro',
format => 'avro'
)
Чтение и запись файлов Avro
Используйте API Кадра данных Apache Spark, если необходимо считывать или записывать файлы Avro для нижестоящей системы, применять преобразования перед загрузкой или параметры управления, такие как секционирование и схема во время записи.
В следующих примерах используется пример набора данных Wanderbricks .
Питон
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.{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;
Дополнительные ресурсы
- Чтение и запись файлов Parquet: если ваша рабочая нагрузка носит преимущественно аналитический характер и ориентирована в основном на чтение, а не на потоковую обработку данных или интенсивную запись, столбцовый формат Parquet обеспечивает более высокую производительность запросов по сравнению со строко-ориентированным хранением Avro.