Apache Pulsar'dan veri akışı

Önemli

Bu özellik Genel Önizlemededir.

Databricks Runtime 14.1 ve üzerinde Yapılandırılmış Akış'ı kullanarak Azure Databricks üzerinde Apache Pulsar'dan veri akışı yapabilirsiniz.

Yapılandırılmış Akış, Pulsar kaynaklarından okunan veriler için tam olarak bir kez işleme semantiği sağlar.

Söz dizimi örneği

Aşağıda, Pulsar'dan okumak için Yapılandırılmış Akış kullanmanın temel bir örneği verilmiştir:

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 konularından okumak için bir service.url ve aşağıdaki seçeneklerden birini sağlamalısınız:

  • topic
  • topics
  • topicsPattern

Seçeneklerin tam listesi için Pulsar akışından okuma için seçenekleri yapılandırma bölümüne bakın.

Pulsar'da kimlik doğrulaması

Azure Databricks, Pulsar'da truststore ve keystore kimlik doğrulamasını destekler. Databricks, yapılandırma ayrıntılarını depolamak için gizli dizileri kullanmanızı önerir.

Kimlik doğrulama seçeneklerinin tam listesi için bkz. Kimlik doğrulaması.

Example

Aşağıdaki örnekte kimlik doğrulama seçeneklerinin yapılandırılması gösterilmektedir:

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 şeması

Pulsar'dan okuduğunuzda, satırların şeması kaynağın konu başlıklarının şemalarına bağlıdır.

  • Avro veya JSON şemasıyla ilgili konular için, sonuçta elde edilen Spark DataFrame'de alan adları ve alan türleri korunur.
  • Pulsar'da şemasız veya basit bir veri türüne sahip konular için yük bir value sütuna yüklenir.
  • Akışı farklı şemalarla birden çok konuyu okuyacak şekilde yapılandırıyorsanız ham içeriği bir allowDifferentTopicSchemas sütuna yüklenecek şekilde ayarlayınvalue.

Pulsar kayıtları aşağıdaki meta veri alanlarına sahiptir:

Sütun Tip
__key binary
__topic string
__messageId binary
__publishTime timestamp
__eventTime timestamp
__messageProperties map<String, String>

Pulsar akış okuma seçeneklerini yapılandırma

Seçeneklerin tam listesi için bkz . Pulsar.

Başlangıç uzaklıklarını oluşturma JSON

JSON olarak bir uzaklık belirten özel ileti kimliğini seçeneğiyle startingOffsets kullanmak için aşağıdaki örneğe bakın:

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()