Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
На этой странице показаны наиболее распространенные методы проверки подлинности соединителя Kafka на Azure Databricks.
Полный список поддерживаемых методов проверки подлинности можно найти в документации Kafka. Справочные сведения о параметрах проверки подлинности см. в разделе "Проверка подлинности".
Connect to Центры событий Azure с субъектом-службой
Azure Databricks поддерживает аутентификацию заданий Spark для служб Центров событий с использованием OAuth и Microsoft Entra ID.
Подключение по учетным данным службы каталога 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.mechanismkafka.sasl.jaas.configkafka.security.protocolkafka.sasl.client.callback.handler.classkafka.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
Создайте кластер Kafka HDInsight.
Инструкции см. в статье Connect to Kafka в HDInsight с помощью Azure Virtual Network.
Настройте брокеры Kafka для объявления правильного адреса.
Следуйте инструкциям из раздела Настройка Kafka для рекламы IP-адресов. Если вы управляете Kafka самостоятельно на Виртуальные машины Azure, убедитесь, что конфигурация
advertised.listenersброкеров настроена на внутренний IP-адрес узлов.Создайте кластер Azure Databricks.
Настройте пиринговое подключение кластера 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.