from_avro

將 Avro 格式的二進位欄位轉換為對應的催化劑值。 指定的結構必須與讀取資料相符,否則行為未定義:可能失敗或回傳任意結果。

若jsonFormatSchema未提供,但同時subjectschemaRegistryAddress提供與,函式會將 Schema Registry 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 力量,選用 Avro 的 JSON 字串格式結構。
options DICT,選擇性 控制 Avro 記錄如何解析及架構登錄用戶端的設定選項。
subject 力量,選用 資料所屬的結構登錄檔主體。
schemaRegistryAddress 力量,選用 結構登錄的位址(主機與埠)。

選項

Option 價值觀 說明
mode FAILFAST、PERMISSIVE 錯誤處理模式。 預設值:FAILFAST。 在 PERMISSIVE 模式下,損壞紀錄會被設為 , NULL 而不是產生錯誤。
compression uncompressed、snappy、deflate、bzip2、xz、zstandard 用於編碼 Avro 資料的壓縮編解碼器。
avroSchemaEvolutionMode none、restart 圖式演化模式。 預設值:none。 當設定為 restart時,查詢會拋出 和 UnknownFieldException 當結構改變時。 重新啟動工作以使用新的架構。 請參見 「使用模式演化模式搭配 from_avro」。
recursiveFieldMaxDepth 距離: -1 到 15 沿著單一遞迴路徑的最大遞迴深度。 預設值: -1,不限制遞迴深度。
當共享型別可從多個不同的結構路徑達成時,模式擴展可能導致驅動程式記憶體不足,因為此選項只限制了一條路徑的深度。 解決方法:

退貨

pyspark.sql.Column: 一個新欄位,包含已解序列化的 Avro 資料作為對應的催化劑值。

Examples

範例 1:使用 JSON 架構反序列化 Avro 二進位欄位

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}}      |
+------------------+