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