Aller au contenu principal

Créer un connecteur personnalisé pour LakeFlow Connect

info

Bêta

Cette fonctionnalité est en bêta. Les administrateurs de Workspace peuvent contrôler l'accès à cette fonctionnalité à partir de la page Aperçus . Voir Gérer les prévisualisations Databricks.

Les connecteurs personnalisés vous permettent d'ingérer des données à partir d'une source que Lakeflow Connect ne prend pas en charge avec un connecteur géré. Vous créez et testez votre connecteur, puis vous le déployez et l'exécutez dans votre propre workspace Databricks. Vous n'avez pas besoin de l'enregistrer auprès de la communauté ou de le contribuer à un repository partagé pour l'utiliser.

Développez votre connecteur à l'aide des outils et templates disponibles dans le repository Lakeflow Community Connectors sur GitHub. Le repository inclut des outils de développement optimisés par l'IA pour vous assister à chaque phase, notamment la recherche de sources, la configuration de l'authentification, l'implémentation et les tests.

Si vous souhaitez partager votre connecteur avec d'autres utilisateurs ultérieurement, vous pouvez le contribuer à la communauté. Pour utiliser un connecteur communautaire existant, consultez Connecteurs communautaires dans Lakeflow Connect.

Exigences

Avant de start, assurez-vous d’avoir :

  • Python 3.10 ou version ultérieure
  • Un workspace Databricks avec Unity Catalog activé
  • Identifiants API pour la source à laquelle vous souhaitez vous connecter
  • Git installé localement

Configurer le repository

Clonez le Lakeflow Communauté Connectors repository et installez les dépendances de développement.

  1. Cloner le repository :

    Bash
    git clone https://github.com/databrickslabs/lakeflow-community-connectors.git
    cd lakeflow-community-connectors
  2. Créez un environnement virtuel et installez les dépendances :

    Bash
    python -m venv .venv
    source .venv/bin/activate
    pip install -e ".[dev]"
  3. Examinez les implémentations de connecteurs existantes dans src/databricks/labs/community_connector/sources/, puis start développer votre connecteur dans un nouveau répertoire sous ce chemin. Suivez les commandes et compétences de développement assisté par IA du repository. Pour le workflow recommandé, utilisez :

    Text
    /develop-connector <your-source>
    /validate-connector <your-source>

Implémenter l’interface LakeflowConnect

Chaque connecteur implémente l’interface LakeflowConnect, qui définit la manière dont votre connecteur s’authentifie, découvre les tables, renvoie les schémas et lit les données.

Python
class LakeflowConnect:
def __init__(self, options: dict[str, str]) -> None:
"""Initialize with connection parameters"""

def list_tables(self) -> list[str]:
"""Return names of all tables supported by this connector."""

def get_table_schema(self, table_name: str, table_options: dict[str, str]) -> StructType:
"""Return the Spark schema for a table."""

def read_table_metadata(self, table_name: str, table_options: dict[str, str]) -> dict:
"""Return metadata: primary_keys, cursor_field, ingestion_type
(snapshot|cdc|cdc_with_deletes|append)."""

def read_table(self, table_name: str, start_offset: dict,
table_options: dict[str, str]) -> (Iterator[dict], dict):
"""Yield records as JSON dicts and return the next offset
for incremental reads."""

def read_table_deletes(self, table_name: str, start_offset: dict,
table_options: dict[str, str]) -> (Iterator[dict], dict):
"""Optional: Only required if ingestion_type is 'cdc_with_deletes'."""

Descriptions des méthodes

Le tableau suivant décrit chaque méthode dans l'interface LakeflowConnect :

Méthode

Description

__init__

Reçoit les paramètres de connexion sous forme de dictionnaire et initialise le client API pour votre source.

list_tables

Renvoie les noms de toutes les tables (ou Endpoint API) que votre connecteur expose. Databricks utilise cette liste pour remplir l’interface utilisateur de sélection de table.

get_table_schema

Renvoie un Spark StructType décrivant le schéma de la table donnée. Appelé avant la première exécution du pipeline et à chaque exécution lorsque l’évolution des schémas est activée.

read_table_metadata

Renvoie un dictionnaire avec primary_keys, cursor_field et ingestion_type. Le ingestion_type doit être l'un des éléments suivants : snapshot, cdc, cdc_with_deletes ou append.

read_table

Produit des enregistrements sous forme de dictionnaires Python et renvoie le décalage suivant pour les lectures incrémentielles. Lors de la première exécution, start_offset est vide. Lors des exécutions ultérieures, il contient le décalage renvoyé par l'exécution précédente.

read_table_deletes

Facultatif. N’implémentez cette méthode que si ingestion_type est cdc_with_deletes. Produit les clés d’enregistrement supprimées et renvoie le décalage suivant.

Méthode

Description

__init__

Reçoit les paramètres de connexion sous forme de dictionnaire et initialise le client API pour votre source.

list_tables

Renvoie les noms de toutes les tables (ou Endpoint API) que votre connecteur expose. Databricks utilise cette liste pour remplir l’interface utilisateur de sélection de table.

get_table_schema

Renvoie un Spark StructType décrivant le schéma de la table donnée. Appelé avant la première exécution du pipeline et à chaque exécution lorsque l’évolution des schémas est activée.

read_table_metadata

Renvoie un dictionnaire avec primary_keys, cursor_field et ingestion_type. Le ingestion_type doit être l'un des éléments suivants : snapshot, cdc, cdc_with_deletes ou append.

read_table

Produit des enregistrements sous forme de dictionnaires Python et renvoie le décalage suivant pour les lectures incrémentielles. Lors de la première exécution, start_offset est vide. Lors des exécutions ultérieures, il contient le décalage renvoyé par l'exécution précédente.

read_table_deletes

Facultatif. N’implémentez cette méthode que si ingestion_type est cdc_with_deletes. Produit les clés d’enregistrement supprimées et renvoie le décalage suivant.

Développer votre connecteur

Suivez ces étapes pour créer et valider un nouveau connecteur :

  1. Recherchez l'API source : étudiez les spécifications de l'API de la source, les mécanismes d'authentification, les limites de débit et les schémas de données disponibles. Identifiez les tables ou les Endpoint à exposer.

  2. Configurer l'authentification : générez la spécification de connexion, configurez les identifiants pour la source et vérifiez la connectivité depuis votre environnement de développement.

  3. Implémentez le connecteur : codez toutes les méthodes d'interface LakeflowConnect requises pour vous connecter à l'API source et renvoyer les données dans le format attendu.

  4. Tester et itérer : exécutez les suites de tests standard sur un système source réel et corrigez les problèmes éventuels. Voir Tester votre connecteur pour plus de détails.

  5. Documenter le connecteur : rédigez un README.md destiné aux utilisateurs et générez le fichier YAML de spécification du connecteur qui décrit les parameter configurables du connecteur.

  6. Générer l'artefact de déploiement : exécutez le script de build pour produire l'artefact à fichier unique pouvant être déployé dans un workspace.

Tester votre connecteur

Le repository propose plusieurs approches de test :

Suite de tests générique (requise)

Cette suite se connecte à une source réelle en utilisant les identifiants que vous avez fournis pour vérifier la fonctionnalité de bout en bout, y compris l’authentification, la découverte de schéma et les lectures de données.

Bash
python -m pytest tests/generic/ --connector <your-source> --credentials credentials.json

Test de réécriture (recommandé)

Les tests de réécriture exécutent des cycles écriture-lecture-vérification pour valider les lectures et suppressions incrémentielles. Ceci confirme que votre suivi de décalage et votre logique CDC fonctionnent correctement.

Bash
python -m pytest tests/writeback/ --connector <your-source> --credentials credentials.json

Tests unitaires

Écrivez des tests unitaires pour toute logique personnalisée complexe dans votre connecteur, telle que la gestion de la pagination, la coercition de type ou la reprise en cas d’erreur.

Générer l’artefact de déploiement

Une fois que votre connecteur a réussi les suites de tests, packagez-le afin qu'un pipeline puisse l'exécuter. Un connecteur se déploie en deux parties :

  • Un artefact source composé d’un seul fichier. Exécutez le script de Merge pour aplatir votre connecteur en un seul fichier Python autonome. Le pipeline utilise ce fichier au moment de l’exécution plutôt que le repository complet.

    Bash
    python tools/scripts/merge_python_source.py --connector <your-source>

    Le script écrit ce fichier dans dist/<your-source>/.

  • Python wheels (.whl) pour les dépendances du connecteur. Le framework de connecteur et toutes les bibliothèques tierces importées par votre connecteur doivent être disponibles pour le pipeline sous forme de wheels stockés dans un volume Unity Catalog. Si ces wheels sont manquants, le pipeline peut échouer lors de la découverte de la source. Vous pouvez les upload vous-même et les référencer dans l’interface utilisateur, ou laisser la CLI du connecteur communautaire les générer et les upload pour vous. Voir Déployer avec la CLI du connecteur communautaire.

Déployez votre connecteur de deux manières :

  • L’interface utilisateur Databricks est le chemin « pointer-cliquer ». Utilisez-le pour un déploiement ponctuel lorsque vous fournissez vous-même les wheels du connecteur dans le champ Dépendances de bibliothèque . Voir Déployer dans l’interface utilisateur Databricks.
  • La CLI community-connector est le chemin scriptable. Utilisez-le lorsque vous développez localement et que vous souhaitez que les wheels du connecteur soient générés et importés pour vous, ou lorsque vous souhaitez un déploiement reproductible que vous pouvez automatiser. Voir Déployer avec la CLI Community Connector.

Déployer dans l'interface utilisateur Databricks

Déployez votre connecteur dans l'interface utilisateur Databricks en deux phases : ajoutez le connecteur, puis créez le pipeline.

Ajouter le connecteur personnalisé

Tout d’abord, ajoutez votre connecteur afin qu’il apparaisse sous forme de vignette sur la page Ajouter des données :

  1. Dans la barre latérale de votre Databricks workspace, cliquez sur +Nouveau > Ajouter ou upload des données , puis, sous Connecteurs communautaires , ajoutez un connecteur personnalisé.
  2. Pour le Nom de la source , saisissez le nom de votre connecteur. Cela doit correspondre au nom du répertoire qui contient le code source de votre connecteur (sources/<source-name>).
  3. Pour Nom d'affichage , saisissez un nom convivial pour le connecteur. Si vous laissez ce champ vide, le nom source sera utilisé par default.
  4. Pour les dépendances de bibliothèque , ajoutez les fichiers Python wheel (.whl) dont votre connecteur a besoin à partir d’un volume Unity Catalog. Voir Générer l’artefact de déploiement.
  5. Pour Connection specification , collez la spécification de connexion du connecteur au format YAML, en faisant correspondre son fichier connector_spec.yaml.
  6. Cliquez sur Enregistrer . Le connecteur apparaît sous la forme d’une tuile Personnalisé sous Connecteurs communautaires .

Créer le pipeline d’ingestion

Créez ensuite le pipeline qui ingère les données à partir de votre source :

  1. Sélectionnez la tuile de votre connecteur pour ouvrir l’assistant Ingest data .
  2. À l’étape Connection , cliquez sur + Create connection ou sélectionnez une connexion existante, saisissez les détails de connexion de votre source, puis cliquez sur Next .
  3. À l’étape Ingestion setup , saisissez un Pipeline name , définissez l’ Event log location (catalogue et schéma), choisissez un Compute type , puis cliquez sur Create pipeline and continue .
  4. À l’étape Source , sélectionnez les tables à ingérer.
  5. À l’étape Destination , choisissez le catalogue et le schéma où les tables ingérées sont écrites.
  6. À l’étape Calendriers et notifications , définissez un calendrier et des notifications facultatifs, puis terminez.
  7. Exécutez le pipeline manuellement ou selon son calendrier.

Pour configurer davantage le pipeline, vous pouvez modifier ingest.py dans l’éditeur de pipeline. Voir Options de configuration du pipeline.

Options de configuration du pipeline

Vous pouvez configurer les options suivantes dans ingest.py:

Option

Description

connection_name

Obligatoire. Le nom de la connexion qui stocke les identifiants d'authentification pour la source.

objects

Obligatoire. Une liste de tables à ingérer. Chaque entrée a le format {"table": {"source_table": "..."}}. Vous pouvez également spécifier un destination_table facultatif à l'intérieur de l'objet table.

destination_catalog

Le catalogue dans lequel les tables ingérées sont écrites. Utilise default le catalogue défini lors de la création du pipeline.

destination_schema

Le schéma où les tables ingérées sont écrites. Utilise default le schéma défini lors de la création du pipeline.

scd_type

La stratégie de dimension à évolution lente : SCD_TYPE_1, SCD_TYPE_2 ou APPEND_ONLY. La valeur default est SCD_TYPE_1.

primary_keys

Remplacer les clés primaires default d'une table. Fournissez une liste de noms de colonne.

Option

Description

connection_name

Obligatoire. Le nom de la connexion qui stocke les identifiants d'authentification pour la source.

objects

Obligatoire. Une liste de tables à ingérer. Chaque entrée a le format {"table": {"source_table": "..."}}. Vous pouvez également spécifier un destination_table facultatif à l'intérieur de l'objet table.

destination_catalog

Le catalogue dans lequel les tables ingérées sont écrites. Utilise default le catalogue défini lors de la création du pipeline.

destination_schema

Le schéma où les tables ingérées sont écrites. Utilise default le schéma défini lors de la création du pipeline.

scd_type

La stratégie de dimension à évolution lente : SCD_TYPE_1, SCD_TYPE_2 ou APPEND_ONLY. La valeur default est SCD_TYPE_1.

primary_keys

Remplacer les clés primaires default d'une table. Fournissez une liste de noms de colonne.

Déployer avec la CLI du connecteur communautaire

La commande publish de la CLI génère le framework et les wheels du connecteur à partir de votre source locale, les upload vers un volume Unity Catalog, enregistre leurs chemins dans le manifeste du connecteur et publie le connecteur en tant que vignette Personnalisé sur la page Ajouter des données . Pointez-le vers votre spécification de connecteur locale afin qu’il ne recherche pas le connecteur dans le repository amont :

Bash
community-connector publish <your-source> \
--spec src/databricks/labs/community_connector/sources/<your-source>/connector_spec.yaml

Pour réutiliser les wheels que vous avez déjà créés et ignorer l'étape de build, transmettez-les avec --package. Pour remplacer un connecteur que vous avez publié précédemment, ajoutez --overwrite. Pour obtenir la liste complète des options, y compris --package, --volume-path, --catalog et --schema, consultez la référence de commandepublish.

Pour exécuter l'intégralité du workflow à partir de la ligne de commande, notamment la création d'une connexion, la création et la mise à jour du pipeline d'ingestion, la publication et l'annulation de la publication, consultez la référence de la CLI du connecteur de communauté.

Contribuez votre connecteur à la communauté

Votre connecteur s’exécute dans votre workspace, que vous y contribuiez ou non. Si vous souhaitez le partager afin que d’autres utilisateurs puissent le découvrir et l’utiliser, ouvrez une pull request dans le repository Lakeflow Communauté Connectors. Les connecteurs fournis deviennent des connecteurs communautaires, qui sont maintenus par la communauté et ne sont pas couverts par les SLA de Databricks.