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
- Scala
- SQL
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()
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://www.googleapis.com/auth/cloud-platform",
)
val df = spark.read.format("kafka").options(kafkaOptions).load()
SELECT * FROM read_kafka(
bootstrapServers => '<bootstrap-hostname>:9092',
subscribe => '<topic>',
serviceCredential => '<service-credential-name>'
);
Lorsque vous utilisez un identifiant de service Unity Catalog pour vous connecter à Kafka, n'utilisez pas les options suivantes :
kafka.sasl.mechanismkafka.sasl.jaas.configkafka.security.protocolkafka.sasl.client.callback.handler.classkafka.sasl.login.callback.handler.classkafka.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.clientdans 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 pascom.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
- Scala
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()
val kafkaOptions = Map(
"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",
)
val df = spark.readStream.format("kafka").options(kafkaOptions).load()
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
- Scala
- SQL
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()
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()
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>";'
);
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
- Scala
- SQL
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()
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()
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>";'
);
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 |
|---|---|
| Réglez sur |
| Chemin vers le fichier du magasin de confiance contenant les certificats d'autorité de certification (CA) de confiance. |
| Mot de passe du fichier de clés de confiance. |
| Format de fichier du magasin de confiance (default : |
| Chemin d'accès au fichier de magasin de clés contenant le certificat client et la clé privée (requis pour mTLS). |
| Mot de passe du fichier de stockage de clés. |
| Mot de passe de la clé privée dans le magasin de clés. |
| Algorithme de vérification du Hostname. La valeur par défaut est |
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
- Scala
- SQL
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>))
)
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>))
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>')
);
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_optionslors de l'authentificationCette 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.scopemanuellement. Pour la plupart des cas d'usage, définissez-le surhttps://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.clientsur 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.
- Assurez-vous que le compte de service dispose de la liaison de rôle IAM
-
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
startingOffsetsestlatest, qui ne lit que les nouvelles données. DéfinissezstartingOffsetssurearliestpour lire les données existantes.