Streaming dari Apache Pulsar

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:

  • topic
  • topics
  • topicsPattern

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 allowDifferentTopicSchemas untuk memuat konten mentah ke value kolom.

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