驗證

本頁展示了 Azure Databricks 上 Kafka 連接器最常見的認證方法。

完整的支援認證方法清單可在 Kafka 文件中找到。 關於認證選項的參考,請參見認證。

透過服務主體連接Azure 事件中樞

Azure Databricks 支援使用 OAuth 搭配 Microsoft Entra ID 進行 Event Hubs 服務的 Spark 工作認證。

AAD 驗證圖表

連結 Unity Catalog 服務存取憑證

在 Databricks Runtime 16.1 及以上版本中,Azure Databricks 支援 Unity 目錄服務憑證以進行Azure 事件中樞驗證。 如果您是在共用叢集或無伺服器運算環境上執行 Kafka 串流,Databricks 建議採用此做法。

要使用 Unity Catalog 服務憑證進行認證,請執行以下步驟:

  • 建立新的 Unity 目錄服務認證。 請參閱 建立服務認證。
    • 確認連接到你服務憑證的存取連接器是否擁有正確的權限,可以連接到 Azure 事件中樞。
  • 將來源選項 databricks.serviceCredential 設為你的服務憑證名稱。

以下範例是利用服務憑證將 Kafka 配置為來源:

Python

kafka_options = {
  "kafka.bootstrap.servers": "<bootstrap-hostname>:9092",
  "subscribe": "<topic>",
  "databricks.serviceCredential": "<service-credential-name>",
  # Optional: set this only if Databricks can't infer the scope for your Kafka service.
  # "databricks.serviceCredential.scope": "https://<event-hubs-server>/.default",
}

df = spark.readStream.format("kafka").options(**kafka_options).load()

Scala

val kafkaOptions = Map(
  "kafka.bootstrap.servers" -> "<bootstrap-hostname>:9092",
  "subscribe" -> "<topic>",
  "databricks.serviceCredential" -> "<service-credential-name>",
  // Optional: set this only if Databricks can't infer the scope for your Kafka service.
  // "databricks.serviceCredential.scope" -> "https://<event-hubs-server>/.default",
)

val df = spark.readStream.format("kafka").options(kafkaOptions).load()

SQL

SELECT * FROM read_kafka(
  bootstrapServers => '<bootstrap-hostname>:9092',
  subscribe => '<topic>',
  serviceCredential => '<service-credential-name>'
);

備註

當你使用 Unity Catalog 服務憑證連接 Kafka 時,請不要使用以下選項:

  • kafka.sasl.mechanism
  • kafka.sasl.jaas.config
  • kafka.security.protocol
  • kafka.sasl.client.callback.handler.class
  • kafka.sasl.oauthbearer.token.endpoint.url

透過客戶端 ID 和密鑰進行連接

Azure Databricks 支援以下運算環境中以客戶端 ID 與秘密密碼進行 Microsoft Entra ID 認證:

  • 在設定為專用存取模式的運算資源上執行的 Databricks Runtime 12.2 LTS 及以上版本。
  • Databricks Runtime 14.3 LTS 及以上版本,計算時設定為標準存取模式。
  • Lakeflow 管線設定時未使用 Unity Catalog。

Azure Databricks 在任何計算環境或以 Unity Catalog 設定的 Lakeflow pipelines 中,都不支援 Microsoft Entra ID 憑證認證。

此認證在標準存取模式的運算或 Unity Catalog Lakeflow 管線中無法運作。

要使用 Microsoft Entra ID 進行驗證,您必須具備以下數值:

  • 租戶識別碼。 你可以在Microsoft Entra ID服務標籤中找到這些資料。

  • 客戶識別碼(clientID),也稱為應用程式識別碼(Application ID)。

  • 用戶端密碼。 將此作為秘密加入你的 Databricks 工作空間。 請參閱機密管理。

  • EventHubs 的主題。 您可以在特定 事件中樞命名空間 頁面的 [實體] 區段底下,找到 [事件中樞] 區段的主題清單。 若要處理多個主題,您可以在 Event Hubs 層級設定 IAM 角色。

  • EventHubs 伺服器。 您可以在特定 [事件中樞命名空間] 的概觀頁面上找到此項:

    事件中樞命名空間

要使用 Entra ID,您必須設定 Kafka 以使用 OAuth SASL:

  • 將 kafka.security.protocol 設定為 SASL_SSL
  • 將 kafka.sasl.mechanism 設定為 OAUTHBEARER
  • 將 kafka.sasl.login.callback.handler.class 設為Java類別的完全合格名稱。 限定名稱為 kafkashaded,也是 Databricks shaded Kafka 類別的登入回呼處理常式。 如需確切的類別,請參閱下列範例。

SASL 是一種通用認證協定,而 OAuth 則是 SASL 機制。

以下範例設定 Kafka 以 Microsoft Entra ID 認證連接 Azure 事件中樞,並附有客戶端 ID 與秘密:

Python

# This is the only section you need to modify for auth purposes
# ------------------------------
tenant_id = "..."
client_id = "..."
client_secret = dbutils.secrets.get("your-scope", "your-secret-name")

event_hubs_server = "..."
event_hubs_topic = "..."
# -------------------------------

sasl_config = f'kafkashaded.org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginModule required clientId="{client_id}" clientSecret="{client_secret}" scope="https://{event_hubs_server}/.default" ssl.protocol="SSL";'

kafka_options = {
    "kafka.bootstrap.servers": f"{event_hubs_server}:9093", # Port 9093 is the EventHubs Kafka port
    "kafka.sasl.jaas.config": sasl_config,
    "kafka.sasl.oauthbearer.token.endpoint.url": f"https://login.microsoft.com/{tenant_id}/oauth2/v2.0/token",
    "subscribe": event_hubs_topic,

    # You should not need to modify these
    "kafka.security.protocol": "SASL_SSL",
    "kafka.sasl.mechanism": "OAUTHBEARER",
    "kafka.sasl.login.callback.handler.class": "kafkashaded.org.apache.kafka.common.security.oauthbearer.secured.OAuthBearerLoginCallbackHandler"
}

df = spark.readStream.format("kafka").options(**kafka_options)

display(df)

Scala

// This is the only section you need to modify for auth purposes
// -------------------------------
val tenantId = "..."
val clientId = "..."
val clientSecret = dbutils.secrets.get("your-scope", "your-secret-name")

val eventHubsServer = "..."
val eventHubsTopic = "..."
// -------------------------------

val saslConfig = s"""kafkashaded.org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginModule required clientId="$clientId" clientSecret="$clientSecret" scope="https://$eventHubsServer/.default" ssl.protocol="SSL";"""

val kafkaOptions = Map(
  "kafka.bootstrap.servers" -> s"$eventHubsServer:9093", // Port 9093 is the EventHubs Kafka port
  "kafka.sasl.jaas.config" -> saslConfig,
  "kafka.sasl.oauthbearer.token.endpoint.url" -> s"https://login.microsoft.com/$tenantId/oauth2/v2.0/token",
  "subscribe" -> eventHubsTopic,

  // You should not need to modify these
  "kafka.security.protocol" -> "SASL_SSL",
  "kafka.sasl.mechanism" -> "OAUTHBEARER",
  "kafka.sasl.login.callback.handler.class" -> "kafkashaded.org.apache.kafka.common.security.oauthbearer.secured.OAuthBearerLoginCallbackHandler"
)

val scalaDF = spark.readStream
  .format("kafka")
  .options(kafkaOptions)
  .load()

display(scalaDF)

SQL

CREATE OR REFRESH STREAMING TABLE <table_name>
AS
SELECT * FROM STREAM read_kafka(
  bootstrapServers => '<event-hubs-server>:9093',
  subscribe => '<event-hubs-topic>',
  `kafka.security.protocol` => 'SASL_SSL',
  `kafka.sasl.mechanism` => 'OAUTHBEARER',
  `kafka.sasl.jaas.config` => 'kafkashaded.org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginModule required clientId="<client-id>" clientSecret="<client-secret>" scope="https://<event-hubs-server>/.default" ssl.protocol="SSL";',
  `kafka.sasl.oauthbearer.token.endpoint.url` => 'https://login.microsoft.com/<tenant-id>/oauth2/v2.0/token',
  `kafka.sasl.login.callback.handler.class` => 'kafkashaded.org.apache.kafka.common.security.oauthbearer.secured.OAuthBearerLoginCallbackHandler'
);

使用 SASL/PLAIN 來認證

若要使用 SASL/PLAIN(使用者名稱與密碼)認證連接 Kafka,請設定以下選項。 使用陰影 PlainLoginModule 類別名稱:

Python

kafka_options = {
  "kafka.bootstrap.servers": "<bootstrap-server>:9093",
  "subscribe": "<topic>",
  "kafka.security.protocol": "SASL_SSL",
  "kafka.sasl.mechanism": "PLAIN",
  "kafka.sasl.jaas.config":
    'kafkashaded.org.apache.kafka.common.security.plain.PlainLoginModule required username="<username>" password="<password>";',
}

df = spark.readStream.format("kafka").options(**kafka_options).load()

Scala

val kafkaOptions = Map(
  "kafka.bootstrap.servers" -> "<bootstrap-server>:9093",
  "subscribe" -> "<topic>",
  "kafka.security.protocol" -> "SASL_SSL",
  "kafka.sasl.mechanism" -> "PLAIN",
  "kafka.sasl.jaas.config" ->
    """kafkashaded.org.apache.kafka.common.security.plain.PlainLoginModule required username="<username>" password="<password>";""",
)

val df = spark.readStream.format("kafka").options(kafkaOptions).load()

SQL

SELECT * FROM STREAM read_kafka(
  bootstrapServers => '<bootstrap-server>:9093',
  subscribe => '<topic>',
  `kafka.security.protocol` => 'SASL_SSL',
  `kafka.sasl.mechanism` => 'PLAIN',
  `kafka.sasl.jaas.config` => 'kafkashaded.org.apache.kafka.common.security.plain.PlainLoginModule required username="<username>" password="<password>";'
);

Azure Databricks 建議你把密碼當作秘密,而不是直接寫在程式碼裡。 如需詳細資訊,請參閱 秘密管理。

使用 SASL/SCRAM 來驗證

若要使用 SASL/SCRAM(SCRAM-SHA-256 或 SCRAM-SHA-512)連接 Kafka,請設定以下選項。 使用陰影 ScramLoginModule 類別名稱:

Python

kafka_options = {
  "kafka.bootstrap.servers": "<bootstrap-server>:9093",
  "subscribe": "<topic>",
  "kafka.security.protocol": "SASL_SSL",
  "kafka.sasl.mechanism": "SCRAM-SHA-512",
  "kafka.sasl.jaas.config":
    'kafkashaded.org.apache.kafka.common.security.scram.ScramLoginModule required username="<username>" password="<password>";',
}

df = spark.readStream.format("kafka").options(**kafka_options).load()

Scala

val kafkaOptions = Map(
  "kafka.bootstrap.servers" -> "<bootstrap-server>:9093",
  "subscribe" -> "<topic>",
  "kafka.security.protocol" -> "SASL_SSL",
  "kafka.sasl.mechanism" -> "SCRAM-SHA-512",
  "kafka.sasl.jaas.config" ->
    """kafkashaded.org.apache.kafka.common.security.scram.ScramLoginModule required username="<username>" password="<password>";""",
)

val df = spark.readStream.format("kafka").options(kafkaOptions).load()

SQL

SELECT * FROM STREAM read_kafka(
  bootstrapServers => '<bootstrap-server>:9093',
  subscribe => '<topic>',
  `kafka.security.protocol` => 'SASL_SSL',
  `kafka.sasl.mechanism` => 'SCRAM-SHA-512',
  `kafka.sasl.jaas.config` => 'kafkashaded.org.apache.kafka.common.security.scram.ScramLoginModule required username="<username>" password="<password>";'
);

備註

如果你的 Kafka 叢集設定為使用 SCRAM-SHA-256,則替換SCRAM-SHA-512為SCRAM-SHA-256

Azure Databricks 建議你把密碼當作秘密,而不是直接寫在程式碼裡。 如需詳細資訊,請參閱 秘密管理。

使用 SSL 連接 Azure Databricks 到 Kafka

要啟用 Kafka 的 SSL/TLS 連線,請設定kafka.security.protocolSSL並提供 信任儲存與金鑰儲存的設定選項,前綴為 kafka.。 對於只需要伺服器驗證(單向 TLS)的 SSL 連線,您必須使用信任存放區。 對於 Kafka broker 也會驗證用戶端身分的雙向 TLS(mTLS),你必須同時使用信任存放區和金鑰存放區。

以下 SSL/TLS 選項可用。 完整 SSL 屬性列表,請參閱 Apache Kafka SSL 設定文件 及 Confluent 文件中的 SSL 加密與認證 。

Option 說明
kafka.security.protocol 設定為 SSL 啟用 TLS 加密。
kafka.ssl.truststore.location 通往包含受信任 CA 憑證的信任憑證庫檔案的路徑。
kafka.ssl.truststore.password 信任庫檔案的密碼。
kafka.ssl.truststore.type 信任儲存檔案格式(預設: JKS)。
kafka.ssl.keystore.location 通往包含用戶端憑證與私鑰(mTLS 必備)的金鑰儲存檔案路徑。
kafka.ssl.keystore.password 金鑰儲存檔的密碼。
kafka.ssl.key.password 用於金鑰庫中私鑰的密碼。
kafka.ssl.endpoint.identification.algorithm 主機名稱驗證演算法。 預設為 https。 設定為空字串以停用。

如果你使用 SSL,Databricks 建議你:

  • 把你的憑證存放在 Unity 目錄卷中。 有權讀取磁碟卷的使用者可以使用你的 Kafka 憑證。 如需詳細資訊,請參閱 什麼是 Unity 目錄磁碟區?。
  • 把你的憑證密碼當作秘密儲存在秘密範圍內。 如需詳細資訊,請參閱 管理秘密範圍。

下列範例使用物件儲存位置和 Databricks 祕密來啟用 SSL 連線:

Python

df = (spark.readStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "<bootstrap-server>:9093")
  .option("kafka.security.protocol", "SSL")
  .option("kafka.ssl.truststore.location", <truststore-location>)
  .option("kafka.ssl.keystore.location", <keystore-location>)
  .option("kafka.ssl.keystore.password", dbutils.secrets.get(scope=<certificate-scope-name>,key=<keystore-password-key-name>))
  .option("kafka.ssl.truststore.password", dbutils.secrets.get(scope=<certificate-scope-name>,key=<truststore-password-key-name>))
)

Scala

val df = spark.readStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "<bootstrap-server>:9093")
  .option("kafka.security.protocol", "SSL")
  .option("kafka.ssl.truststore.location", <truststore-location>)
  .option("kafka.ssl.keystore.location", <keystore-location>)
  .option("kafka.ssl.keystore.password", dbutils.secrets.get(scope = <certificate-scope-name>, key = <keystore-password-key-name>))
  .option("kafka.ssl.truststore.password", dbutils.secrets.get(scope = <certificate-scope-name>, key = <truststore-password-key-name>))

SQL

SELECT * FROM read_kafka(
  bootstrapServers => '<bootstrap-server>:9093',
  subscribe => '<topic>',
  `kafka.security.protocol` => 'SSL',
  `kafka.ssl.truststore.location` => '<truststore-location>',
  `kafka.ssl.keystore.location` => '<keystore-location>',
  `kafka.ssl.keystore.password` => secret('<certificate-scope-name>', '<keystore-password-key-name>'),
  `kafka.ssl.truststore.password` => secret('<certificate-scope-name>', '<truststore-password-key-name>')
);

將 HDInsight 上的 Kafka 連接到 Azure Databricks

  1. 建立 HDInsight Kafka 叢集。

    請參閱 透過 Azure 虛擬網路 連線到 HDInsight 上的 Kafka 以獲取說明。

  2. 設定 Kafka 訊息代理程式來公告正確的位址。

    請遵循設定 Kafka 進行 IP 公告中的指示。 如果你自己在 Azure 虛擬機器 上管理 Kafka,請確保代理的 advertised.listeners 設定設為主機的內部 IP。

  3. 建立 Azure Databricks 叢集。

  4. 將 Kafka 叢集對等連線至 Azure Databricks 叢集。

    請遵循對等互連虛擬網路中的指示。

使用 Databricks 著色的 Kafka 類別名稱

Azure Databricks 包含專有的、特殊處理版本的 Kafka 用戶端程式庫。 你在認證設定選項中引用的所有 Kafka 用戶端類別名稱,必須使用著色類別名稱前綴,而非標準的開源類別名稱。 這適用於任何在選項中引用的類別,如 kafka.sasl.jaas.config、 kafka.sasl.login.callback.handler.class和 kafka.sasl.client.callback.handler.class。

如果您使用未經 shade 處理的類別名稱,您的程式碼會引發 RESTRICTED_STREAMING_OPTION_PERMISSION_ENFORCED 錯誤。 詳情請參閱 常見問題 集。

處理潛在錯誤

  • 無法建立新的 KafkaAdminClient

    此內部 Kafka 錯誤會拋出,若以下任何認證選項錯誤:

    • 用戶端識別碼 (也稱為應用程式識別碼)
    • 租戶識別碼
    • Event Hubs 伺服器

    若要解決錯誤,請確認這些選項的值正確無誤。 此外,如果您修改範例中預設提供的設定選項(例如 kafka.security.protocol),可能會出現此錯誤。

  • 沒有回傳紀錄

    如果您嘗試顯示或處理 DataFrame 但未取得結果,您會在 UI 中看到下列內容。

    無結果訊息

    此訊息表示驗證成功,但 EventHubs 未傳回任何數據。 一些可能的 (儘管並不詳盡) 原因包括:

    • 指定的 EventHubs 主題錯誤。
    • Kafka 的預設組態選項 startingOffsets 是 latest,而且您目前尚未透過主题接收任何資料。 您可以將 設定 startingOffsets 為 earliest 以從 Kafka 最早位移開始讀取數據。