重要
這項功能處於公開預覽狀態。
在 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 和下列選項之一:
topictopicstopicsPattern
如需選項的完整清單,請參閱 設定 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()