Aller au contenu principal

Référence du connecteur Kafka

Cette page contient la documentation de référence du connecteur 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

topics

Liste de chaînes

Liste des noms de sujets auxquels s'abonner. Mutuellement exclusif avec topic_pattern. topics ou topic_pattern est requis.

topic_pattern

Chaîne

Noms de rubriques correspondant aux expressions régulières Java auxquelles s'abonner. Exclusif mutuellement avec topics.

starting_offset

Chaîne

latest

Où commencer la lecture lorsqu'aucun point de contrôle n'existe (première exécution uniquement). Valeurs valides : latest, earliest.

key_transformer

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 BINARY. Consultez les options du Transformer.

value_transformer

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 BINARY. Consultez les options du Transformer.

Option

Type

Par défaut

Description

topics

Liste de chaînes

Liste des noms de sujets auxquels s'abonner. Mutuellement exclusif avec topic_pattern. topics ou topic_pattern est requis.

topic_pattern

Chaîne

Noms de rubriques correspondant aux expressions régulières Java auxquelles s'abonner. Exclusif mutuellement avec topics.

starting_offset

Chaîne

latest

Où commencer la lecture lorsqu'aucun point de contrôle n'existe (première exécution uniquement). Valeurs valides : latest, earliest.

key_transformer

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 BINARY. Consultez les options du Transformer.

value_transformer

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 BINARY. Consultez les options du Transformer.

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

format

Tout

Chaîne

Format de sérialisation des données. Valeurs valides : STRING, JSON. STRING ne requiert aucune option supplémentaire. Si aucun json_options n’est spécifié sur le transformateur, la valeur est analysée en tant que variante par default. Voir Format de données variant pour plus d’informations.

json_options.schema

JSON

Chaîne

Schéma en ligne au format DDL Spark (par exemple, "id BIGINT, name STRING"). Mutuellement exclusif avec schema_file_path.

json_options.schema_file_path

JSON

Chaîne

Chemin d'accès à un fichier de schéma .ddl. Mutuellement exclusif avec schema. Prend en charge les chemins d'accès (/Volumes/...) des Volumes Unity Catalog.

json_options.schema_evolution_mode

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_options.schema_hints

JSON

Chaîne

Paires "column_name type" séparées par des virgules pour influencer l'inférence de schéma (par exemple, "id BIGINT, ts TIMESTAMP"). Requiert que schema_evolution_mode soit défini. Voir Ignorer l'inférence de schéma à l'aide d'indications de schéma.

Option

S'applique à

Type

Par défaut

Description

format

Tout

Chaîne

Format de sérialisation des données. Valeurs valides : STRING, JSON. STRING ne requiert aucune option supplémentaire. Si aucun json_options n’est spécifié sur le transformateur, la valeur est analysée en tant que variante par default. Voir Format de données variant pour plus d’informations.

json_options.schema

JSON

Chaîne

Schéma en ligne au format DDL Spark (par exemple, "id BIGINT, name STRING"). Mutuellement exclusif avec schema_file_path.

json_options.schema_file_path

JSON

Chaîne

Chemin d'accès à un fichier de schéma .ddl. Mutuellement exclusif avec schema. Prend en charge les chemins d'accès (/Volumes/...) des Volumes Unity Catalog.

json_options.schema_evolution_mode

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_options.schema_hints

JSON

Chaîne

Paires "column_name type" séparées par des virgules pour influencer l'inférence de schéma (par exemple, "id BIGINT, ts TIMESTAMP"). Requiert que schema_evolution_mode soit défini. Voir Ignorer l'inférence de schéma à l'aide d'indications de schéma.

Options de fanout

info

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

fanout_by

Chaîne

Obligatoire. Expression SQL évaluée par rapport à l’enregistrement source brut (avec les colonnes Kafka key et value) dont le résultat détermine le nom de la table de destination. La valeur devient le segment final du nom de la table : {destination_catalog}.{destination_schema}.{value}. La colonne value est binaire, convertissez-la en chaîne de caractères avant d'en extraire un champ, par exemple cast(value as string):event_type::string. L'expression doit correspondre à une valeur STRING non nulle, et cette chaîne est utilisée tel quel comme segment de nom de table, sans guillemets ni nettoyage. Il doit donc s'agir d'un identifiant de table valide non entre guillemets. Les valeurs avec des espaces, des points ou d'autres caractères non valides dans un identifiant non entre guillemets échouent à l'écriture, tout comme une valeur entièrement composée de chiffres (par exemple, 123). Un chiffre initial est autorisé si la valeur contient également une lettre ou un trait de soulignement, tel que 2024_events. Les tables de destination sont créées automatiquement si elles n'existent pas déjà.

transforms

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 fanout_by voit. Un seul transform est autorisé, et il doit utiliser format: JSON. Voir les options de transformation Fanout.

Option

Type

Par défaut

Description

fanout_by

Chaîne

Obligatoire. Expression SQL évaluée par rapport à l’enregistrement source brut (avec les colonnes Kafka key et value) dont le résultat détermine le nom de la table de destination. La valeur devient le segment final du nom de la table : {destination_catalog}.{destination_schema}.{value}. La colonne value est binaire, convertissez-la en chaîne de caractères avant d'en extraire un champ, par exemple cast(value as string):event_type::string. L'expression doit correspondre à une valeur STRING non nulle, et cette chaîne est utilisée tel quel comme segment de nom de table, sans guillemets ni nettoyage. Il doit donc s'agir d'un identifiant de table valide non entre guillemets. Les valeurs avec des espaces, des points ou d'autres caractères non valides dans un identifiant non entre guillemets échouent à l'écriture, tout comme une valeur entièrement composée de chiffres (par exemple, 123). Un chiffre initial est autorisé si la valeur contient également une lettre ou un trait de soulignement, tel que 2024_events. Les tables de destination sont créées automatiquement si elles n'existent pas déjà.

transforms

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 fanout_by voit. Un seul transform est autorisé, et il doit utiliser format: JSON. Voir les options de transformation Fanout.

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

format

Chaîne

JSON

Facultatif. Format de sérialisation de la transformation. Seul JSON est pris en charge pour les transformations fanout, et c'est le default en cas d'omission. STRING, AVRO et PROTOBUF ne sont pas pris en charge.

input_column

Chaîne

La colonne à partir de laquelle la transformation lit et réécrit (par exemple, value). La transformation JSON en éventail analyse cette colonne sur place.

Option

Type

Par défaut

Description

format

Chaîne

JSON

Facultatif. Format de sérialisation de la transformation. Seul JSON est pris en charge pour les transformations fanout, et c'est le default en cas d'omission. STRING, AVRO et PROTOBUF ne sont pas pris en charge.

input_column

Chaîne

La colonne à partir de laquelle la transformation lit et réécrit (par exemple, value). La transformation JSON en éventail analyse cette colonne sur place.

remarque

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, broker1:9092,broker2:9092).

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.

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, broker1:9092,broker2:9092).

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, broker1:9092,broker2:9092).

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.

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, broker1:9092,broker2:9092).

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

key

BINARY

Le contenu binaire brut de la clé de message Kafka.

value

BINARY

Le contenu binaire brut de la valeur de message Kafka.

Colonne

Type

Description

key

BINARY

Le contenu binaire brut de la clé de message Kafka.

value

BINARY

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_options n'est spécifiée, la colonne est écrite comme VARIANT.
  • Si json_options.schema ou json_options.schema_file_path est spécifié, le JSON est analysé en colonnes typées correspondant au schéma.
  • Si json_options.schema_evolution_mode est défini, l'inférence de schéma est utilisée et le schéma évolue automatiquement.
remarque

Les colonnes de métadonnées Kafka (notamment topic, partition, offset, timestamp, timestampType et headers) ne sont pas disponibles dans la table de destination en version bêta.