Lire et écrire des données Avro en streaming
Apache Avro est un système de sérialisation de données couramment utilisé dans le monde du streaming. Une solution typique consiste à placer les données au format Avro dans Apache Kafka, les métadonnées dans Confluent Schema Registry, puis à exécuter des requêtes avec un framework de streaming qui se connecte à Kafka et à Schema Registry.
Databricks prend en charge les fonctions from_avro et to_avro pour créer des pipelines de streaming avec des données Avro dans Kafka et des métadonnées dans Schema Registry. La fonction to_avro encode une colonne au format binaire Avro et from_avro décode les données binaires Avro en une colonne. Les deux fonctions transforment une colonne en une autre colonne, et le type de données SQL d'entrée/sortie peut être un type complexe ou un type primitif.
Les fonctions from_avro et to_avro :
- Sont disponibles en Python, Scala et Java.
- Peut être transmis aux fonctions SQL dans les requêtes batch et streaming.
Voir aussi source de données de fichier Avro. Pour une liste complète des options from_avro et to_avro, consultez Avro.
Exemple de schéma spécifié manuellement
De manière similaire à from_json et to_json, vous pouvez utiliser from_avro et to_avro avec n'importe quelle colonne binaire. Vous pouvez spécifier le schéma Avro manuellement, comme dans l'exemple suivant :
import org.apache.spark.sql.avro.functions._
import org.apache.avro.SchemaBuilder
// When reading the key and value of a Kafka topic, decode the
// binary (Avro) data into structured data.
// The schema of the resulting DataFrame is: <key: string, value: int>
val df = spark
.readStream
.format("kafka")
.option("kafka.bootstrap.servers", servers)
.option("subscribe", "t")
.load()
.select(
from_avro($"key", SchemaBuilder.builder().stringType()).as("key"),
from_avro($"value", SchemaBuilder.builder().intType()).as("value"))
// Convert structured data to binary from string (key column) and
// int (value column) and save to a Kafka topic.
dataDF
.select(
to_avro($"key").as("key"),
to_avro($"value").as("value"))
.writeStream
.format("kafka")
.option("kafka.bootstrap.servers", servers)
.option("topic", "t")
.start()
Exemple jsonFormatSchema
Vous pouvez également spécifier un schéma sous forme de chaîne JSON. Par exemple, si /tmp/user.avsc est :
{
"namespace": "example.avro",
"type": "record",
"name": "User",
"fields": [
{ "name": "name", "type": "string" },
{ "name": "favorite_color", "type": ["string", "null"] }
]
}
Vous pouvez créer une chaîne JSON :
from pyspark.sql.avro.functions import from_avro, to_avro
jsonFormatSchema = open("/tmp/user.avsc", "r").read()
Utilisez ensuite le schéma dans from_avro:
# 1. Decode the Avro data into a struct.
# 2. Filter by column "favorite_color".
# 3. Encode the column "name" in Avro format.
output = df\
.select(from_avro("value", jsonFormatSchema).alias("user"))\
.where('user.favorite_color == "red"')\
.select(to_avro("user.name").alias("value"))
Exemple avec le registre des schémas
Si votre cluster dispose d'un service de registre de schémas, from_avro peut fonctionner avec celui-ci afin que vous n'ayez pas besoin de spécifier le schéma Avro manuellement.
L'exemple suivant montre la lecture d'un sujet Kafka « t », en supposant que la clé et la valeur sont déjà enregistrées dans le Registre de schémas en tant que sujets « t-key » et « t-value » de types STRING et INT:
import org.apache.spark.sql.avro.functions._
val schemaRegistryAddr = "https://myhost:8081"
val df = spark
.readStream
.format("kafka")
.option("kafka.bootstrap.servers", servers)
.option("subscribe", "t")
.load()
.select(
from_avro(data = $"key", subject = "t-key", schemaRegistryAddress = schemaRegistryAddr).as("key"),
from_avro(data = $"value", subject = "t-value", schemaRegistryAddress = schemaRegistryAddr).as("value"))
Pour to_avro, le schéma de sortie Avro default peut ne pas correspondre au schéma du sujet cible dans le service du registre de schémas pour les raisons suivantes :
- Le mappage du type Spark SQL au schéma Avro n'est pas un à un. Consultez la documentation Fichiers Avro.
- Si le schéma Avro de sortie converti est de type enregistrement, le nom de l'enregistrement est
topLevelRecordet il n'y a pas d'espace de noms par default.
Si le schéma de sortie default de to_avro correspond au schéma du sujet cible, vous pouvez effectuer les opérations suivantes :
// The converted data is saved to Kafka as a Kafka topic "t".
dataDF
.select(
to_avro($"key", lit("t-key"), schemaRegistryAddr).as("key"),
to_avro($"value", lit("t-value"), schemaRegistryAddr).as("value"))
.writeStream
.format("kafka")
.option("kafka.bootstrap.servers", servers)
.option("topic", "t")
.start()
Sinon, vous devez fournir le schéma du sujet cible dans la fonction to_avro :
// The Avro schema of subject "t-value" in JSON string format.
val avroSchema = ...
// The converted data is saved to Kafka as a Kafka topic "t".
dataDF
.select(
to_avro($"key", lit("t-key"), schemaRegistryAddr).as("key"),
to_avro($"value", lit("t-value"), schemaRegistryAddr, avroSchema).as("value"))
.writeStream
.format("kafka")
.option("kafka.bootstrap.servers", servers)
.option("topic", "t")
.start()
S’authentifier auprès d’un registre de schémas Confluent externe
Dans Databricks Runtime 12.2 LTS et versions ultérieures, vous pouvez vous authentifier auprès d'un registre de schémas Confluent externe. Les exemples suivants montrent comment configurer vos options de registre de schémas pour inclure les informations d’identification d’authentification et les clés API.
- Scala
- Python
import org.apache.spark.sql.avro.functions._
import scala.collection.JavaConverters._
val schemaRegistryAddr = "https://confluent-schema-registry-endpoint"
val schemaRegistryOptions = Map(
"confluent.schema.registry.basic.auth.credentials.source" -> "USER_INFO",
"confluent.schema.registry.basic.auth.user.info" -> "confluentApiKey:confluentApiSecret")
val df = spark
.readStream
.format("kafka")
.option("kafka.bootstrap.servers", servers)
.option("subscribe", "t")
.load()
.select(
from_avro(data = $"key", subject = "t-key", schemaRegistryAddress = schemaRegistryAddr, options = schemaRegistryOptions.asJava).as("key"),
from_avro(data = $"value", subject = "t-value", schemaRegistryAddress = schemaRegistryAddr, options = schemaRegistryOptions.asJava).as("value"))
// The converted data is saved to Kafka as a Kafka topic "t".
dataDF
.select(
to_avro($"key", lit("t-key"), schemaRegistryAddr, schemaRegistryOptions.asJava).as("key"),
to_avro($"value", lit("t-value"), schemaRegistryAddr, schemaRegistryOptions.asJava).as("value"))
.writeStream
.format("kafka")
.option("kafka.bootstrap.servers", servers)
.option("topic", "t")
.save()
// The Avro schema of subject "t-value" in JSON string format.
val avroSchema = ...
// The converted data is saved to Kafka as a Kafka topic "t".
dataDF
.select(
to_avro($"key", lit("t-key"), schemaRegistryAddr, schemaRegistryOptions.asJava).as("key"),
to_avro($"value", lit("t-value"), schemaRegistryAddr, schemaRegistryOptions.asJava, avroSchema).as("value"))
.writeStream
.format("kafka")
.option("kafka.bootstrap.servers", servers)
.option("topic", "t")
.save()
from pyspark.sql.functions import col, lit
from pyspark.sql.avro.functions import from_avro, to_avro
schema_registry_address = "https://confluent-schema-registry-endpoint"
schema_registry_options = {
"confluent.schema.registry.basic.auth.credentials.source": 'USER_INFO',
"confluent.schema.registry.basic.auth.user.info": f"{key}:{secret}"
}
df = (spark
.readStream
.format("kafka")
.option("kafka.bootstrap.servers", servers)
.option("subscribe", "t")
.load()
.select(
from_avro(
data = col("key"),
jsonFormatSchema = None,
options = schema_registry_options,
subject = "t-key",
schemaRegistryAddress = schema_registry_address
).alias("key"),
from_avro(
data = col("value"),
jsonFormatSchema = None,
options = schema_registry_options,
subject = "t-value",
schemaRegistryAddress = schema_registry_address
).alias("value")
)
)
# The converted data is saved to Kafka as a Kafka topic "t".
data_df
.select(
to_avro(
data = col("key"),
subject = lit("t-key"),
schemaRegistryAddress = schema_registry_address,
options = schema_registry_options
).alias("key"),
to_avro(
data = col("value"),
subject = lit("t-value"),
schemaRegistryAddress = schema_registry_address,
options = schema_registry_options
).alias("value")
)
.writeStream
.format("kafka")
.option("kafka.bootstrap.servers", servers)
.option("topic", "t")
.save()
# The Avro schema of subject "t-value" in JSON string format.
avro_schema = ...
# The converted data is saved to Kafka as a Kafka topic "t".
data_df
.select(
to_avro(
data = col("key"),
subject = lit("t-key"),
schemaRegistryAddress = schema_registry_address,
options = schema_registry_options
).alias("key"),
to_avro(
data = col("value"),
subject = lit("t-value"),
schemaRegistryAddress = schema_registry_address,
options = schema_registry_options,
jsonFormatSchema = avro_schema).alias("value"))
.writeStream
.format("kafka")
.option("kafka.bootstrap.servers", servers)
.option("topic", "t")
.save()
Utiliser les fichiers truststore et keystore dans les volumes Unity Catalog
Dans Databricks Runtime 14.3 LTS et versions ultérieures, vous pouvez utiliser les fichiers truststore et keystore dans les volumes Unity Catalog pour vous authentifier auprès d'un registre de schémas Confluent. Mettez à jour la configuration dans l'exemple précédent à l'aide de la syntaxe suivante :
val schemaRegistryAddr = "https://confluent-schema-registry-endpoint"
val schemaRegistryOptions = Map(
"confluent.schema.registry.ssl.truststore.location" -> "/Volumes/<catalog_name>/<schema_name>/<volume_name>/truststore.jks",
"confluent.schema.registry.ssl.truststore.password" -> "truststorePassword",
"confluent.schema.registry.ssl.keystore.location" -> "/Volumes/<catalog_name>/<schema_name>/<volume_name>/keystore.jks",
"confluent.schema.registry.ssl.truststore.password" -> "keystorePassword",
"confluent.schema.registry.ssl.key.password" -> "keyPassword")
Utiliser le mode d'évolution des schémas avec from_avro
Dans Databricks Runtime 14.2 et versions ultérieures, vous pouvez utiliser le mode d'évolution des schémas avec from_avro. L'activation du mode d'évolution des schémas fait en sorte que le Job lève une UnknownFieldException après la détection de l'évolution des schémas. Databricks recommande de configurer les Jobs avec le mode d'évolution des schémas pour un redémarrage automatique en cas d'échec de la tâche. Voir Considérations relatives à la production pour Structured Streaming.
L'évolution des schémas est utile si vous vous attendez à ce que le schéma de votre source de données évolue au fil du temps et ingère tous les champs de votre source de données. Si vos requêtes spécifient déjà explicitement les champs à interroger dans votre source de données, les champs ajoutés sont ignorés quelle que soit l'évolution des schémas.
Utilisez l'option avroSchemaEvolutionMode pour activer l'évolution des schémas. Le tableau suivant décrit les options du mode d'évolution des schémas :
Option | Comportement |
|---|---|
| default . Ignore l'évolution des schémas et le Job continue. |
| Lance une |
Vous pouvez modifier cette configuration entre les Jobs de streaming et réutiliser le même point de contrôle. La désactivation de l'évolution des schémas peut entraîner la suppression de colonnes.
Configurer le mode d'analyse
Vous pouvez configurer le mode d'analyse pour déterminer si vous souhaitez échouer ou émettre des enregistrements nuls lorsque le mode d'évolution des schémas est désactivé et que le schéma évolue d'une manière non rétrocompatible. Avec les paramètres default, from_avro échoue lorsqu'il observe des modifications de schéma incompatibles.
Utilisez l’option mode pour spécifier le mode d’analyse. Le tableau suivant décrit l’option du mode d’analyse :
Option | Comportement |
|---|---|
| default . Une erreur d’analyse génère un |
| Une erreur d'analyse est ignorée et un enregistrement nul est émis. |
Avec l'évolution des schémas activée, FAILFAST lève uniquement des exceptions si un enregistrement est corrompu.
Exemple d'utilisation de l'évolution des schémas et de définition du mode d'analyse
L'exemple suivant illustre l'activation de l'évolution des schémas et la spécification du mode d'analyse FAILFAST avec un Confluent Schema Registry :
- Scala
- Python
import org.apache.spark.sql.avro.functions._
import scala.collection.JavaConverters._
val schemaRegistryAddr = "https://confluent-schema-registry-endpoint"
val schemaRegistryOptions = Map(
"confluent.schema.registry.basic.auth.credentials.source" -> "USER_INFO",
"confluent.schema.registry.basic.auth.user.info" -> "confluentApiKey:confluentApiSecret",
"avroSchemaEvolutionMode" -> "restart",
"mode" -> "FAILFAST")
val df = spark
.readStream
.format("kafka")
.option("kafka.bootstrap.servers", servers)
.option("subscribe", "t")
.load()
.select(
// We read the "key" binary column from the subject "t-key" in the schema
// registry at schemaRegistryAddr. We provide schemaRegistryOptions,
// which has avroSchemaEvolutionMode -> "restart". This instructs from_avro
// to fail the query if the schema for the subject t-key evolves.
from_avro(
data = $"key",
subject = "t-key",
schemaRegistryAddress = schemaRegistryAddr,
options = schemaRegistryOptions.asJava).as("key"))
from pyspark.sql.functions import col, lit
from pyspark.sql.avro.functions import from_avro, to_avro
schema_registry_address = "https://confluent-schema-registry-endpoint"
schema_registry_options = {
"confluent.schema.registry.basic.auth.credentials.source": 'USER_INFO',
"confluent.schema.registry.basic.auth.user.info": f"{key}:{secret}",
"avroSchemaEvolutionMode": "restart",
"mode": "FAILFAST",
}
df = (spark
.readStream
.format("kafka")
.option("kafka.bootstrap.servers", servers)
.option("subscribe", "t")
.load()
.select(
from_avro(
data = col("key"),
jsonFormatSchema = None,
options = schema_registry_options,
subject = "t-key",
schemaRegistryAddress = schema_registry_address
).alias("key")
)
)