Référence du connecteur Kafka
Cette page documente les options de connecteur, les options de configuration de table et les paramètres de transformateur JSON pour le connecteur Apache Kafka géré dans Lakeflow Connect.
Options du connecteur
Les options suivantes configurent la source Kafka pour chaque table de destination dans le pipeline d’ingestion. Spécifiez ces options sous connector_options.kafka_options dans votre définition de pipeline. Consultez Exemples pour obtenir des exemples de pipelines complets.
Option | Type | Par défaut | Description |
|---|---|---|---|
| Liste de chaînes | — | Liste des noms de sujets auxquels s'abonner. Mutuellement exclusif avec |
| Chaîne | — | Noms de rubriques correspondant aux expressions régulières Java auxquelles s'abonner. Exclusif mutuellement avec |
| Chaîne |
| Où commencer la lecture lorsqu'aucun point de contrôle n'existe (première exécution uniquement). Valeurs valides : |
| Transformer | — | Configuration du désérialiseur pour les clés de message. S’il n’est pas défini, la colonne clé est conservée en tant que |
| Transformer | — | Configuration du désérialiseur pour les valeurs de message. S'il n'est pas défini, la colonne de valeur est conservée en tant que |
Options du Transformer
Les transformateurs définissent comment les clés et les valeurs des messages binaires Kafka sont désérialisées en colonnes structurées. Spécifiez le format de sérialisation et les options spécifiques au format correspondantes sous key_transformer ou value_transformer. Vous pouvez configurer un transformateur pour la clé, la valeur ou les deux indépendamment. Si aucun transformateur n'est défini, la colonne est conservée telle quelle BINARY.
Pour le JSON, vous pouvez fournir un schéma explicite, utiliser l'inférence de schéma avec évolution, ou omettre json_options entièrement pour stocker la valeur comme une colonne VARIANT.
Option | S'applique à | Type | Par défaut | Description |
|---|---|---|---|---|
| Tout | Chaîne | — | Format de sérialisation des données. Valeurs valides : |
| JSON | Chaîne | — | Schéma en ligne au format DDL Spark (par exemple, |
| JSON | Chaîne | — | Chemin d'accès à un fichier de schéma |
| JSON | Chaîne | — | Mode d'évolution des schémas pour l'inférence automatique du schéma. Consulter Modes d'évolution des schémas. |
| JSON | Chaîne | — | Paires |
Options Avro
Définissez ces options sous avro_options lorsque format: AVRO. Fournissez le schéma de manière dynamique, à partir d’un fichier ou à partir d’un registre de schémas.
Option | Type | Par défaut | Description |
|---|---|---|---|
| Chaîne | — | Schéma Avro inclus au format JSON. S’exclut mutuellement avec |
| Chaîne | — | Chemin d’accès vers un fichier de schémas |
| Objet | — | Résoudre le schéma à partir d'un registre de schémas au moment de l'exécution au lieu de |
| Chaîne |
| Comment traiter les enregistrements dont la désérialisation échoue. Valeurs valides : |
Options Protobuf
Définissez ces options sous protobuf_options lorsque format: PROTOBUF. Fournissez un fichier de jeu de descripteurs compilé (.desc) et un nom de message, ou résolvez le schéma à partir d'un registre de schémas.
Option | Type | Par défaut | Description |
|---|---|---|---|
| Chaîne | — | Chemin d’accès à un fichier d’ensemble de descripteurs Protobuf compilé ( |
| Chaîne | — | Nom du type de message Protobuf entièrement qualifié (par exemple, |
| Objet | — | Résoudre le schéma à partir d’un registre de schémas à l’exécution au lieu de |
| Entier | — | Profondeur d’expansion maximale pour les champs Protobuf récursifs, que Spark SQL ne prend pas en charge en mode natif. Valeurs valides : |
| Chaîne |
| Comment traiter les enregistrements dont la désérialisation échoue. Valeurs valides : |
Options du registre de schémas
Définissez schema_registry sous avro_options ou protobuf_options pour résoudre le schéma à l'exécution à partir d'un registre de schémas compatible Confluent. By default, le pipeline s'authentifie auprès du registre à l'aide de la connexion source Kafka du pipeline, qui stocke l'URL du registre et la clé d'API. Pour vous authentifier avec une connexion Unity Catalog différente, configurez connection_name. Voir Propriétés de connexion.
Option | Type | Par défaut | Description |
|---|---|---|---|
| Chaîne | — | Obligatoire. Le sujet à résoudre dans le registre de schémas compatible Confluent. |
| Chaîne | — | Connexion Unity Catalog utilisée pour s'authentifier auprès du registre. default, utilise la connexion source Kafka du pipeline. Définissez cette option lorsque le registre utilise des identifiants différents. |
| Chaîne | — | Protobuf uniquement. Sélectionne un message lorsque le sujet définit plusieurs messages Protobuf. Simple ( |
Options de configuration de la table
Les options suivantes sont définies sous table_configuration sur un objet table, un frère de connector_options. Voir Exemples pour des exemples complets de pipelines.
Option | Type | Par défaut | Description |
|---|---|---|---|
| Chaîne | — | Nom d'une colonne de type struct ajoutée à la table de destination qui contient les métadonnées de la source Kafka pour chaque enregistrement. Voir Colonne de métadonnées source. Le nom ne doit pas être |
Options de fanout
Aperçu
Cette fonctionnalité est en préversion privée. Pour l'essayer, veuillez contacter votre interlocuteur Databricks.
Les options de fanout acheminent chaque enregistrement d'une source Kafka unique vers l'une des nombreuses tables de destination. Spécifiez ces options sous fanout_options sur un objet de schéma (pas un objet de table) dans votre définition de pipeline. Consultez Acheminement des enregistrements vers plusieurs tables (fanout) pour un exemple complet de pipeline et Limitations de fanout pour les contraintes.
Option | Type | Par défaut | Description |
|---|---|---|---|
| Chaîne | — | Obligatoire. Expression SQL évaluée par rapport à l’enregistrement source brut (avec les colonnes Kafka |
| Liste de Transformer | — | Une transformation appliquée à chaque enregistrement acheminé après l'acheminement, avant l'écriture dans sa table de destination. Parce qu'il s'exécute après le routage, cela n'affecte pas la valeur que |
Options de transformation Fanout
Chaque entrée dans transforms utilise les options suivantes. Seul le format JSON est pris en charge pour les transformations de fanout.
Option | Type | Par défaut | Description |
|---|---|---|---|
| Chaîne |
| Facultatif. Format de sérialisation de la transformation. Seul |
| Chaîne | — | La colonne à partir de laquelle la transformation lit et réécrit (par exemple, |
L'option de transformation output_column n'est pas appliquée aux transformations d'éclatement. La transformation JSON en éventail écrit toujours son résultat dans input_column sur place ; toute valeur output_column est ignorée.
Propriétés de connexion
Lorsque vous créez la connexion Kafka Unity Catalog dans Catalog Explorer, vous devez spécifier les propriétés suivantes selon la méthode d'authentification. Consultez Créer une connexion Kafka pour les étapes de création de connexion.
Nom d'utilisateur et mot de passe
Propriété | Description |
|---|---|
Nom de la connexion | Un nom unique pour la connexion dans Unity Catalog. |
Type de connexion | Select Kafka . |
Type d'authentification | Sélectionnez Nom d'utilisateur et mot de passe . |
Nom d'utilisateur | Nom d'utilisateur utilisé pour s'authentifier auprès du cluster Kafka. |
Mot de passe | Le mot de passe utilisé pour s'authentifier auprès du cluster Kafka. |
Serveurs d'amorçage | L’adresse du serveur bootstrap du cluster Kafka (par exemple, |
URL du registre de schémas (facultatif) | L’URL de votre registre de schémas. |
**Clé API du registre de schémas** (facultatif) | La clé API de votre registre de schémas. |
Secret de l'API du registre de schémas (facultatif) | Le secret d'API pour votre registre de schémas. |
Identifiant de service
Propriété | Description |
|---|---|
Nom de la connexion | Un nom unique pour la connexion dans Unity Catalog. |
Type de connexion | Select Kafka . |
Type d'authentification | Sélectionner l' identifiant de service . |
Identifiant de service | Sélectionnez un identifiant de service Unity Catalog existant ou cliquez sur Créer un nouvel identifiant de service . |
Serveurs d'amorçage | L’adresse du serveur bootstrap du cluster Kafka (par exemple, |
URL du registre de schémas (facultatif) | L’URL de votre registre de schémas. |
**Clé API du registre de schémas** (facultatif) | La clé API de votre registre de schémas. |
Secret de l'API du registre de schémas (facultatif) | Le secret d'API pour votre registre de schémas. |
Schéma de table de destination
Le connecteur Kafka écrit dans des tables de streaming (en mode ajout uniquement). Les colonnes écrites dans la table de destination dépendent de la configuration des transformateurs.
Sans transformateurs (binaire brut)
Lorsque key_transformer ou value_transformer n'est pas configuré, la table de destination contient les colonnes suivantes :
Colonne | Type | Description |
|---|---|---|
|
| Le contenu binaire brut de la clé de message Kafka. |
|
| Le contenu binaire brut de la valeur de message Kafka. |
Avec un transformateur STRING
Lorsque format: STRING est défini sur un transformateur, la colonne correspondante est écrite en tant que STRING au lieu de BINARY.
Avec un transformateur JSON
Lorsque format: JSON est défini sur un transformateur :
- Si aucune
json_optionsn'est spécifiée, la colonne est écrite commeVARIANT. - Si
json_options.schemaoujson_options.schema_file_pathest spécifié, le JSON est analysé en colonnes typées correspondant au schéma. - Si
json_options.schema_evolution_modeest défini, l'inférence de schéma est utilisée et le schéma évolue automatiquement.
Avec un transformateur Avro
Lorsque format: AVRO est défini sur un transformateur, la colonne est analysée en colonnes typées correspondant au schéma Avro. En mode PERMISSIVE (default), une colonne _corrupt_record de type BINARY est également ajoutée ; elle contient les octets bruts de tout enregistrement dont la désérialisation échoue et est null pour les enregistrements qui sont analysés avec succès.
Avec un transformateur Protobuf
Lorsque format: PROTOBUF est défini sur un transformateur, la colonne est analysée en colonnes typées correspondant à la définition de message Protobuf. En mode PERMISSIVE (default), une colonne _corrupt_record de type BINARY est également ajoutée ; elle contient les octets bruts de tout enregistrement dont la désérialisation a échoué et prend la valeur null pour les enregistrements analysés avec succès.
Colonne des métadonnées source
Définissez source_metadata_column sous table_configuration pour ajouter une colonne struct de ce nom à la table de destination. La structure contient les champs de métadonnées source Kafka suivants pour chaque enregistrement : topic, partition, offset, timestamp, timestampType et headers. Le nom de la colonne ne doit pas être key ou value. Voir les options de configuration de table.