from_avro

Преобразует двоичный столбец формата Avro в соответствующее значение катализатора. Указанная схема должна соответствовать данным чтения, в противном случае поведение не определено: оно может завершиться ошибкой или возвратом произвольного результата.

Если jsonFormatSchema это не указано, но subject оба и schemaRegistryAddress предоставлено, функция преобразует двоичный столбец формата Avro реестра схем в соответствующее значение катализатора.

Синтаксис

from pyspark.sql.avro.functions import from_avro

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

Параметры

Параметр Тип Описание
data pyspark.sql.Column или str Двоичный столбец, содержащий данные в кодировке Avro.
jsonFormatSchema str, необязательный Схема Avro в формате строки JSON.
options дикт, необязательный Параметры для управления анализом и настройкой записи Avro для клиента реестра схем.
subject str, необязательный Тема в реестре схем, к которой принадлежат данные.
schemaRegistryAddress str, необязательный Адрес (узел и порт) реестра схем.

Опции

Опция Ценности Описание
mode FAILFAST, PERMISSIVE Режим обработки ошибок. По умолчанию: FAILFAST. В PERMISSIVE режиме поврежденные записи задаются вместо того, чтобы NULL вызывать ошибку.
compression uncompressed, , snappydeflatebzip2xz,zstandard Кодек сжатия для кодирования данных Avro.
avroSchemaEvolutionMode none, restart Режим эволюции схемы. По умолчанию: none. Если задано значение restart, запрос создает исключение UnknownFieldException при изменении схемы. Перезапустите задание, чтобы использовать новую схему. См. раздел "Использование режима эволюции схемы" с from_avro.
recursiveFieldMaxDepth Диапазон: -1 до 15 Максимальная глубина рекурсии вдоль одного рекурсивного пути. Значение по умолчанию: -1не ограничивает глубину рекурсии.
Если общий тип доступен из многих разных путей схемы, расширение схемы может привести к нехватке памяти драйвера, так как этот параметр ограничивает глубину только одного пути. Для обходного решения:

Возвраты

pyspark.sql.Column: новый столбец, содержащий десериализированные данные Avro в качестве соответствующего значения катализатора.

Примеры

Пример 1. Десериализация двоичного столбца Avro с помощью схемы JSON

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)
+------------------+
|value             |
+------------------+
|{{2, Alice}}      |
+------------------+