Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Буферы протокола (protobuf) — это нейтральный от языка формат двоичной сериализации, разработанный Google. Azure Databricks пользователи чаще всего сталкиваются с ним при обработке двоичных закодированных записей из систем потоковой передачи событий, таких как Apache Kafka. Azure Databricks поддерживает чтение и запись данных protobuf в Apache Spark с помощью функций from_protobuf и to_protobuf, которые преобразуют двоичные данные protobuf в структурные типы Spark SQL и обратно как для потоковых, так и для пакетных рабочих нагрузок.
Prerequisites
Функции Protobuf требуют Databricks Runtime 12.2 LTS и более поздних версий.
Синтаксис функции
Используйте from_protobuf, чтобы привести двоичный столбец к структуре, и to_protobuf, чтобы привести столбец со структурой к двоичному типу. Необходимо указать либо файл дескриптора, descFilePath определенный аргументом, либо реестр схем, указанный в аргументе options . Полный список параметров см. в разделе Protobuf.
Питон
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])
Options
Передайте параметры в from_protobuf и to_protobuf с помощью аргумента options. Полный список поддерживаемых параметров см. в разделе Protobuf.
Параметры реестра схем
Следующие параметры относятся к использованию реестра схем и не рассматриваются в справочнике по общим параметрам.
| Option | Обязательное поле | По умолчанию | Description |
|---|---|---|---|
schema.registry.schema.evolution.mode |
нет | "restart" |
Как изменения схемы обрабатываются при обнаружении более нового идентификатора схемы в входящей записи.
"restart" завершает запрос с ошибкой UnknownFieldException; настройте задания на перезапуск при сбое, чтобы изменения были подхвачены.
"none" игнорирует изменения идентификатора схемы и анализирует более новые записи с исходной схемой. |
confluent.schema.registry.<option> |
нет | — | Передайте любую опцию клиента Confluent Schema Registry, используя префикс "confluent.schema.registry". Например, задайте для "confluent.schema.registry.basic.auth.credentials.source" значение "USER_INFO", а для "confluent.schema.registry.basic.auth.user.info" — "<KEY>:<SECRET>", чтобы настроить базовую аутентификацию. |
Usage
В следующих примерах используется набор данных Wanderbricks, чтобы продемонстрировать сериализацию структур Apache Spark в двоичный формат protobuf с помощью to_protobuf() и десериализацию двоичных записей protobuf с помощью from_protobuf().
Используйте protobuf с реестром схем Confluent
Azure Databricks поддерживает использование реестра схем Confluent для определения Protobuf.
Питон
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()
Аутентификация во внешнем реестре схем Confluent
Чтобы выполнить проверку подлинности во внешнем реестре схем Confluent, обновите параметры реестра схем, чтобы включить учетные данные проверки подлинности и ключи API.
Питон
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 и хранилища ключей в томах каталога Unity
В Databricks Runtime 14.3 LTS и более поздних версиях можно использовать файлы truststore и keystore в разделах каталога Unity для аутентификации в реестр схем Confluent. Обновите параметры реестра схем в соответствии со следующим примером:
Питон
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>"
)
Использование Protobuf с дескрипторным файлом
Вы также можете ссылаться на файл дескриптора protobuf, доступный для вычислительного кластера. Убедитесь, что у вас есть правильные разрешения на чтение файла в зависимости от его расположения.
Питон
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()
Дополнительные ресурсы
-
Чтение и запись потоковых данных Avro. Если рабочая нагрузка потоковой передачи использует сериализацию Avro, а не Protobuf, см. функции потоковой передачи Avro для эквивалентных
from_avroиto_avroфункций.