Catatan
Akses ke halaman ini memerlukan otorisasi. Anda dapat mencoba masuk atau mengubah direktori.
Akses ke halaman ini memerlukan otorisasi. Anda dapat mencoba mengubah direktori.
Penting
Fitur ini ada di Pratinjau Publik.
Di Databricks Runtime 14.1 ke atas, Anda dapat menggunakan Streaming Terstruktur untuk mengalirkan data dari Apache Pulsar di Azure Databricks.
Structured Streaming menyediakan semantik pemrosesan sekali saja untuk data yang dibaca dari sumber data Pulsar.
Contoh sintaksis
Berikut ini adalah contoh dasar menggunakan Streaming Terstruktur untuk membaca dari 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()
Untuk membaca dari topik Pulsar, Anda harus menyediakan service.url dan salah satu opsi berikut:
topictopicstopicsPattern
Untuk daftar lengkap opsi, lihat Mengonfigurasi opsi untuk streaming Pulsar membaca.
Mengautentikasi ke Pulsar
Azure Databricks mendukung autentikasi truststore dan keystore untuk Pulsar. Databricks merekomendasikan agar Anda menggunakan rahasia untuk menyimpan detail konfigurasi.
Untuk daftar lengkap opsi autentikasi, lihat Autentikasi.
Example
Contoh berikut menunjukkan konfigurasi opsi autentikasi:
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()
Skema Pulsar
Saat Anda membaca dari Pulsar, skema baris tergantung pada skema topik sumber.
- Untuk topik dengan skema Avro atau JSON, nama bidang dan jenis bidang dipertahankan dalam Spark DataFrame yang dihasilkan.
- Untuk topik tanpa skema atau dengan jenis data sederhana di Pulsar, payload dimuat ke kolom
value. - Jika Anda mengonfigurasi aliran untuk membaca beberapa topik dengan skema yang berbeda, atur
allowDifferentTopicSchemasuntuk memuat konten mentah kevaluekolom.
Rekaman Pulsar memiliki bidang metadata berikut:
| Kolom | Tipe |
|---|---|
__key |
binary |
__topic |
string |
__messageId |
binary |
__publishTime |
timestamp |
__eventTime |
timestamp |
__messageProperties |
map<String, String> |
Konfigurasikan opsi untuk pembacaan streaming Pulsar
Untuk daftar lengkap opsi, lihat Pulsar.
Susun JSON offset awal
Untuk menggunakan ID pesan kustom yang menentukan offset, sebagai JSON, dengan startingOffsets opsi , lihat contoh berikut:
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()