Проверка подлинности

На этой странице показаны наиболее распространенные методы проверки подлинности соединителя Kafka на Azure Databricks.

Полный список поддерживаемых методов проверки подлинности можно найти в документации Kafka. Справочные сведения о параметрах проверки подлинности см. в разделе "Проверка подлинности".

Connect to Центры событий Azure с субъектом-службой

Azure Databricks поддерживает аутентификацию заданий Spark для служб Центров событий с использованием OAuth и Microsoft Entra ID.

Схема проверки подлинности AAD

Подключение по учетным данным службы каталога Unity

В Databricks Runtime 16.1 и более поздних версиях Azure Databricks поддерживает учетные данные службы Unity Catalog для аутентификации в Центры событий Azure. Databricks рекомендует этот подход, если вы запускаете потоковую передачу Kafka в общих кластерах или бессерверных вычислениях.

Чтобы использовать учетные данные службы каталога Unity для проверки подлинности, выполните следующие действия.

  • Создайте учетные данные службы каталога 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 для подключения к Kafka не используйте следующие параметры:

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

Подключитесь, используя идентификатор клиента и секретный ключ.

Azure Databricks поддерживает проверку подлинности Microsoft Entra ID с идентификатором клиента и секретом в следующих вычислительных средах:

  • Databricks Runtime 12.2 LTS и выше в средах вычислений, настроенных с использованием выделенного режима доступа.
  • Databricks Runtime 14.3 LTS и более поздних версий для вычислений, настроенных в стандартном режиме доступа.
  • Конвейеры Lakeflow, настроенные без каталога Unity.

Azure Databricks не поддерживает проверку подлинности Microsoft Entra ID с помощью сертификата в любой вычислительной среде или конвейерах Lakeflow, настроенных с помощью каталога Unity.

Эта аутентификация не работает для вычислительных ресурсов со стандартным режимом доступа или конвейеров Lakeflow в Unity Catalog.

Чтобы выполнить проверку подлинности с помощью Microsoft Entra ID, необходимо иметь следующие значения:

  • Идентификатор клиента. Это можно найти на вкладке служб Microsoft Entra ID.

  • Идентификатор клиента, также известный как идентификатор приложения.

  • Секрет клиента. Добавьте это как секрет в рабочее пространство Databricks. Дополнительные сведения см. в разделе Управление секретами.

  • Тема EventHubs. Список тем вы можете найти в разделе "Центры событий" в разделе "Сущности" на конкретной странице пространства имен "Центры событий". Чтобы работать с несколькими темами, можно задать роль IAM на уровне Event Hubs.

  • Сервер EventHubs. Это можно найти на странице общей информации вашего конкретного пространства имен для Центров событий.

    Пространство имен Центров событий

Чтобы использовать Entra ID, необходимо настроить Kafka для использования OAuth SASL:

  • Установите kafka.security.protocol на SASL_SSL
  • Установите kafka.sasl.mechanism на OAUTHBEARER
  • Укажите для kafka.sasl.login.callback.handler.class полностью определённое имя класса Java. Полное имя — kafkashaded и обработчик обратного вызова входа класса Kafka с затенением Databricks. См. следующий пример для данного класса.

SASL — это универсальный протокол проверки подлинности, а OAuth — это механизм SASL.

Следующий пример настраивает Kafka для подключения к Центры событий Azure с помощью проверки подлинности Microsoft Entra 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 для проверки подлинности

Чтобы подключиться к Kafka с помощью проверки подлинности SASL/PLAIN (имя пользователя и пароль), настройте следующие параметры. Используйте имя затеняемого 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 для проверки подлинности

Чтобы подключиться к Kafka с помощью SASL/SCRAM (SCRAM-SHA-256 или SCRAM-SHA-512), настройте следующие параметры. Используйте имя затеняемого 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>";'
);

Замечание

Замените SCRAM-SHA-512 на SCRAM-SHA-256, если кластер Kafka настроен на использование SCRAM-SHA-256.

Azure Databricks рекомендует хранить пароль как секрет, а не включать его непосредственно в код. Дополнительные сведения см. в разделе "Управление секретами".

Использовать SSL для подключения Azure Databricks к Kafka

Чтобы включить подключения SSL/TLS к Kafka, установите kafka.security.protocol на SSL и укажите параметры конфигурации хранилища доверия и хранилища ключей с префиксом kafka.. Для SSL-подключений, требующих только проверки подлинности сервера (односторонняя проверка подлинности TLS), необходимо использовать хранилище доверия. При использовании взаимной аутентификации TLS (mTLS), когда брокер Kafka также аутентифицирует клиента, необходимо использовать и хранилище доверенных сертификатов, и хранилище ключей.

Доступны следующие параметры SSL/TLS. Полный список свойств SSL см. в документации по конфигурации SSL Apache Kafkaи шифрованию и проверке подлинности с помощью SSL в документации confluent.

Опция Описание
kafka.security.protocol Установите значение SSL для включения шифрования TLS.
kafka.ssl.truststore.location Путь к файлу хранилища доверия, содержащему доверенные сертификаты Центра сертификации.
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>')
);

Подключение Kafka на HDInsight к Azure Databricks

  1. Создайте кластер Kafka HDInsight.

    Инструкции см. в статье Connect to Kafka в HDInsight с помощью Azure Virtual Network.

  2. Настройте брокеры Kafka для объявления правильного адреса.

    Следуйте инструкциям из раздела Настройка Kafka для рекламы IP-адресов. Если вы управляете Kafka самостоятельно на Виртуальные машины Azure, убедитесь, что конфигурация advertised.listeners брокеров настроена на внутренний IP-адрес узлов.

  3. Создайте кластер Azure Databricks.

  4. Настройте пиринговое подключение кластера Kafka к кластеру Azure Databricks.

    Следуйте инструкциям из раздела Виртуальные одноранговые сети.

Используйте имена классов Kafka с затенением Databricks

Azure Databricks пакеты собственных, затеняемых версий клиентских библиотек Kafka. Все имена клиентских классов Kafka, на которые вы ссылаетесь в параметрах конфигурации проверки подлинности, должны использовать префикс имени затенённого класса вместо стандартного имени класса с открытым исходным кодом. Это относится к любому классу, на который ссылается ссылка в таких параметрах, как kafka.sasl.jaas.config, kafka.sasl.login.callback.handler.classи kafka.sasl.client.callback.handler.class.

Если вы используете незамеченные имена классов, код вызывает ошибку RESTRICTED_STREAMING_OPTION_PERMISSION_ENFORCED . Дополнительные сведения см. в разделе часто задаваемых вопросов .

Обработка потенциальных ошибок

  • Не удалось создать новый KafkaAdminClient

    Эта внутренняя ошибка Kafka возникает, если какие-либо из следующих параметров проверки подлинности неверны:

    • Идентификатор клиента (также известный как идентификатор приложения)
    • Идентификатор арендатора
    • Сервер Центров событий

    Чтобы устранить ошибку, убедитесь, что значения верны для этих параметров. Кроме того, эта ошибка может появиться при изменении параметров конфигурации, предоставленных по умолчанию в примере (например kafka.security.protocol, ).

  • Записи не найдены

    Если вы пытаетесь отобразить или обработать кадр данных, но не получаете результаты, вы увидите следующее в пользовательском интерфейсе.

    Нет сообщения о результатах

    Это сообщение означает, что проверка подлинности прошла успешно, но EventHubs не возвращала никаких данных. Некоторые возможные причины (хотя и не являются исчерпывающими) являются:

    • Вы указали неправильный раздел EventHubs .
    • Для параметра конфигурации Kafka startingOffsets по умолчанию используется latest, и вы пока не получаете никаких данных через тему. Вы можете установить startingOffsets, чтобы начать чтение данных с earliest, начиная с самых ранних смещений Kafka.