Catatan
Akses ke halaman ini memerlukan otorisasi. Anda dapat mencoba masuk atau mengubah direktori.
Akses ke halaman ini memerlukan otorisasi. Anda dapat mencoba mengubah direktori.
Mengonversi kolom biner format Avro menjadi nilai katalis yang sesuai. Skema yang ditentukan harus cocok dengan data baca, jika tidak, perilaku tidak terdefinisi: mungkin gagal atau mengembalikan hasil arbitrer.
Jika jsonFormatSchema tidak disediakan tetapi keduanya subject dan schemaRegistryAddress disediakan, fungsi mengonversi kolom biner format Schema Registry Avro menjadi nilai katalis yang sesuai.
Sintaksis
from pyspark.sql.avro.functions import from_avro
from_avro(data, jsonFormatSchema=None, options=None, subject=None, schemaRegistryAddress=None)
Parameter-parameternya
| Parameter | Tipe | Deskripsi |
|---|---|---|
data |
pyspark.sql.Column atau str |
Kolom biner yang berisi data yang dikodekan Avro. |
jsonFormatSchema |
str, opsional | Skema Avro dalam format string JSON. |
options |
dict, opsional | Opsi untuk mengontrol bagaimana rekaman Avro diurai dan konfigurasi untuk klien registri skema. |
subject |
str, opsional | Subjek dalam Registri Skema tempat data berada. |
schemaRegistryAddress |
str, opsional | Alamat (host dan port) Registri Skema. |
Opsi
| Option | Nilai | Deskripsi |
|---|---|---|
mode |
FAILFAST, PERMISSIVE |
Mode penanganan kesalahan. Standar: FAILFAST. Dalam PERMISSIVE mode, rekaman yang rusak diatur ke NULL alih-alih meningkatkan kesalahan. |
compression |
uncompressed, , snappydeflate, bzip2, , xz,zstandard |
Kodek kompresi untuk mengodekan data Avro. |
avroSchemaEvolutionMode |
none, restart |
Mode evolusi skema. Standar: none. Saat diatur ke restart, kueri melempar UnknownFieldException saat skema berubah. Mulai ulang pekerjaan untuk menggunakan skema baru. Lihat Menggunakan mode evolusi skema dengan from_avro. |
recursiveFieldMaxDepth |
Rentang: -1 ke 15 |
Kedalaman rekursi maksimum di sepanjang satu jalur rekursif. Default: -1, yang tidak membatasi kedalaman rekursi.Ketika jenis bersama dapat dijangkau dari banyak jalur skema yang berbeda, ekspansi skema dapat menyebabkan driver kehabisan memori karena opsi ini mengikat kedalaman pada satu jalur saja. Untuk solusinya:
|
Pengembalian Barang
pyspark.sql.Column: Kolom baru yang berisi data Avro yang dideserialisasi sebagai nilai katalis yang sesuai.
Examples
Contoh 1: Mendeserialisasi kolom biner Avro menggunakan skema 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}} |
+------------------+