read_pulsar 串流數據表值函式

適用於:核取為「是」 Databricks SQL 核取為「是」 Databricks Runtime 14.1 和以上版本

重要

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

傳回資料表,其中包含從 Pulsar 讀取的記錄。

此數據表值函式僅支援串流處理,而不支援批次查詢。

語法

read_pulsar ( { option_key => option_value } [, ...] )

引數

此函式需要使用具名參數調用選項鍵。

serviceUrl和topic是必要選項。

這裡的參數描述很簡短。 如需擴充描述,請參閱 結構化串流 Pulsar 檔。

選項 類型 預設 描述
服務網址 (serviceUrl) 字串 必要 Pulsar 服務的 URI。
主題 字串 必要 要從中讀取的主題。
預定義訂閱 字串 無 連接器用來追蹤 Spark 應用程式進度的預先定義訂用帳戶名稱。
訂閱前綴 字串 無 連接器用以產生隨機訂閱以追踪 Spark 應用程式進度的前綴。
pollTimeoutMs LONG 120000 從 Pulsar 讀取訊息的逾時時間,以毫秒為單位。
發生數據丟失時失敗 BOOLEAN true 控制當資料遺失時是否使查詢失敗(例如,主題被刪除,或由於保留政策訊息被刪除)。
起始偏移量 字串 最新 查詢啟動時的起點,可以是最早、最新或指定特定位移的 JSON 字串。 如果最新,讀取器會在開始執行之後讀取最新的記錄。 如果是最早的話,讀取器會從最早的位移開始讀取。 使用者也可以指定一個設定特定偏移量的 JSON 字串。
開始時間 字串 無 指定時,Pulsar 來源會讀取從指定 startingTime 位置開始的訊息。

下列自變數用於驗證脈衝星用戶端:

選項 類型 預設 描述
pulsarClientAuthPlugin類別名稱 (驗證插件類別名稱) 字串 無 驗證外掛程式的名稱。
pulsarClientAuthParams 字串 無 驗證外掛程式的參數。
pulsarClientUseKeyStoreTls 字串 無 是否要使用 KeyStore 進行 Tls 驗證。
pulsarClientTlsTrustStoreType 字串 無 TLS 驗證的 TrustStore 檔案類型。
pulsarClientTlsTrustStorePath 字串 無 TLS 驗證的 TrustStore 檔案路徑。
pulsarClientTlsTrustStorePassword 字串 無 TLS 驗證的 TrustStore 密碼。

這些參數用於配置和驗證 Pulsar 准入控制,僅在啟用准入控制時(即設定 maxBytesPerTrigger 時)才需要進行 Pulsar 管理員配置。

選項 類型 預設 描述
觸發器最大位元組數 (maxBytesPerTrigger) BIGINT 無 我們希望處理的每個微批次最大位元組數的軟性限制。 如果指定這個,也必須指定 admin.url。
adminUrl 字串 無 Pulsar 服務的 HttpUrl 配置。 只有在指定 maxBytesPerTrigger 時才需要。
pulsarAdminAuthPlugin 字串 無 驗證外掛程式的名稱。
pulsar 管理員認證參數 字串 無 驗證外掛程式的參數。
pulsarClientUseKeyStoreTls 字串 無 是否要使用 KeyStore 進行 Tls 驗證。
pulsarAdminTlsTrustStoreType 字串 無 TLS 驗證的 TrustStore 檔案類型。
pulsarAdminTlsTrustStorePath 字串 無 TLS 驗證的 TrustStore 檔案路徑。
pulsarAdminTlsTrustStorePassword 字串 無 TLS 驗證的 TrustStore 密碼。

退貨

具有下列架構的脈衝星記錄數據表。

  • __key STRING NOT NULL: Pulsar 消息鍵。

  • value BINARY NOT NULL: Pulsar 訊息值。

    注意:針對包含 Avro 或 JSON 架構的主題,內容不會被載入為二進位值,而是會展開以保留 Pulsar 主題的欄位名稱和欄位類型。

  • __topic STRING NOT NULL: Pulsar 主題名稱。

  • __messageId BINARY NOT NULL: Pulsar 訊息ID。

  • __publishTime TIMESTAMP NOT NULL: Pulsar 訊息發佈時間。

  • __eventTime TIMESTAMP NOT NULL: Pulsar 訊息事件時間。

  • __messageProperties MAP<STRING, STRING>: Pulsar 訊息屬性。

範例

-- Streaming from Pulsar
> CREATE STREAMING TABLE testing.streaming_table AS
  SELECT * FROM STREAM read_pulsar(
      serviceUrl => 'pulsar://broker.example.com:6650',
      startingOffsets => 'earliest',
      topic => 'my-topic');

-- Streaming Ingestion from Pulsar with authentication
> CREATE STREAMING TABLE testing.streaming_table AS
  SELECT * FROM STREAM read_pulsar(
        serviceUrl => 'pulsar://broker.example.com:6650',
        startingOffsets => 'earliest',
        topic => 'my-topic',
        pulsarClientAuthPluginClassName => 'org.apache.pulsar.client.impl.auth.AuthenticationKeyStoreTls',
        pulsarClientAuthParams => 'keyStoreType:JKS,keyStorePath:/var/private/tls/client.keystore.jks,keyStorePassword:clientpw'
        );

The data can now to be queried from the testing.streaming_table for further analysis.