Créer un connecteur personnalisé
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 des Template disponibles dans le repository Lakeflow Communauté Connectors sur GitHub. Le repository inclut des outils de développement optimisés par l'IA pour assister chaque phase, notamment la recherche de sources, la configuration de l'authentification, l'implémentation et les tests. L'utilisation de ces outils ne fait pas de votre connecteur un connecteur Communauté. Le repository fournit le framework et les exemples, et votre connecteur reste dans votre Workspace à moins que vous ne choisissiez de le contribuer.
Si vous souhaitez partager votre connecteur avec d'autres utilisateurs ultérieurement, vous pouvez éventuellement le contribuer à la Communauté. Pour utiliser un connecteur de la Communauté existant, consultez Connecteurs de la Communauté 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.
-
Cloner le repository :
Bashgit clone https://github.com/databrickslabs/lakeflow-community-connectors.git
cd lakeflow-community-connectors -
Créez un environnement virtuel et installez les dépendances :
Bashpython -m venv .venv
source .venv/bin/activate
pip install -e ".[dev]" -
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.
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 |
|---|---|
| Reçoit les paramètres de connexion sous forme de dictionnaire et initialise le client API pour votre source. |
| 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. |
| Renvoie un Spark |
| Renvoie un dictionnaire avec |
| 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, |
| Facultatif. N’implémentez cette méthode que si |
Développer votre connecteur
Suivez ces étapes pour créer et valider un nouveau connecteur :
-
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.
-
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.
-
Implémentez le connecteur : codez toutes les méthodes d'interface
LakeflowConnectrequises pour vous connecter à l'API source et renvoyer les données dans le format attendu. -
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.
-
Documenter le connecteur : rédigez un
README.mddestiné aux utilisateurs et générez le fichier YAML de spécification du connecteur qui décrit les parameter configurables du connecteur. -
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)
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.
python -m pytest tests/generic/ --connector <your-source> --credentials credentials.json
Test de réécriture (recommandé)
Exécute des cycles d'écriture-lecture-vérification pour valider les lectures et suppressions incrémentales. Ceci confirme que votre suivi de décalage et votre logique CDC fonctionnent correctement.
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, exécutez le script de Merge pour générer un artefact de déploiement en un seul fichier. Le pipeline utilise ce fichier au moment de l'exécution plutôt que le repository complet.
python tools/scripts/merge_python_source.py --connector <your-source>
Cela produit un fichier Python autonome dans dist/<your-source>/ qui inclut tout le code du connecteur et ses dépendances.
Créer un pipeline d'ingestion
Déployez et exécutez votre connecteur dans votre propre workspace Databricks :
-
Dans la barre latérale de votre workspace Databricks, cliquez sur +New > Add or upload data , puis choisissez l'option permettant d'ajouter un connecteur personnalisé.
-
Pour Nom de la source , saisissez le nom de votre connecteur.
-
Pour URL du repository GitHub , saisissez l’URL du repository GitHub qui héberge le code source de votre connecteur.
-
Cliquez sur Ajouter un connecteur .
-
Cliquez sur + Créer une connexion ou sélectionnez une connexion existante, puis cliquez sur Suivant .
-
Pour Nom du pipeline , saisissez un nom pour le pipeline.
-
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 default.
-
Pour Root path , saisissez votre chemin d'accès au Workspace (par exemple,
/Workspace/Users/<your-email>/connectors). Databricks clone et stocke le code source du connecteur ici. -
Cliquez sur Créer un pipeline .
-
Dans l'éditeur de pipeline, ouvrez
ingest.pyet modifiez le champ objects pour inclure les tables que vous souhaitez ingérer. Par exemple :Pythonfrom 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) -
Exécutez le pipeline manuellement ou planifiez-le.
Options de configuration du pipeline
Vous pouvez configurer les options suivantes dans ingest.py:
Option | Description |
|---|---|
| Obligatoire. Le nom de la connexion qui stocke les identifiants d'authentification pour la source. |
| Obligatoire. Une liste de tables à ingérer. Chaque entrée a le format |
| Le catalogue dans lequel les tables ingérées sont écrites. Utilise default le catalogue défini lors de la création du pipeline. |
| 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. |
| La stratégie de dimension à évolution lente : |
| Remplacer les clés primaires default d'une table. Fournissez une liste de noms de colonne. |
Contribuez votre connecteur à la communauté
La contribution de votre connecteur à la Communauté est facultative. Votre connecteur s'exécute dans votre Workspace sans cela. Si vous souhaitez partager votre connecteur afin que d'autres utilisateurs puissent le découvrir et l'utiliser, ouvrez une pull request dans le repository Lakeflow Communauté Connectors. Les connecteurs contribués deviennent des connecteurs de la Communauté, qui sont maintenus par la Communauté et ne sont pas couverts par les SLA de Databricks.