使用內建的連接器訂閱 Google Pub/Sub。 此連接器對來自訂閱者的資料列提供恰好一次處理語意。
備註
Pub/Sub 可能會發布重複的資料列,或者資料列到達訂閱者時可能不按順序。 你必須寫程式碼來處理重複和亂序的列。
設定 Pub/Sub 串流
以下程式碼範例展示了如何配置從 Pub/Sub 讀取結構化串流並使用私鑰進行認證。
Python
auth_options = {
"clientId": client_id,
"clientEmail": client_email,
"privateKey": private_key,
"privateKeyId": private_key_id
}
query = (spark.readStream
.format("pubsub")
.option("subscriptionId", "mysub")
.option("topicId", "mytopic")
.option("projectId", "myproject")
.options(auth_options)
.load()
)
Scala
val authOptions: Map[String, String] =
Map("clientId" -> clientId,
"clientEmail" -> clientEmail,
"privateKey" -> privateKey,
"privateKeyId" -> privateKeyId)
val query = spark.readStream
.format("pubsub")
// Creates a Pub/Sub subscription if one does not already exist with this ID
.option("subscriptionId", "mysub")
.option("topicId", "mytopic")
.option("projectId", "myproject")
.options(authOptions)
.load()
SQL
CREATE OR REFRESH STREAMING TABLE pubsub_raw
AS SELECT * FROM STREAM read_pubsub(
subscriptionId => 'mysub',
projectId => 'myproject',
topicId => 'mytopic',
clientEmail => secret('pubsub-scope', 'clientEmail'),
clientId => secret('pubsub-scope', 'clientId'),
privateKeyId => secret('pubsub-scope', 'privateKeyId'),
privateKey => secret('pubsub-scope', 'privateKey')
);
如需更多配置選項,請參閱 Pub/Sub 串流讀取的配置選項。
設定對 Pub/Sub 的存取
您的資格必須具備以下職務:
| 角色 | 必要或選用 | 角色的運用方式 |
|---|---|---|
roles/pubsub.viewer 或 roles/viewer |
必要 | 檢查訂閱是否存在並獲取訂閱。 |
roles/pubsub.subscriber |
必要 | 從訂閱中擷取資料。 |
roles/pubsub.editor 或 roles/editor |
選擇性 | 如果不存在訂閱,可以啟用建立訂閱的功能,並在資料流終止時使用 deleteSubscriptionOnStreamStop 來刪除訂閱。 |
備註
如果您是在資源層級而非專案層級授與 roles/pubsub.viewer 和 roles/pubsub.subscriber,則必須將這兩個角色同時套用至主題和訂閱項目。 如果您未使用選用的 roles/pubsub.editor 或 roles/editor 角色,僅在該主題上授予必要角色仍不足夠。
Databricks 建議使用金鑰時使用秘密。 授權連線需要下列選項:
clientEmailclientIdprivateKeyprivateKeyId
了解發佈/訂閱架構
串流的結構與從 Pub/Sub 抓取的列相符,詳見下表:
| 欄位 | 類型 | Description |
|---|---|---|
messageId |
StringType |
Pub/Sub 在發佈訊息時指派給該訊息的唯一 ID。 |
payload |
ArrayType[ByteType] |
訊息資料以原始位元組形式呈現。 |
attributes |
StringType |
訊息的屬性,作為一個由鍵值對組成的 JSON 字串。 |
publishTimestampInMillis |
LongType |
Pub/Sub 發布該訊息的時間,自 Unix 時代(1970 年 1 月 1 日 00:00:00 UTC)以來的毫秒計。 |
配置 Pub/Sub 串流讀取選項
Pub/Sub 的某些設定選項使用擷取而非微批次。 這是內部實作細節,選項的運作方式與其他結構化串流連接器相似,只是列先取回再處理。
完整選項列表請參見 出版/訂閱。
使用增量批次處理搭配 Pub/Sub
你可以使用 Trigger.AvailableNow 以增量批次方式取用來自 Pub/Sub 來源的可用資料列。
Azure Databricks 會在您使用 Trigger.AvailableNow 設定開始進行讀取操作時記錄時間戳。 由該批次處理的資料列包含所有先前已擷取的資料,以及任何其時間戳記早於已記錄的開始時間戳記的新發佈資料列。 欲了解更多資訊,請參閱 AvailableNow:增量批次處理。
監控發佈/訂閱串流指標
結構化串流進度指標報告已擷取並準備處理的資料列數、已擷取並準備處理的資料列大小,以及自串流開始以來所見的重複數量。
以下是 Pub/Sub 指標的範例:
"metrics" : {
"numDuplicatesSinceStreamStart" : "1",
"numRecordsReadyToProcess" : "1",
"sizeOfRecordsReadyToProcess" : "8"
}
限制
Pub/Sub 不支援搭配 spark.speculation 的推測執行。