Aller au contenu principal

Authentification

Cette page présente les méthodes d'authentification les plus courantes pour le connecteur Kafka sur Databricks.

La liste complète des méthodes d’authentification prises en charge se trouve dans la documentation Kafka. Pour la référence des options d’authentification, consultez Authentification.

Se connecter au service géré Google Cloud pour Apache Kafka

Pour vous authentifier auprès de Google Cloud Managed Service for Apache Kafka, utilisez un identifiant de service Unity Catalog. Vous pouvez soit laisser Databricks configurer l'authentification automatiquement, soit configurer manuellement les options SASL et OAUTHBEARER.

Se connecter avec les identifiants de service Unity Catalog

Dans Databricks Runtime 18,0 et versions ultérieures, Databricks prend en charge les identifiants de service Unity Catalog pour authentifier l'accès au service géré Google Cloud pour Apache Kafka.

Les clusters Google Cloud Managed Service for Apache Kafka s'exécutent dans un Virtual Private Cloud (VPC). Votre compute Databricks doit s'exécuter dans un Virtual Private Cloud (VPC) qui peut atteindre le cluster Kafka géré. Pour activer la connectivité, configurez le peering de Virtual Private Cloud (VPC) ou Private Service Connect. Consultez Virtual Private Cloud (VPC) Network Peering ou Private Service Connect dans la documentation Google Cloud. Le compute Serverless n’est pas pris en charge, car il s’exécute dans le plan de contrôle géré par Databricks.

Pour l'authentification, accordez au compte de service Google Cloud, créé pour votre identifiant de service Unity Catalog, le rôle suivant dans le projet Google Cloud qui contient votre cluster Kafka géré :

  • roles/managedkafka.client

Pour obtenir des instructions, consultez Octroi, modification et révocation de l'accès dans la documentation Google Cloud.

L'exemple suivant configure Kafka comme source à l'aide d'un identifiant 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://www.googleapis.com/auth/cloud-platform",
}

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

Lorsque vous utilisez un identifiant de service Unity Catalog pour vous connecter à Kafka, n'utilisez pas les options suivantes :

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

Se connecter avec les options SASL

Vous pouvez vous authentifier auprès de Google Cloud Managed Service pour Apache Kafka sans utiliser l'option de source databricks.serviceCredential en spécifiant les options SASL/OAUTHBEARER. N'utilisez cette approche que si vous devez définir explicitement les options SASL Kafka. Par exemple, pour utiliser un gestionnaire de rappel spécifique.

Pour se connecter avec SASL, vous devez utiliser un identifiant de service Unity Catalog :

  • Le compte de service des identifiants doit avoir roles/managedkafka.client dans le projet Google Cloud qui contient votre cluster Kafka géré.
  • Utilisez le gestionnaire de rappel de connexion fourni par Databricks org.apache.spark.sql.kafka010.CustomGCPOAuthBearerLoginCallbackHandler. N'utilisez pas com.google.cloud.hosted.kafka.auth.GcpLoginCallbackHandler, qui n'est pas inclus dans Databricks Runtime.

L'authentification avec les comptes de service Google sans Unity Catalog n'est pas prise en charge.

L'exemple suivant configure Kafka comme source avec SASL :

Python
kafka_options = {
"kafka.bootstrap.servers": "<bootstrap-hostname>:9092",
"subscribe": "<topic>",
"startingOffsets": "earliest",
"kafka.security.protocol": "SASL_SSL",
"kafka.sasl.mechanism": "OAUTHBEARER",
"kafka.sasl.jaas.config":
"kafkashaded.org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginModule required;",
"kafka.sasl.login.callback.handler.class":
"org.apache.spark.sql.kafka010.CustomGCPOAuthBearerLoginCallbackHandler",
# Required by the login callback handler:
"kafka.databricks.serviceCredential": "<service-credential-name>",
# Optional:
# "kafka.databricks.serviceCredential.scope": "https://www.googleapis.com/auth/cloud-platform",
}

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

Ne définissez pas databricks.serviceCredential lorsque vous spécifiez les options SASL. Si vous définissez databricks.serviceCredential, Databricks configure automatiquement l'authentification Kafka et interdit de spécifier les options kafka.sasl.*.

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 classe ombré 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()

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 en utilisant SASL/SCRAM (SCRAM-SHA-256 ou SCRAM-SHA-512), configurez les options suivantes. Utilisez le nom de classe ombré 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()
remarque

Remplacez SCRAM-SHA-512 par SCRAM-SHA-256 si votre cluster Kafka est configuré pour utiliser SCRAM-SHA-256.

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.

Utilisez SSL pour connecter Databricks à Kafka

Pour activer les connexions SSL/TLS à Kafka, définissez kafka.security.protocol sur SSL et fournissez les options de configuration du magasin de confiance et du magasin de clés préfixées par kafka.. Pour les connexions SSL qui ne nécessitent qu'une authentification du serveur (TLS unidirectionnel), vous devez utiliser un magasin de confiance. Pour le TLS mutuel (mTLS) où le broker Kafka authentifie également le client, vous devez utiliser à la fois un magasin de confiance et un magasin de clés.

Les options SSL/TLS suivantes sont disponibles. Pour la liste complète des propriétés SSL, consultez la documentation de configuration SSL d'Apache Kafka et la documentation sur le chiffrement et l'authentification avec SSL dans la documentation Confluent.

Option

Description

kafka.security.protocol

Réglez sur SSL pour activer le chiffrement TLS.

kafka.ssl.truststore.location

Chemin vers le fichier du magasin de confiance contenant les certificats d'autorité de certification (CA) de confiance.

kafka.ssl.truststore.password

Mot de passe du fichier de clés de confiance.

kafka.ssl.truststore.type

Format de fichier du magasin de confiance (default : 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 du fichier de stockage 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 Hostname. La valeur par défaut est https. Définir sur une chaîne vide pour désactiver.

Option

Description

kafka.security.protocol

Réglez sur SSL pour activer le chiffrement TLS.

kafka.ssl.truststore.location

Chemin vers le fichier du magasin de confiance contenant les certificats d'autorité de certification (CA) de confiance.

kafka.ssl.truststore.password

Mot de passe du fichier de clés de confiance.

kafka.ssl.truststore.type

Format de fichier du magasin de confiance (default : 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 du fichier de stockage 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 Hostname. La valeur par défaut est https. Définir sur une chaîne vide pour désactiver.

Si vous utilisez SSL, Databricks vous recommande de :

  • Stockez vos certificats dans un volume Unity Catalog. Les utilisateurs qui ont accès en lecture au volume peuvent utiliser vos certificats Kafka. Pour plus d’informations, consultez Que sont les volumes Unity Catalog ?.
  • Stockez vos mots de passe de certificat en tant que secrets dans un Secret Scope. Pour plus d’informations, consultez Manage secret scopes.

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>))
)

Utiliser les noms de classe Kafka obscurcis de Databricks

Databricks regroupe des versions propriétaires et « shaded » des bibliothèques clientes Kafka. Tous les noms de classes clients Kafka que vous référencez dans les options de configuration d'authentification doivent utiliser le préfixe de nom de classe masquée au lieu du nom de classe open source standard. Ceci s'applique à toute classe référencée dans des options comme kafka.sasl.jaas.config, kafka.sasl.login.callback.handler.class et kafka.sasl.client.callback.handler.class.

Si vous utilisez des noms de classe non ombrés, votre code génère une erreur RESTRICTED_STREAMING_OPTION_PERMISSION_ENFORCED. Consultez la FAQ pour plus de détails.

Gestion des erreurs potentielles

  • Exception missing gcp_options lors de l'authentification

    Cette exception est levée si le gestionnaire de rappel ne peut pas déduire la portée de l'URL d'amorçage. Définissez databricks.serviceCredential.scope manuellement. Pour la plupart des cas d'usage, définissez-le sur https://www.googleapis.com/auth/cloud-platform.

  • Aucune URL d’amorçage résolvable trouvée

    Cela signifie que le cluster de compute est incapable de résoudre le hostname d’amorçage. Vérifiez la configuration de votre Virtual Private Cloud (VPC) pour vous assurer que le cluster compute peut atteindre le cluster Kafka géré. Configurez le peering Virtual Private Cloud (VPC) ou Private Service Connect si nécessaire.

  • Problèmes d'autorisations

    • Assurez-vous que le compte de service dispose de la liaison de rôle IAM roles/managedkafka.client sur le projet auquel le cluster Kafka appartient.
    • Assurez-vous que le cluster Kafka dispose des listes de contrôle d'accès (ACL) appropriées définies pour les comptes de rubrique et de service. Voir le contrôle d'accès avec les ACL Kafka dans la documentation Google Cloud.
  • Aucun enregistrement renvoyé

    Si l'authentification réussit mais qu'aucune donnée n'est renvoyée :

    • Vérifiez que vous êtes abonné au nom de rubrique correct.
    • Le default startingOffsets est latest, qui ne lit que les nouvelles données. Définissez startingOffsets sur earliest pour lire les données existantes.