訂閱Google Pub/Sub

使用內建的連接器訂閱 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.viewerroles/viewer 必要 檢查訂閱是否存在並獲取訂閱。
roles/pubsub.subscriber 必要 從訂閱中擷取資料。
roles/pubsub.editorroles/editor 選擇性 如果不存在訂閱,可以啟用建立訂閱的功能,並在資料流終止時使用 deleteSubscriptionOnStreamStop 來刪除訂閱。

備註

如果您是在資源層級而非專案層級授與 roles/pubsub.viewerroles/pubsub.subscriber,則必須將這兩個角色同時套用至主題和訂閱項目。 如果您未使用選用的 roles/pubsub.editorroles/editor 角色,僅在該主題上授予必要角色仍不足夠。

Databricks 建議使用金鑰時使用秘密。 授權連線需要下列選項:

  • clientEmail
  • clientId
  • privateKey
  • privateKeyId

了解發佈/訂閱架構

串流的結構與從 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 的推測執行。