從 Apache Pulsar 串流

重要

這項功能處於公開預覽狀態。

在 Databricks Runtime 14.1 及以上版本中,您可以使用 Structured Streaming,在 Azure Databricks 上從 Apache Pulsar 串流處理資料。

結構化串流針對從 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 DataFrame 中。
  • 對於在 Pulsar 中沒有結構或具有簡單數據類型的主題,負載會被載入至 value 資料行。
  • 如果你設定串流能讀取多個不同結構的主題,設定 allowDifferentTopicSchemas 將原始內容載入欄位 value 。

Pulsar 記錄具有下列元數據欄位:

列 類型
__key binary
__topic string
__messageId binary
__publishTime timestamp
__eventTime timestamp
__messageProperties map<String, String>

設定 Pulsar 串流讀取選項

完整選項列表請參見 Pulsar。

建構起始位移 JSON

若要使用指定偏移量的自訂訊息 ID,如 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()