將欄位轉換成 Avro 格式的二進位檔。
若 subject 同時提供 和 schemaRegistryAddress ,函式會將欄位轉換為 Schema Registry Avro 格式的二進位。 輸入資料結構必須已在結構登錄檔中註冊給該主體,否則查詢在執行時會失敗。
語法
from pyspark.sql.avro.functions import to_avro
to_avro(data, jsonFormatSchema=None, subject=None, schemaRegistryAddress=None, options=None)
參數
| 參數 | 類型 | 說明 |
|---|---|---|
data |
pyspark.sql.Column 或 str |
資料欄位要序列化。 |
jsonFormatSchema |
力量,選用 | 使用者指定的 Avro 架構,格式為 JSON 字串。 |
subject |
pyspark.sql.Column 或是力量,選用 |
資料所屬的結構登錄檔主體。 |
schemaRegistryAddress |
力量,選用 | 結構登錄的位址(主機與埠)。 |
options |
DICT,選擇性 | 控制 Avro 記錄序列化的選項,以及架構登錄用戶端的設定。 |
退貨
pyspark.sql.Column:一個包含 Avro 編碼二進位資料的新欄位。
Examples
範例 1:將字串欄位轉換為 Avro 二進位格式
from pyspark.sql.avro.functions import to_avro
data = ['SPADES']
df = spark.createDataFrame(data, "string")
df.select(to_avro(df.value).alias("avro")).show(truncate=False)
+--------------------+
|avro |
+--------------------+
|[00 0C 53 50 41 4...|
+--------------------+
範例 2:使用自訂 JSON 架構將字串欄位轉換成 Avro
from pyspark.sql.avro.functions import to_avro
data = ['SPADES']
df = spark.createDataFrame(data, "string")
json_format_schema = '''["null", {"type": "enum", "name": "value",
"symbols": ["SPADES", "HEARTS", "DIAMONDS", "CLUBS"]}]'''
df.select(to_avro(df.value, json_format_schema).alias("avro")).show(truncate=False)
+--------+
|avro |
+--------+
|[02 00] |
+--------+