Remarque
L’accès à cette page requiert une autorisation. Vous pouvez essayer de vous connecter ou de modifier des répertoires.
L’accès à cette page requiert une autorisation. Vous pouvez essayer de modifier des répertoires.
Cette page présente les méthodes d’authentification les plus courantes pour le connecteur Kafka sur Azure Databricks.
Vous trouverez la liste complète des méthodes d’authentification prises en charge dans la documentation Kafka. Pour obtenir la référence des options d’authentification, consultez Authentification.
Se connecter à Azure Event Hubs avec un principal de service
Azure Databricks prend en charge l’authentification des travaux Spark avec les services Event Hubs à l’aide d’OAuth avec Microsoft Entra ID.
Se connecter avec les informations d’identification du service catalogue Unity
Dans Databricks Runtime 16.1 et versions ultérieures, Azure Databricks prend en charge les informations d’identification du service catalogue Unity pour l’authentification auprès de Azure Event Hubs. Databricks recommande cette approche si vous exécutez du streaming Kafka sur des clusters partagés ou avec du calcul sans serveur.
Pour utiliser les informations d’identification d’un service catalogue Unity pour l’authentification, procédez comme suit :
- Créez un identifiant de service pour Unity Catalog. Consultez Créer des informations d’identification de service.
- Vérifiez que le connecteur d’accès attaché à vos informations d’identification de service dispose des autorisations appropriées pour se connecter à Azure Event Hubs.
- Définissez l’option source
databricks.serviceCredentialsur le nom de votre identifiant de service.
L’exemple suivant configure Kafka en tant que source à l’aide d’informations d’identification de service :
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>'
);
Note
Lorsque vous utilisez des informations d’identification du service catalogue Unity pour vous connecter à Kafka, n’utilisez pas les options suivantes :
kafka.sasl.mechanismkafka.sasl.jaas.configkafka.security.protocolkafka.sasl.client.callback.handler.classkafka.sasl.oauthbearer.token.endpoint.url
Se connecter avec un ID client et un secret
Azure Databricks prend en charge l’authentification Microsoft Entra ID avec un ID client et un secret dans les environnements de calcul suivants :
- Databricks Runtime 12.2 LTS et versions ultérieures sur des ressources de calcul configurées en mode d’accès dédié.
- Databricks Runtime 14.3 LTS et versions ultérieures sur les ressources de calcul configurées avec le mode d’accès standard.
- Pipelines Lakeflow configurés sans catalogue Unity.
Azure Databricks ne prend pas en charge l’authentification Microsoft Entra ID avec un certificat dans n’importe quel environnement de calcul ou dans les pipelines Lakeflow configurés avec le catalogue Unity.
Cette authentification ne fonctionne pas sur le calcul avec le mode d’accès standard ou sur les pipelines Unity Catalog Lakeflow.
Pour effectuer l’authentification avec Microsoft Entra ID, vous devez avoir les valeurs suivantes :
ID de locataire. Vous pouvez le trouver sous l’onglet services Microsoft Entra ID.
Id client, également appelé ID d’application.
Une clé secrète client. Ajoutez-le comme secret à votre espace de travail Databricks. Consultez Gestion des secrets.
Une rubrique EventHubs. Vous pouvez trouver la liste des rubriques dans la section Event Hubs sous la section Entités, sur une page spécifique d’un Espace de noms Event Hubs. Pour utiliser plusieurs rubriques, vous pouvez définir un rôle IAM au niveau d’Event Hubs.
Un serveur EventHubs. Vous pouvez le trouver sur la page de présentation de votre espace de noms Event Hubs spécifique :
Pour utiliser Entra ID, vous devez configurer Kafka pour utiliser la SAPL OAuth :
- Paramétrez
kafka.security.protocolsurSASL_SSL - Paramétrez
kafka.sasl.mechanismsurOAUTHBEARER - Définissez
kafka.sasl.login.callback.handler.classcomme nom complet de la classe Java. Le nom qualifié estkafkashadedet le gestionnaire de rappel de connexion de la classe Kafka ombrée de Databricks. Consultez l’exemple suivant pour connaître la classe exacte.
SASL est un protocole d’authentification générique et OAuth est un mécanisme SASL.
L’exemple suivant configure Kafka pour se connecter à Azure Event Hubs à l’aide de l’authentification Microsoft Entra ID avec un ID client et un secret :
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'
);
Utiliser SASL/PLAIN pour s’authentifier
Pour vous connecter à Kafka à l’aide de l’authentification SASL/PLAIN (nom d’utilisateur et mot de passe), configurez les options suivantes. Utilisez le nom de la classe ombrée 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 vous recommande de stocker votre mot de passe en tant que secret plutôt que de l’inclure directement dans votre code. Pour plus d’informations, consultez Gestion des secrets.
Utiliser SASL/SCRAM pour s’authentifier
Pour vous connecter à Kafka à l’aide de SASL/SCRAM (SCRAM-SHA-256 ou SCRAM-SHA-512), configurez les options suivantes. Utilisez le nom de la classe ombrée 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>";'
);
Note
Remplacez par SCRAM-SHA-512 le cas SCRAM-SHA-256 où votre cluster Kafka est configuré pour utiliser SCRAM-SHA-256.
Azure Databricks vous recommande de stocker votre mot de passe en tant que secret plutôt que de l’inclure directement dans votre code. Pour plus d’informations, consultez Gestion des secrets.
Use SSL pour connecter Azure Databricks à Kafka
Pour activer les connexions SSL/TLS à Kafka, définissez kafka.security.protocol sur SSL et fournissez les options de configuration du truststore et du keystore précédées de kafka.. Pour les connexions SSL qui nécessitent uniquement l’authentification du serveur (TLS unidirectionnel), vous devez utiliser un magasin d’approbations. Pour le TLS mutuel (mTLS), dans lequel le broker Kafka authentifie également le client, vous devez utiliser à la fois un magasin de certificats de confiance et un magasin de clés.
Les options SSL/TLS suivantes sont disponibles. Pour obtenir la liste complète des propriétés SSL, consultez la documentation de configuration d’Apache Kafka SSL et le chiffrement et l’authentification avec SSL dans la documentation Confluent.
| Choix | Description |
|---|---|
kafka.security.protocol |
Définissez sur SSL pour activer le chiffrement TLS. |
kafka.ssl.truststore.location |
Chemin vers le fichier du magasin de confiance contenant les certificats d’autorités de certification (CA) approuvées. |
kafka.ssl.truststore.password |
Mot de passe pour le fichier du magasin de confiance. |
kafka.ssl.truststore.type |
Format de fichier du magasin d’approbations (par défaut : JKS). |
kafka.ssl.keystore.location |
Chemin d’accès au fichier de magasin de clés contenant le certificat client et la clé privée (requis pour mTLS). |
kafka.ssl.keystore.password |
Mot de passe pour le fichier du magasin de clés. |
kafka.ssl.key.password |
Mot de passe de la clé privée dans le magasin de clés. |
kafka.ssl.endpoint.identification.algorithm |
Algorithme de vérification du nom d’hôte. La valeur par défaut est https. Définissez une chaîne vide pour désactiver. |
Si vous utilisez SSL, Databricks vous recommande :
- Stockez vos certificats dans un volume de catalogue Unity. Les utilisateurs qui ont accès à la lecture à partir du volume peuvent utiliser vos certificats Kafka. Pour plus d'informations, consultez Qu'est-ce que les volumes Unity Catalog ?.
- Stockez vos mots de passe de certificat en tant que secrets dans un domaine secret. Pour plus d’informations, consultez Gérer les étendues de secrets.
L’exemple suivant utilise des emplacements de stockage d’objets et des secrets Databricks pour activer une connexion 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>')
);
Connecter Kafka sur HDInsight à Azure Databricks
Créez un cluster Kafka sur HDInsight.
Pour obtenir des instructions, consultez Se connecter à Kafka sur HDInsight via un réseau virtuel Azure.
Configurez les répartiteurs Kafka pour qu’ils publient l’adresse correcte.
Suivez les instructions fournies dans Configurer Kafka pour la publication d’adresses IP. Si vous gérez Kafka vous-même sur Machines virtuelles Azure, assurez-vous que la configuration
advertised.listenersdes répartiteurs est définie sur l’adresse IP interne des hôtes.Créez un cluster Azure Databricks.
Appairez le cluster Kafka au cluster Azure Databricks.
Suivez les instructions fournies dans Appairer des réseaux virtuels.
Utiliser les noms de classes Kafka ombrées de Databricks
Azure Databricks regroupe des versions propriétaires et ombrées des bibliothèques clientes Kafka. Tous les noms de classe client Kafka que vous référencez dans les options de configuration d’authentification doivent utiliser le préfixe de nom de classe ombré au lieu du nom de classe open source standard. Cela s’applique à n’importe quelle classe référencée dans les options telles que kafka.sasl.jaas.config, kafka.sasl.login.callback.handler.classet kafka.sasl.client.callback.handler.class.
Si vous utilisez des noms de classes sans ombrage, votre code génère une erreur RESTRICTED_STREAMING_OPTION_PERMISSION_ENFORCED. Pour plus d’informations, consultez le FAQ .
Gestion des erreurs potentielles
Échec de la création d’un nouveau
KafkaAdminClientCette erreur Kafka interne est levée si l’une des options d’authentification suivantes est incorrecte :
- ID client (également appelé ID d’application)
- ID du locataire
- Serveur Event Hubs
Pour résoudre cette erreur, vérifiez que les valeurs sont correctes pour ces options. En outre, vous pouvez voir cette erreur si vous modifiez les options de configuration fournies par défaut dans l’exemple (telles que
kafka.security.protocol).Aucun enregistrement retourné
Si vous essayez d’afficher ou de traiter votre DataFrame, mais que vous n’obtenez pas de résultats, vous verrez ce qui suit dans l’interface utilisateur.
Ce message signifie que l’authentification a réussi, mais EventHubs n’a retourné aucune donnée. Parmi les causes possibles (non exhaustives) :
- Vous avez spécifié une rubrique EventHubs incorrecte.
- L’option de configuration Kafka par défaut est
startingOffsetslatest, et vous ne recevez actuellement aucune donnée via la rubrique. Vous pouvez définirstartingOffsetssurearliestpour commencer à lire des données à partir des décalages les plus anciens de Kafka.