Hinweis
Für den Zugriff auf diese Seite ist eine Autorisierung erforderlich. Sie können versuchen, sich anzumelden oder das Verzeichnis zu wechseln.
Für den Zugriff auf diese Seite ist eine Autorisierung erforderlich. Sie können versuchen, das Verzeichnis zu wechseln.
Von Bedeutung
Dieses Feature befindet sich in der Betaversion. Arbeitsbereichsadministratoren können den Zugriff auf dieses Feature über die Vorschauseite steuern. Siehe Manage Azure Databricks Previews.
Benutzerdefinierte Connectors ermöglichen es Ihnen, Daten von einer Quelle zu importieren, die Lakeflow Connect mit einem verwalteten Connector nicht unterstützt. Du baust und testest deinen Connector, dann deployst und führst ihn in deinem eigenen Azure Databricks Workspace aus. Du musst es nicht bei der Community registrieren oder in einem gemeinsamen Repository beitragen, um es zu nutzen.
Entwickeln Sie Ihren Connector mit den Tools und Vorlagen im Lakeflow Community Connectors-Repository auf GitHub. Das Repository enthält KI-basierte Entwicklungstools, die bei jeder Phase unterstützt werden, einschließlich Quellforschung, Authentifizierungseinrichtung, Implementierung und Tests.
Wenn du später deinen Connector mit anderen Nutzern teilen möchtest, kannst du ihn der Community beisteuern. Um einen bestehenden Community Connector zu verwenden, siehe Community Connectors in Lakeflow Connect.
Anforderungen
Bevor Sie beginnen, stellen Sie sicher, dass Sie über Folgendes verfügen:
- Python 3.10 oder höher
- Ein Azure Databricks Arbeitsbereich mit aktiviertem Unity-Katalog
- API-Anmeldeinformationen für die Quelle, mit der Sie eine Verbindung herstellen möchten
- Git lokal installiert
Einrichten des Repositorys
Klonen Sie das Lakeflow Community Connectors-Repository , und installieren Sie die Entwicklungsabhängigkeiten.
Klonen Sie das Repository:
git clone https://github.com/databrickslabs/lakeflow-community-connectors.git cd lakeflow-community-connectorsErstellen Sie eine virtuelle Umgebung, und installieren Sie Abhängigkeiten:
python -m venv .venv source .venv/bin/activate pip install -e ".[dev]"Überprüfen Sie die vorhandenen Connectorimplementierungen in
src/databricks/labs/community_connector/sources/, und beginnen Sie dann mit der Entwicklung Ihres Connectors in einem neuen Verzeichnis unter diesem Pfad. Folgen Sie den KI-unterstützten Entwicklungsbefehlen und -fähigkeiten des Repositorys. Verwenden Sie für den empfohlenen Workflow Folgendes:/develop-connector <your-source> /validate-connector <your-source>
Implementieren der LakeflowConnect Schnittstelle
Jeder Connector implementiert die LakeflowConnect Schnittstelle, die definiert, wie Ihr Connector authentifiziert, Tabellen entdeckt, Schemata zurückgibt und Daten liest.
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'."""
Methodenbeschreibungen
Die folgende Tabelle beschreibt jede Methode in der LakeflowConnect Schnittstelle:
| Methode | Description |
|---|---|
__init__ |
Empfängt die Verbindungsparameter als Wörterbuch und initialisiert den API-Client für Ihre Quelle. |
list_tables |
Gibt die Namen aller Tabellen (oder API-Endpunkte) zurück, die Ihr Connector verfügbar macht. Azure Databricks verwendet diese Liste, um die Tabellenauswahl-UI aufzufüllen. |
get_table_schema |
Gibt einen Spark StructType zurück, der das Schema für die angegebene Tabelle beschreibt. Wird vor der ersten Pipelineausführung und bei jeder Ausführung aufgerufen, wenn die Schemaentwicklung aktiviert ist. |
read_table_metadata |
Gibt ein Wörterbuch mit primary_keys, cursor_fieldund ingestion_type. Das ingestion_type muss eine von snapshot, cdc, cdc_with_deletes oder append sein. |
read_table |
Liefert Datensätze als Python-Wörterbücher und gibt den nächsten Offset für inkrementelle Lesevorgänge zurück. Bei der ersten Ausführung start_offset ist leer. Bei aufeinanderfolgenden Ausführungen enthält dies den Offset, der von der vorherigen Ausführung zurückgegeben wird. |
read_table_deletes |
Optional. Implementieren Sie diese Methode nur, wenn ingestion_typecdc_with_deletes ist. Gibt gelöschte Datensatzschlüssel zurück und gibt den nächsten Offset zurück. |
Entwickeln Sie Ihren Connector
Führen Sie die folgenden Schritte aus, um einen neuen Connector zu erstellen und zu überprüfen:
Recherchieren Sie die Quell-API: Untersuchen Sie die API-Spezifikationen, Authentifizierungsmechanismen, Ratelimits und verfügbare Datenschemas. Identifizieren Sie, welche Tabellen oder Endpunkte freigegeben werden sollen.
Authentifizierung einrichten: Generieren Sie die Verbindungsspezifikation, konfigurieren Sie Anmeldeinformationen für die Quelle, und überprüfen Sie die Konnektivität aus Ihrer Entwicklungsumgebung.
Implementieren Sie den Connector: Codieren Sie alle erforderlichen
LakeflowConnectSchnittstellenmethoden, um eine Verbindung mit der Quell-API herzustellen und Daten im erwarteten Format zurückzugeben.Testen und iterieren: Führen Sie die standardmäßigen Testsuites mit einem echten Quellsystem aus und beheben Sie alle Probleme. Details finden Sie unter "Testen des Connectors ".
Dokumentieren Sie den Verbinder: Schreiben Sie eine benutzerorientierte
README.mdDatei, und generieren Sie die YaML-Datei der Konnektorspezifikation, die die konfigurierbaren Parameter des Connectors beschreibt.Erstellen Sie das Bereitstellungsartefakt: Führen Sie das Buildskript aus, um das Einzeldateiartefakt zu erzeugen, das in einem Arbeitsbereich bereitgestellt werden kann.
Testen des Verbinders
Das Repository bietet mehrere Testansätze:
Generische Testsuite (erforderlich)
Diese Suite verbindet sich mit einer echten Quelle, die Ihre bereitgestellten Zugangsdaten nutzt, um die End-to-End-Funktionalität zu überprüfen, einschließlich Authentifizierung, Schema-Erkennung und Datenlesungen.
python -m pytest tests/generic/ --connector <your-source> --credentials credentials.json
Zurückschreiben von Tests (empfohlen)
Write-Back-Tests führen Write-Read-Verify-Zyklen durch, um inkrementelle Lese- und Löschvorgänge zu validieren. Dadurch wird bestätigt, dass die Offsetnachverfolgung und die CDC-Logik ordnungsgemäß funktionieren.
python -m pytest tests/writeback/ --connector <your-source> --credentials credentials.json
Komponententests
Schreiben Sie Komponententests für jede komplexe benutzerdefinierte Logik in Ihrem Connector, z. B. Paginierungsbehandlung, Typkoersion oder Fehlerwiederherstellung.
Bereitstellungsartefakt erstellen
Nachdem dein Connector die Testsuiten bestanden hat, paketiere ihn so, dass eine Pipeline ihn ausführen kann. Ein Konnektor wird in zwei Teilen bereitgestellt:
Ein Quellartefakt aus einer Einzeldatei. Führe das Merge-Skript aus, um deinen Connector in eine eigenständige Python-Datei abzuflachen. Die Pipeline verwendet diese Datei zur Laufzeit anstelle des vollständigen Repositorys.
python tools/scripts/merge_python_source.py --connector <your-source>Das Skript schreibt diese Datei nach
dist/<your-source>/.Python-Wheels (
.whl) für die Abhängigkeiten des Konnektors. Das Konnektorframework und alle Drittanbieterbibliotheken, die Ihr Konnektor importiert, müssen für die Pipeline als in einem Unity Catalog-Volume gespeicherte Wheels verfügbar sein. Wenn diese Wheels fehlen, kann die Pipeline während der Quellerkennung fehlschlagen. Du kannst sie selbst hochladen und in der Benutzeroberfläche referenzieren lassen oder die Community Connector CLI erstellen und für dich hochladen. Siehe Deploy with the Community Connector CLI.
Setzen Sie Ihren Stecker auf eine von zwei Arten ein:
- Die Azure Databricks UI ist der Point-and-Click-Pfad. Verwenden Sie dies für eine einmalige Bereitstellung, wenn Sie die Wheels des Konnektors selbst im Feld Bibliotheksabhängigkeiten bereitstellen. Siehe Deploy in der Databricks-Benutzeroberfläche.
-
Die
community-connectorCLI ist der skriptbare Weg. Nutze es, wenn du lokal entwickelst und die Räder des Connectors für dich bauen und hochladen möchtest, oder wenn du eine wiederholbare Bereitstellung möchtest, die du automatisieren kannst. Siehe Deploy with the Community Connector CLI.
In der Databricks-Benutzeroberfläche bereitstellen
Stellen Sie Ihren Connector in der Azure Databricks UI in zwei Phasen bereit: Fügen Sie den Connector hinzu und erstellen Sie dann die Pipeline.
Fügen Sie den benutzerdefinierten Stecker hinzu
Füge zuerst deinen Connector hinzu, sodass er als Kachel auf der Seite Daten hinzufügen erscheint:
- In der Seitenleiste Ihres Azure Databricks-Arbeitsbereichs klicken Sie auf +Neues>hinzufügen oder Daten hochladen und fügen dann unter Community-Connectors einen benutzerdefinierten Connector hinzu.
- Geben Sie für den Quellnamen den Namen des Connectors ein. Dies muss mit dem Verzeichnisnamen übereinstimmen, der den Quellcode Ihres Connectors enthält (
sources/<source-name>). - Geben Sie für Anzeigename einen leicht erkennbaren Namen für den Konnektor ein. Wenn du das leer lässt, wird standardmäßig der Quellname verwendet.
- Für Bibliotheksabhängigkeiten fügen Sie die Python Wheel (
.whl)-Dateien hinzu, die Ihr Connector benötigt, aus einem Unity-Katalog-Volume. Siehe Erstellen des Bereitstellungsartefakts. - Für Verbindungsspezifikation fügen Sie die Verbindungsspezifikation des Konnektors in YAML ein, passend zu seiner
connector_spec.yaml-Datei. - Klicke auf Speichern. Der Stecker erscheint als benutzerdefinierte Kachel unter Community-Verbinder.
Erstellen der Aufnahmepipeline
Erstellen Sie dann die Pipeline, die Daten von Ihrer Quelle aufnimmt:
- Wählen Sie die Kachel Ihres Connectors aus, um den Assistenten zum Erfassen von Daten zu öffnen.
- Im Verbindungsschritt klicken Sie auf + Verbindung erstellen oder wählen Sie eine bestehende Verbindung aus, geben Sie die Verbindungsdetails Ihrer Quelle ein und klicken Sie dann auf Nächst.
- Im Ingestion-Einrichtungsschritt geben Sie einen Pipeline-Namen ein, legen Sie den Standort des Ereignisprotokolls (Katalog und Schema) fest, wählen Sie einen Compute-Typ, klicken Sie dann auf Pipeline erstellen und fahren fort.
- Wählen Sie im Source-Schritt die zu importierenden Tabellen aus.
- Im Ziel-Schritt wählen Sie den Katalog und das Schema aus, in dem die eingesessenen Tabellen geschrieben sind.
- Im Schritt Zeitpläne und Benachrichtigungen stellen Sie einen optionalen Zeitplan und Benachrichtigungen ein und beenden Sie den Termin.
- Führen Sie die Pipeline manuell oder gemäß Zeitplan aus.
Um die Pipeline weiter zu konfigurieren, kannst du ingest.py im Pipeline-Editor bearbeiten. Siehe Pipeline-Konfigurationsoptionen.
Pipelinekonfigurationsoptionen
Sie können die folgenden Optionen konfigurieren in ingest.py:
| Auswahl | Description |
|---|---|
connection_name |
Erforderlich. Der Name der Verbindung, die Authentifizierungsanmeldeinformationen für die Quelle speichert. |
objects |
Erforderlich. Eine Liste der zu aufzunehmenden Tabellen. Jeder Eintrag weist das Format {"table": {"source_table": "..."}}auf. Sie können auch ein optionales destination_table innerhalb des table Objekts angeben. |
destination_catalog |
Der Katalog, in dem aufgenommen Tabellen abgelegt werden. Standardeinstellung für den Katalog, der während der Pipelineerstellung festgelegt wurde. |
destination_schema |
Das Schema, in das eingelesene Tabellen geschrieben werden. Standardeinstellung für das Schema, das während der Pipelineerstellung festgelegt wurde. |
scd_type |
Die langsam ändernde Dimensionstrategie: SCD_TYPE_1, SCD_TYPE_2 oder APPEND_ONLY. Wird standardmäßig auf SCD_TYPE_1 festgelegt. |
primary_keys |
Überschreiben Sie die voreingestellten Primärschlüssel einer Tabelle. Geben Sie eine Liste von Spaltennamen an. |
Bereitstellung mit der Community Connector CLI
Der publish-Befehl der CLI erstellt die Framework- und Connector-Wheel-Pakete aus Ihrem lokalen Quellcode, lädt sie in ein Unity Catalog-Volume hoch, hinterlegt deren Pfade im Connector-Manifest und veröffentlicht den Connector als benutzerdefinierte Kachel auf der Seite Daten hinzufügen. Verweisen Sie sie auf Ihre lokale Konnektorspezifikation, damit nicht im Upstream-Repository nach dem Konnektor gesucht wird:
community-connector publish <your-source> \
--spec src/databricks/labs/community_connector/sources/<your-source>/connector_spec.yaml
Um bereits erstellte Wheels wiederzuverwenden und den Schritt der Erstellung zu überspringen, übergeben Sie sie mit --package. Um einen bereits veröffentlichten Konnektor zu ersetzen, fügen Sie --overwrite hinzu. Für die vollständige Liste der Optionen, einschließlich --package, --volume-path, --catalog, und --schema, siehe die publish Befehlsreferenz.
Um den gesamten Workflow über die Kommandozeile auszuführen, einschließlich Verbindungsaufbau, Erstellung und Aktualisierung der Eingabepipeline, Veröffentlichen und Entveröffentlichen, siehe die Community Connector CLI-Referenz.
Trage deinen Connector zur Community bei
Dein Konnektor läuft in deinem Arbeitsbereich, egal, ob du ihn beisteuerst oder nicht. Wenn Sie es teilen möchten, damit andere Nutzer es entdecken und nutzen können, eröffnen Sie eine Pull Request im Lakeflow Community Connectors-Repository . Beigetragene Connectors werden zu Community-Connectors, die von der Community gepflegt werden und nicht durch Databricks-SLAs abgedeckt sind.