Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Внимание
Эта функция предоставляется в режиме общедоступной предварительной версии.
В Databricks Runtime 14.1 и более поздних версиях можно использовать структурированную потоковую передачу для потоковой передачи данных из Apache Pulsar в Azure Databricks.
Structured Streaming обеспечивает семантику обработки «ровно один раз» для данных, считываемых из источников Pulsar.
Пример синтаксиса
Ниже приведен базовый пример использования структурированной потоковой передачи для чтения из Pulsar:
Python
query = (spark.readStream
.format("pulsar")
.option("service.url", "pulsar://broker.example.com:6650")
.option("topics", "topic1,topic2")
.load()
)
Scala
val query = spark.readStream
.format("pulsar")
.option("service.url", "pulsar://broker.example.com:6650")
.option("topics", "topic1,topic2")
.load()
Чтобы считывать данные из топиков Pulsar, необходимо указать service.url и один из следующих вариантов:
topictopicstopicsPattern
Для полного списка параметров см. раздел Настройка параметров потоковой передачи Pulsar.
Аутентифицируйтесь в Pulsar
Azure Databricks поддерживает аутентификацию для Pulsar с использованием хранилища доверенных сертификатов и хранилища ключей. Databricks рекомендует использовать секреты для хранения сведений о конфигурации.
Полный список параметров проверки подлинности см. в разделе "Проверка подлинности".
Example
В следующем примере демонстрируется настройка параметров проверки подлинности:
Python
client_auth_params = dbutils.secrets.get(scope="pulsar", key="clientAuthParams")
client_pw = dbutils.secrets.get(scope="pulsar", key="clientPw")
# clientAuthParams is a comma-separated list of key-value pairs, such as:
# "keyStoreType:JKS,keyStorePath:/var/private/tls/client.keystore.jks,keyStorePassword:clientpw"
query = (spark.readStream
.format("pulsar")
.option("service.url", "pulsar://broker.example.com:6650")
.option("topics", "topic1,topic2")
.option("startingOffsets", starting_offsets)
.option("pulsar.client.authPluginClassName", "org.apache.pulsar.client.impl.auth.AuthenticationKeyStoreTls")
.option("pulsar.client.authParams", client_auth_params)
.option("pulsar.client.useKeyStoreTls", "true")
.option("pulsar.client.tlsTrustStoreType", "JKS")
.option("pulsar.client.tlsTrustStorePath", trust_store_path)
.option("pulsar.client.tlsTrustStorePassword", client_pw)
.load()
)
Scala
val clientAuthParams = dbutils.secrets.get(scope = "pulsar", key = "clientAuthParams")
val clientPw = dbutils.secrets.get(scope = "pulsar", key = "clientPw")
// clientAuthParams is a comma-separated list of key-value pairs, such as:
// "keyStoreType:JKS,keyStorePath:/var/private/tls/client.keystore.jks,keyStorePassword:clientpw"
val query = spark.readStream
.format("pulsar")
.option("service.url", "pulsar://broker.example.com:6650")
.option("topics", "topic1,topic2")
.option("startingOffsets", startingOffsets)
.option("pulsar.client.authPluginClassName", "org.apache.pulsar.client.impl.auth.AuthenticationKeyStoreTls")
.option("pulsar.client.authParams", clientAuthParams)
.option("pulsar.client.useKeyStoreTls", "true")
.option("pulsar.client.tlsTrustStoreType", "JKS")
.option("pulsar.client.tlsTrustStorePath", trustStorePath)
.option("pulsar.client.tlsTrustStorePassword", clientPw)
.load()
Схема Pulsar
При чтении из Pulsar схема строк зависит от схем тем источника.
- Для разделов с схемой Avro или JSON имена полей и типы полей сохраняются в результирующем кадре данных Spark.
- Для разделов без схемы или с простым типом данных в Pulsar полезные данные загружаются в столбец
value. - Если вы настроите поток для чтения нескольких топиков с разными схемами, задайте для
allowDifferentTopicSchemasзагрузку необработанного содержимого в столбецvalue.
Записи Pulsar имеют следующие поля метаданных:
| Столбец | Тип |
|---|---|
__key |
binary |
__topic |
string |
__messageId |
binary |
__publishTime |
timestamp |
__eventTime |
timestamp |
__messageProperties |
map<String, String> |
Настройте параметры потокового чтения Pulsar
Полный список параметров см. в разделе Pulsar.
Формирование JSON начальных смещений
Чтобы использовать пользовательский идентификатор сообщения, который задаёт смещение, в формате JSON с параметром startingOffsets, см. следующий пример:
import org.apache.spark.sql.pulsar.JsonUtils
import org.apache.pulsar.client.api.MessageId
import org.apache.pulsar.client.impl.MessageIdImpl
val topic = "my-topic"
val msgId: MessageId = new MessageIdImpl(ledgerId, entryId, partitionIndex)
val startOffsets = JsonUtils.topicOffsets(Map(topic -> msgId))
query = spark.readStream
.format("pulsar")
.option("service.url", "pulsar://broker.example.com:6650")
.option("topic", topic)
.option("startingOffsets", startOffsets)
.load()