Megjegyzés
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhat bejelentkezni vagy módosítani a címtárat.
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhatja módosítani a címtárat.
A protokollpufferek (protobuf) a Google által kifejlesztett nyelvsemleges bináris szerializálási formátum. Azure Databricks felhasználók leggyakrabban akkor találkoznak, ha binárisan kódolt rekordokat dolgoznak fel eseménystreamelési rendszerekből, például az Apache Kafkából. Azure Databricks támogatja a protobuf-adatok olvasását és írását az Apache Sparktal a from_protobufto_protobuf függvényeken keresztül, amelyek bináris protobuf és Spark SQL-struktúratípusok közötti konvertálást támogatnak a streamelési és kötegelt számítási feladatokhoz.
Prerequisites
A Protobuf-függvényekhez a Databricks Runtime 12.2 LTS és újabb verziók szükségesek.
Függvényszintaxis
A bináris oszlop struktúrává alakításához használja a from_protobuf elemet, a struktúraoszlop binárissá alakításához pedig a to_protobuf elemet. Meg kell adnia vagy a(z) descFilePath argumentum által azonosított leírófájlt, vagy a(z) options argumentummal megadott sémaregisztrációt. A lehetőségek teljes listáját a Protobufban találja.
Python
from_protobuf(data: 'ColumnOrName', messageName: Optional[str] = None, descFilePath: Optional[str] = None, options: Optional[Dict[str, str]] = None)
to_protobuf(data: 'ColumnOrName', messageName: Optional[str] = None, descFilePath: Optional[str] = None, options: Optional[Dict[str, str]] = None)
Scala
// While using with Schema registry:
from_protobuf(data: Column, options: Map[String, String])
// Or with Protobuf descriptor file:
from_protobuf(data: Column, messageName: String, descFilePath: String, options: Map[String, String])
// While using with Schema registry:
to_protobuf(data: Column, options: Map[String, String])
// Or with Protobuf descriptor file:
to_protobuf(data: Column, messageName: String, descFilePath: String, options: Map[String, String])
Beállítások
Adja át a beállításokat a(z) from_protobuf és to_protobuf számára a(z) options argumentum használatával. A támogatott lehetőségek teljes listáját a Protobufban találja.
A sémaregiszter beállításai
Az alábbi beállítások a sémaregisztrációs adatbázis használatára vonatkoznak, és nem vonatkoznak az általános beállításokra vonatkozó hivatkozásra.
| Option | Kötelező | Alapértelmezett | Leírás |
|---|---|---|---|
schema.registry.schema.evolution.mode |
No | "restart" |
A sémaváltozások kezelése, amikor egy beérkező rekordban újabb sémaazonosító észlelhető.
"restart" megszakítja a lekérdezést egy UnknownFieldException elemmel; a feladatokat úgy kell konfigurálni, hogy hiba esetén újrainduljanak a módosítások érvénybe léptetéséhez.
"none" figyelmen kívül hagyja a sémaazonosító módosításait, és elemzi az újabb rekordokat az eredeti sémával. |
confluent.schema.registry.<option> |
No | — | Adja meg a Confluent Schema Registry kliens bármely beállítását a(z) "confluent.schema.registry" előtag használatával. Például az alapszintű hitelesítés konfigurálásához állítsa a(z) "confluent.schema.registry.basic.auth.credentials.source" értékét "USER_INFO"-ra, a(z) "confluent.schema.registry.basic.auth.user.info" értékét pedig "<KEY>:<SECRET>"-ra. |
Usage
Az alábbi példák a Wanderbricks-adatkészlettel szemléltetik az Apache Spark-szerkezetek bináris protobufra való szerializálását bináris protobuf rekordokkal to_protobuf() , és deszerializálják a bináris protobuf rekordokat a következővel from_protobuf(): .
A protobuf használata a Confluent-sémaregisztrációs adatbázissal
Az Azure Databricks támogatja a Confluent sémaregisztrációs adatbázisának használatát a Protobuf definiálásához.
Python
from pyspark.sql.protobuf.functions import to_protobuf, from_protobuf
from pyspark.sql.functions import struct
schema_registry_options = {
"schema.registry.subject" : "app-events-value",
"schema.registry.address" : "https://schema-registry:8081/"
}
# Serialize Wanderbricks reviews to binary Protobuf using schema registry
reviews_df = spark.read.table("samples.wanderbricks.reviews")
proto_bytes_df = reviews_df.select(
to_protobuf(struct("review_id", "rating", "comment"), options=schema_registry_options).alias("proto_bytes")
)
# Deserialize binary Protobuf records back to a struct
reviews_restored_df = proto_bytes_df.select(
from_protobuf("proto_bytes", options=schema_registry_options).alias("proto_event")
)
display(reviews_restored_df)
Scala
import org.apache.spark.sql.protobuf.functions._
import org.apache.spark.sql.functions.struct
import scala.collection.JavaConverters._
val schemaRegistryOptions = Map(
"schema.registry.subject" -> "app-events-value",
"schema.registry.address" -> "https://schema-registry:8081/"
)
// Serialize Wanderbricks reviews to binary Protobuf using schema registry
val reviewsDF = spark.read.table("samples.wanderbricks.reviews")
val protoBytesDF = reviewsDF.select(
to_protobuf(struct($"review_id", $"rating", $"comment"), options = schemaRegistryOptions.asJava)
.as("proto_bytes")
)
// Deserialize binary Protobuf records back to a struct
val reviewsRestoredDF = protoBytesDF.select(
from_protobuf($"proto_bytes", options = schemaRegistryOptions.asJava)
.as("proto_event")
)
reviewsRestoredDF.show()
Hitelesítés külső Confluent-sémaregisztrációs adatbázisba
Ha külső Confluent-sémaregisztrációs adatbázisba szeretne hitelesítést végezni, frissítse a sémaregisztrációs beállításokat úgy, hogy azok tartalmazzák a hitelesítési hitelesítő adatokat és az API-kulcsokat.
Python
schema_registry_options = {
"schema.registry.subject" : "app-events-value",
"schema.registry.address" : "https://remote-schema-registry-endpoint",
"confluent.schema.registry.basic.auth.credentials.source" : "USER_INFO",
"confluent.schema.registry.basic.auth.user.info" : "confluentApiKey:confluentApiSecret"
}
Scala
val schemaRegistryOptions = Map(
"schema.registry.subject" -> "app-events-value",
"schema.registry.address" -> "https://remote-schema-registry-endpoint",
"confluent.schema.registry.basic.auth.credentials.source" -> "USER_INFO",
"confluent.schema.registry.basic.auth.user.info" -> "confluentApiKey:confluentApiSecret"
)
Truststore- és keystore-fájlok használata Unity Catalog-kötetekben
A Databricks Runtime 14.3 LTS-ben és újabb verziókban a Unity Catalog-kötetekben található truststore- és keystore-fájlokat használhatja a Confluent-sémaregisztrációs adatbázisba való hitelesítéshez. Frissítse a sémaregisztrációs adatbázis beállításait az alábbi példának megfelelően:
Python
schema_registry_options = {
"schema.registry.subject" : "app-events-value",
"schema.registry.address" : "https://remote-schema-registry-endpoint",
"confluent.schema.registry.ssl.truststore.location" : "/Volumes/<catalog_name>/<schema_name>/<volume_name>/kafka.client.truststore.jks",
"confluent.schema.registry.ssl.truststore.password" : "<password>",
"confluent.schema.registry.ssl.keystore.location" : "/Volumes/<catalog_name>/<schema_name>/<volume_name>/kafka.client.keystore.jks",
"confluent.schema.registry.ssl.keystore.password" : "<password>",
"confluent.schema.registry.ssl.key.password" : "<password>"
}
Scala
val schemaRegistryOptions = Map(
"schema.registry.subject" -> "app-events-value",
"schema.registry.address" -> "https://remote-schema-registry-endpoint",
"confluent.schema.registry.ssl.truststore.location" -> "/Volumes/<catalog_name>/<schema_name>/<volume_name>/kafka.client.truststore.jks",
"confluent.schema.registry.ssl.truststore.password" -> "<password>",
"confluent.schema.registry.ssl.keystore.location" -> "/Volumes/<catalog_name>/<schema_name>/<volume_name>/kafka.client.keystore.jks",
"confluent.schema.registry.ssl.keystore.password" -> "<password>",
"confluent.schema.registry.ssl.key.password" -> "<password>"
)
A Protobuf használata leíró fájllal
Hivatkozhat egy protobuf leírófájlra is, amely elérhető a számítási fürt számára. Győződjön meg arról, hogy megfelelő engedélyekkel rendelkezik a fájl olvasásához a helyétől függően.
Python
from pyspark.sql.protobuf.functions import to_protobuf, from_protobuf
from pyspark.sql.functions import struct
descriptor_file = "/path/to/proto_descriptor.desc"
# Serialize Wanderbricks reviews to binary Protobuf using a descriptor file
reviews_df = spark.read.table("samples.wanderbricks.reviews")
proto_bytes_df = reviews_df.select(
to_protobuf(struct("review_id", "rating", "comment"), "Review", descriptor_file).alias("proto_bytes")
)
# Deserialize binary Protobuf records back to a struct
reviews_restored_df = proto_bytes_df.select(
from_protobuf("proto_bytes", "Review", descFilePath=descriptor_file).alias("review")
)
display(reviews_restored_df)
Scala
import org.apache.spark.sql.protobuf.functions._
import org.apache.spark.sql.functions.struct
val descriptorFile = "/path/to/proto_descriptor.desc"
// Serialize Wanderbricks reviews to binary Protobuf using a descriptor file
val reviewsDF = spark.read.table("samples.wanderbricks.reviews")
val protoBytesDF = reviewsDF.select(
to_protobuf(struct($"review_id", $"rating", $"comment"), "Review", descriptorFile).as("proto_bytes")
)
// Deserialize binary Protobuf records back to a struct
val reviewsRestoredDF = protoBytesDF.select(
from_protobuf($"proto_bytes", "Review", descFilePath=descriptorFile).as("review")
)
reviewsRestoredDF.show()
További források
-
Streamelési Avro-adatok olvasása és írása: Ha a streamelési számítási feladat a Protobuf helyett Avro szerializálást használ, tekintse meg az Avro streamelési függvényeit az egyenértékű
from_avroésto_avroa függvényekhez.