Поток из Apache Pulsar

Внимание

Эта функция предоставляется в режиме общедоступной предварительной версии.

В 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 и один из следующих вариантов:

  • topic
  • topics
  • topicsPattern

Для полного списка параметров см. раздел Настройка параметров потоковой передачи 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()