Aller au contenu principal

Créez un connecteur personnalisé

info

Bêta

Cette fonctionnalité est en Bêta. Les administrateurs du Workspace peuvent contrôler l'accès à cette fonctionnalité à partir de la page Previews . Consultez Gérer les aperçus Databricks.

Cette page montre comment créer un connecteur pour une source qui n'est pas encore prise en charge dans Lakeflow Connect. Tout d'abord, créez et testez votre connecteur localement en utilisant les outils et les templates dans le repository Lakeflow Communauté Connectors sur GitHub. Le repository inclut des outils de développement basés sur l'IA pour faciliter chaque phase, y compris la recherche de source, la configuration de l'authentification, l'implémentation et les tests.

Lorsque votre connecteur personnalisé est prêt à être utilisé, essayez-le dans votre Workspace Databricks, puis enregistrez-le auprès de la Communauté en ouvrant une requête de tirage.

Pour utiliser un connecteur de la communauté enregistré, consultez Utiliser un connecteur de la communauté enregistré.

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 repository Lakeflow Communauté Connectors 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 les 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 communautaire implémente l'interface LakeflowConnect, qui définit comment 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

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 d'API) que votre connecteur expose. Databricks utilise cette liste pour remplir l'interface utilisateur de sélection des tables.

get_table_schema

Retourne 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 snapshot, cdc, cdc_with_deletes ou append.

read_table

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

read_table_deletes

Facultatif. Implémentez cette méthode uniquement si ingestion_type est cdc_with_deletes. Génère 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 d'API) que votre connecteur expose. Databricks utilise cette liste pour remplir l'interface utilisateur de sélection des tables.

get_table_schema

Retourne 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 snapshot, cdc, cdc_with_deletes ou append.

read_table

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

read_table_deletes

Facultatif. Implémentez cette méthode uniquement si ingestion_type est cdc_with_deletes. Génère les clés d'enregistrement supprimées et renvoie le décalage suivant.

Développez 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 source, les mécanismes d'authentification, les limites de débit et les schémas de données disponibles. Identifiez les tables ou les Endpoints à 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 standards sur un système source réel et corrigez les problèmes. Consultez Tester votre connecteur pour plus de détails.

  5. Documentez le connecteur : Rédigez un README.md destiné à l'utilisateur et générez le fichier YAML de spécification du connecteur qui décrit les paramètres configurables du connecteur.

  6. Générez l’artefact de déploiement : exécutez le script de build pour produire l’artefact à fichier unique qui peut être déployé dans un workspace.

Testez votre connecteur

Le repository propose plusieurs approches de test :

Suite de tests générique (obligatoire)

Se connecte à une source réelle à l’aide des identifiants que vous avez fournis pour vérifier la fonctionnalité de bout en bout, y compris l’authentification, la découverte de schémas et les lectures de données.

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

Tests de réécriture (recommandé)

Exécute des cycles d'écriture-lecture-vérification pour valider les lectures et les suppressions incrémentielles. Cela confirme que votre suivi des offsets et votre logique CDC fonctionnent correctement.

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

Tests unitaires

Rédigez des tests unitaires pour toute logique personnalisée complexe dans votre connecteur, telle que la gestion de la pagination, la conversion de type ou la récupération d'erreurs.

Créer l'artefact de déploiement

Une fois que votre connecteur a passé les suites de tests, exécutez le script de Merge pour générer un artefact de déploiement à fichier unique. Le pipeline utilise ce fichier lors de l'exécution plutôt que l'intégralité du repository.

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

Ceci produit un fichier Python autonome dans dist/<your-source>/ qui inclut tout le code et les dépendances du connecteur.

Créer un pipeline d’ingestion

Pour essayer votre connecteur :

  1. Dans la barre latérale de votre Workspace Databricks, cliquez sur +Nouveau > Ajouter ou upload de données , puis sélectionnez + Ajouter un connecteur de Communauté sous Connecteurs de Communauté .

  2. Pour Nom de la source , entrez le nom de votre connecteur.

  3. Pour l'**URL du GitHub repository**, saisissez l'URL du GitHub repository qui héberge le code source de votre connecteur.

  4. Cliquez sur Ajouter un connecteur .

  5. Cliquez sur + Créer une connexion ou sélectionnez une connexion existante, puis cliquez sur Suivant .

  6. Pour Nom de la pipeline , saisissez un nom pour la pipeline.

  7. Pour l'emplacement du log des événements , saisissez un nom de catalogue et un nom de schéma. Databricks stocke le Log des événements du pipeline ici. Les tables ingérées sont également écrites ici par default.

  8. Pour le chemin racine , saisissez votre chemin de Workspace (par exemple, /Workspace/Users/<your-email>/connectors). Databricks clone et stocke le code source du connecteur ici.

  9. Cliquez sur **Créer un pipeline**.

  10. Dans l’éditeur de pipeline, ouvrez ingest.py et modifiez le champ objects pour inclure les tables que vous souhaitez ingérer. Par exemple :

    Python
    from databricks.labs.community_connector.pipeline import ingest

    pipeline_spec = {
    "connection_name": "my_connector_connection", # Required: UC connection name
    "objects": [
    {"table": {"source_table": "my_table"}},
    ],
    }

    ingest(spark, pipeline_spec)
  11. Exécutez le pipeline manuellement ou planifiez-le.

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 informations d'identification 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 où les tables ingérées sont écrites. Par défaut, il s’agit du catalogue défini lors de la création du pipeline.

destination_schema

Le schéma où les tables ingérées sont écrites. Default to the schema set during pipeline creation.

scd_type

La stratégie de dimension à évolution lente : SCD_TYPE_1, SCD_TYPE_2 ou APPEND_ONLY. Default to SCD_TYPE_1.

primary_keys

Remplacer les clés principales par default d'une table. Fournissez une liste de noms de colonnes.

Option

Description

connection_name

Obligatoire. Le nom de la connexion qui stocke les informations d'identification 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 où les tables ingérées sont écrites. Par défaut, il s’agit du catalogue défini lors de la création du pipeline.

destination_schema

Le schéma où les tables ingérées sont écrites. Default to the schema set during pipeline creation.

scd_type

La stratégie de dimension à évolution lente : SCD_TYPE_1, SCD_TYPE_2 ou APPEND_ONLY. Default to SCD_TYPE_1.

primary_keys

Remplacer les clés principales par default d'une table. Fournissez une liste de noms de colonnes.

Enregistrez votre connecteur

Après avoir créé et testé votre connecteur, ouvrez une pull request dans le repository Lakeflow Community Connectors pour le rendre disponible à la communauté.