Lire et écrire des fichiers Avro
Le format Apache Avro est un format de sérialisation de données basé sur les lignes qui fournit des structures de données riches et un encodage binaire compact et rapide. Les utilisateurs de Databricks le rencontrent le plus souvent lorsqu'ils ingèrent des données de systèmes de streaming d'événements tels que Apache Kafka et Google Pub/Sub, où Avro est le format de sérialisation dominant. Databricks prend en charge Avro pour la lecture et l'écriture avec Apache Spark, y compris la conversion automatique de schéma entre les types Avro et Spark SQL, le partitionnement, la compression et les noms d'enregistrement 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 streaming, qui couvre les fonctions from_avro et to_avro utilisées pour la désérialisation en streaming.
Prérequis
Databricks ne nécessite pas de configuration supplémentaire pour utiliser les fichiers Avro. Cependant, pour le Stream de fichiers Avro, vous avez besoin de l'Auto Loader.
Options
Utilisez les méthodes .option() et .options() de DataFrameReader et DataFrameWriter pour configurer les sources de données Avro. Pour une liste complète des options prises en charge, consultez DataFrameReader options Avro et DataFrameWriter options Avro.
Utilisation
Les exemples suivants utilisent le dataset Wanderbricks pour illustrer la lecture et l'écriture de fichiers Avro à l'aide de l'API Spark DataFrame et de SQL.
Lisez les fichiers Avro à l'aide de SQL
Pour interroger les fichiers Avro sans enregistrer de table, utilisez read_files. Les autorisations Unity Catalog 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 Apache Spark DataFrame lorsque vous avez besoin de lire ou d'écrire des fichiers Avro pour un système en aval, d'appliquer des transformations avant le chargement, ou de contrôler des options 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
- Scala
- SQL
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")
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")
-- 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 supplémentaires
- Lisez et écrivez des fichiers Parquet: Si votre charge de travail est principalement analytique et à forte lecture plutôt que streaming ou à forte écriture, le Layout en colonnes de Parquet offre des performances de query plus efficaces que le stockage basé sur les lignes d'Avro.