Aller au contenu principal

from_avro

Convertit une colonne binaire de format Avro en sa valeur catalyst correspondante. Le schéma spécifié doit correspondre aux données lues, sinon le comportement n'est pas défini : il peut échouer ou renvoyer un résultat arbitraire.

Si jsonFormatSchema n'est pas fourni, mais que subject et schemaRegistryAddress sont fournis, la fonction convertit une colonne binaire au format Avro du registre de schémas en sa valeur de catalyseur correspondante.

Syntaxe

Python
from pyspark.sql.avro.functions import from_avro

from_avro(data, jsonFormatSchema=None, options=None, subject=None, schemaRegistryAddress=None)

parameter

parameter

Type

Description

data

pyspark.sql.Column ou str

La colonne binaire contenant des données encodées en Avro.

jsonFormatSchema

str, facultatif

Le schéma Avro au format de chaîne JSON.

options

dictionnaire, facultatif

Options permettant de contrôler la façon dont l'enregistrement Avro est analysé et la configuration du client du registre de schémas.

subject

str, facultatif

Le sujet dans le Registre de schémas auquel appartiennent les données.

schemaRegistryAddress

str, facultatif

L'adresse (hôte et port) du registre de schémas.

parameter

Type

Description

data

pyspark.sql.Column ou str

La colonne binaire contenant des données encodées en Avro.

jsonFormatSchema

str, facultatif

Le schéma Avro au format de chaîne JSON.

options

dictionnaire, facultatif

Options permettant de contrôler la façon dont l'enregistrement Avro est analysé et la configuration du client du registre de schémas.

subject

str, facultatif

Le sujet dans le Registre de schémas auquel appartiennent les données.

schemaRegistryAddress

str, facultatif

L'adresse (hôte et port) du registre de schémas.

Options

Option

Valeurs

Description

mode

FAILFAST, PERMISSIVE

Mode de gestion des erreurs. default: FAILFAST. En mode PERMISSIVE, les enregistrements corrompus sont définis sur NULL au lieu de générer une erreur.

compression

uncompressed, snappy, deflate, bzip2, xz, zstandard

Codec de compression pour l'encodage des données Avro.

avroSchemaEvolutionMode

none, restart

Mode évolution des schémas. default: none. Lorsque défini sur restart, la query renvoie une UnknownFieldException lorsque le schéma change. Redémarrez le job pour utiliser le nouveau schéma. Voir Utiliser le mode évolution des schémas avec from_avro.

recursiveFieldMaxDepth

Plage : -1 à 15

Profondeur maximale de récursion le long d'un seul chemin récursif. default: -1, ce qui ne limite pas la profondeur de récursivité.

Lorsqu'un type partagé est accessible à partir de nombreux chemins de schéma distincts, l'expansion de schéma peut entraîner une panne de mémoire du Driver, car cette option ne limite la profondeur que sur un seul chemin. Pour contourner :

Option

Valeurs

Description

mode

FAILFAST, PERMISSIVE

Mode de gestion des erreurs. default: FAILFAST. En mode PERMISSIVE, les enregistrements corrompus sont définis sur NULL au lieu de générer une erreur.

compression

uncompressed, snappy, deflate, bzip2, xz, zstandard

Codec de compression pour l'encodage des données Avro.

avroSchemaEvolutionMode

none, restart

Mode évolution des schémas. default: none. Lorsque défini sur restart, la query renvoie une UnknownFieldException lorsque le schéma change. Redémarrez le job pour utiliser le nouveau schéma. Voir Utiliser le mode évolution des schémas avec from_avro.

recursiveFieldMaxDepth

Plage : -1 à 15

Profondeur maximale de récursion le long d'un seul chemin récursif. default: -1, ce qui ne limite pas la profondeur de récursivité.

Lorsqu'un type partagé est accessible à partir de nombreux chemins de schéma distincts, l'expansion de schéma peut entraîner une panne de mémoire du Driver, car cette option ne limite la profondeur que sur un seul chemin. Pour contourner :

Renvoie

pyspark.sql.Column: une nouvelle colonne contenant les données Avro désérialisées en tant que valeur de catalyseur correspondante.

Exemples

Exemple 1 : désérialisation d'une colonne binaire Avro à l'aide d'un schéma JSON

Python
from pyspark.sql import Row
from pyspark.sql.avro.functions import from_avro, to_avro

data = [(1, Row(age=2, name='Alice'))]
df = spark.createDataFrame(data, ("key", "value"))
avro_df = df.select(to_avro(df.value).alias("avro"))
json_format_schema = '''{"type":"record","name":"topLevelRecord","fields":
[{"name":"avro","type":[{"type":"record","name":"value",
"namespace":"topLevelRecord","fields":[{"name":"age","type":["long","null"]},
{"name":"name","type":["string","null"]}]},"null"]}]}'''
avro_df.select(from_avro(avro_df.avro, json_format_schema).alias("value")).show(truncate=False)
Output
+------------------+
|value |
+------------------+
|{{2, Alice}} |
+------------------+