Not
Bu sayfaya erişim yetkilendirme gerektiriyor. Oturum açmayı veya dizinleri değiştirmeyi deneyebilirsiniz.
Bu sayfaya erişim yetkilendirme gerektiriyor. Dizinleri değiştirmeyi deneyebilirsiniz.
Ö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:
topictopicstopicsPattern
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
valuesütuna yüklenir. - Akışı farklı şemalarla birden çok konuyu okuyacak şekilde yapılandırıyorsanız ham içeriği bir
allowDifferentTopicSchemassü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()