Kommentar
Åtkomst till den här sidan kräver auktorisering. Du kan prova att logga in eller ändra kataloger.
Åtkomst till den här sidan kräver auktorisering. Du kan prova att ändra kataloger.
Apache Avro är ett radbaserat dataserialiseringsformat som ger omfattande datastrukturer och en kompakt, snabb binär kodning. Azure Databricks användare stöter oftast på det när de matar in data från händelseströmningssystem som Apache Kafka och Google Pub/Sub, där Avro är det dominerande serialiseringsformatet. Azure Databricks stöder Avro för både läsning och skrivning med Apache Spark, inklusive automatisk schemakonvertering mellan Avro- och Spark SQL-typer, partitionering, komprimering och anpassade postnamn.
Om du läser Avro-kodade dataposter från Apache Kafka eller en annan meddelandebuss snarare än från filer, se Läsa och skriva strömmande Avro-data, som behandlar funktionerna from_avro och to_avro som används för strömmande deserialisering.
Förutsättningar
Azure Databricks kräver inte ytterligare konfiguration för att använda Avro-filer. Om du vill strömma Avro-filer behöver du dock automatisk inläsning.
Options
Använd metoderna .option() och .options() i DataFrameReader och DataFrameWriter för att konfigurera Avro-datakällor. En fullständig lista över alternativ som stöds finns i DataFrameReader Avro-alternativ och DataFrameWriter Avro-alternativ.
Usage
I följande exempel används Wanderbricks-datamängden för att demonstrera läsning och skrivning av Avro-filer med hjälp av Spark DataFrame API och SQL.
Läsa Avro-filer med SQL
Om du vill köra frågor mot Avro-filer utan att registrera en tabell använder du read_files. Behörigheter för Unity Catalog för den externa platsen gäller automatiskt.
SELECT * FROM read_files(
'/Volumes/<catalog>/<schema>/<volume>/reviews_avro',
format => 'avro'
)
Läsa och skriva Avro-filer
Använd Apache Spark DataFrame-API:et när du behöver läsa eller skriva Avro-filer för ett nedströmssystem, tillämpa transformeringar före inläsning eller kontrollalternativ som partitionering och schema vid skrivtillfället.
I följande exempel används Wanderbricks-exempeldatauppsättningen .
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.{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;
Ytterligare resurser
- Läsa och skriva Parquet-filer: Om din arbetsbelastning främst är analytisk och läsintensiv i stället för strömning eller skrivintensiv, erbjuder Parquets kolumnlayout effektivare frågeprestanda än Avros radbaserade lagring.