Lire et écrire des fichiers Avro

Apache Avro est un format de sérialisation de données basé sur des lignes qui fournit des structures de données enrichies et un encodage binaire compact et rapide. Les utilisateurs d’Azure Databricks y sont le plus souvent confrontés lors de l’ingestion de données depuis des systèmes de streaming d’événements tels qu’Apache Kafka et Google Pub/Sub, où Avro est le format de sérialisation dominant. Azure Databricks prend en charge Avro pour la lecture et l’écriture avec Apache Spark, notamment la conversion automatique de schémas entre les types SQL Avro et Spark, le partitionnement, la compression et les noms d’enregistrements personnalisés.

Si vous lisez des enregistrements encodés Avro à partir d’Apache Kafka ou d’un autre bus de messages plutôt qu’à partir de fichiers, consultez Lire et écrire des données Avro en continu, qui couvrent les from_avro fonctions utilisées to_avro pour la désérialisation de streaming.

Prerequisites

Azure Databricks ne nécessite pas de configuration supplémentaire pour utiliser les fichiers Avro. Toutefois, pour diffuser en continu des fichiers Avro, vous avez besoin d’un chargeur automatique.

Options

Utilisez les méthodes .option() et .options() de DataFrameReader et DataFrameWriter pour configurer les sources de données Avro. Pour obtenir la liste complète des options prises en charge, consultez DataFrameReader les options Avro et DataFrameWriter les options Avro.

Utilisation

Les exemples suivants utilisent le jeu de données Wanderbricks pour illustrer la lecture et l’écriture de fichiers Avro à l’aide de l’API DataFrame Spark et de SQL.

Lire des fichiers Avro à l’aide de SQL

Pour interroger des fichiers Avro sans inscrire une table, utilisez read_files. Les autorisations du catalogue Unity sur l’emplacement externe s’appliquent automatiquement.

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

Lire et écrire des fichiers Avro

Utilisez l’API DataFrame Apache Spark lorsque vous devez lire ou écrire des fichiers Avro pour un système en aval, appliquer des transformations avant le chargement ou des options de contrôle telles que le partitionnement et le schéma au moment de l’écriture.

Les exemples suivants utilisent l’exemple de jeu de données 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")

Langage de programmation 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;

Ressources additionnelles

  • Lire et écrire des fichiers Parquet : si votre charge de travail est principalement analytique et davantage axée sur la lecture que sur le traitement en continu ou une forte écriture, le format colonnaire de Parquet offre de meilleures performances de requête que le stockage en lignes d’Avro.