Construye un conector personalizado para Lakeflow Connect

Importante

Esta característica se encuentra en su versión beta. Los administradores del área de trabajo pueden controlar el acceso a esta característica desde la página Vistas previas . Consulte Administrar versiones preliminares de Azure Databricks.

Los conectores personalizados permiten ingir datos de una fuente que Lakeflow Connect no soporta con un conector gestionado. Construyes y pruebas tu conector, luego lo despliegas y ejecutas en tu propio espacio de trabajo de Azure Databricks. No necesitas registrarlo en la comunidad ni contribuirlo a ningún repositorio compartido para usarlo.

Desarrolla tu conector utilizando las herramientas y plantillas del repositorio Lakeflow Community Connectors en GitHub. El repositorio incluye herramientas de desarrollo con tecnología de inteligencia artificial para ayudar con cada fase, incluida la investigación de origen, la configuración de autenticación, la implementación y las pruebas.

Si más adelante quieres compartir tu conector con otros usuarios, puedes contribuirlo a la comunidad. Para utilizar un conector comunitario existente, consulte Conectores comunitarios en Lakeflow Connect.

Requisitos

Antes de empezar, asegúrese de que dispone de lo siguiente:

  • Python 3.10 o superior
  • Un área de trabajo de Azure Databricks con el catálogo de Unity habilitado
  • Credenciales de API para el origen al que desea conectarse
  • Git instalado localmente

Configuración del repositorio

Clone el repositorio lakeflow Community Connectors e instale las dependencias de desarrollo.

  1. Clone el repositorio:

    git clone https://github.com/databrickslabs/lakeflow-community-connectors.git
    cd lakeflow-community-connectors
    
  2. Cree un entorno virtual e instale las dependencias:

    python -m venv .venv
    source .venv/bin/activate
    pip install -e ".[dev]"
    
  3. Revise las implementaciones del conector existentes en src/databricks/labs/community_connector/sources/y, a continuación, empiece a desarrollar el conector en un nuevo directorio en esa ruta de acceso. Siga los comandos y las habilidades de desarrollo asistido por IA del repositorio. Para el flujo de trabajo recomendado, use:

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

Implementación de la LakeflowConnect interfaz

Cada conector implementa la LakeflowConnect interfaz, que define cómo tu conector se autentica, descubre tablas, devuelve esquemas y lee datos.

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'."""

Descripciones del método

La siguiente tabla describe cada método en la LakeflowConnect interfaz:

Método Description
__init__ Recibe los parámetros de conexión como un diccionario e inicializa el cliente de API para el origen.
list_tables Devuelve los nombres de todas las tablas (o puntos de conexión de API) que expone el conector. Azure Databricks usa esta lista para rellenar la interfaz de usuario de selección de tabla.
get_table_schema Devuelve un spark StructType que describe el esquema de la tabla especificada. Se llama antes de la primera ejecución de canalización y en cada ejecución cuando se habilita la evolución del esquema.
read_table_metadata Devuelve un diccionario con primary_keys, cursor_fieldy ingestion_type. ingestion_type debe ser uno de los elementos snapshot, cdc, cdc_with_deleteso append.
read_table Devuelve registros como diccionarios de Python y devuelve el siguiente desplazamiento para las lecturas incrementales. En la primera ejecución, start_offset está vacío. En ejecuciones posteriores, contiene el desplazamiento devuelto por la ejecución anterior.
read_table_deletes Optional. Implemente este método solo si ingestion_type es cdc_with_deletes. Produce claves de registro eliminadas y devuelve el siguiente desplazamiento.

Desarrollo del conector

Siga estos pasos para compilar y validar un nuevo conector:

  1. Investigación de la API de origen: estudie las especificaciones de api del origen, los mecanismos de autenticación, los límites de velocidad y los esquemas de datos disponibles. Identifique qué tablas o puntos de conexión se van a exponer.

  2. Configuración de la autenticación: genere la especificación de conexión, configure las credenciales para el origen y compruebe la conectividad desde el entorno de desarrollo.

  3. Implemente el conector: codigo todos los métodos de interfaz necesarios LakeflowConnect para conectarse a la API de origen y devolver datos en el formato esperado.

  4. Probar e iterar: corra los conjuntos de pruebas estándar en un sistema de origen real y resuelva cualquier problema. Consulte Probar el conector para obtener más información.

  5. Documente el conector: escriba un archivo YAML orientado al README.md usuario y genere el archivo YAML de especificación del conector que describe los parámetros configurables del conector.

  6. Compilación del artefacto de implementación: ejecute el script de compilación para generar el artefacto de un solo archivo que se puede implementar en un área de trabajo.

Prueba del conector

El repositorio proporciona varios enfoques de prueba:

Conjunto de pruebas genérico (obligatorio)

Esta suite se conecta a una fuente real usando tus credenciales proporcionadas para verificar funcionalidades de extremo a extremo, incluyendo autenticación, descubrimiento de esquemas y lecturas de datos.

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

Las pruebas de reescritura ejecutan ciclos de escritura-lectura-verificación para validar las lecturas y eliminaciones incrementales. Esto confirma que el seguimiento de desplazamiento y la lógica CDC funcionan correctamente.

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

Pruebas unitarias

Escriba pruebas unitarias para cualquier lógica personalizada compleja en el conector, como el control de paginación, la coerción de tipos o la recuperación de errores.

Construir el artefacto de despliegue

Después de que tu conector pase las suites de pruebas, empaquetalo para que una tubería pueda ejecutarlo. Un conector se despliega en dos partes:

  • Un artefacto de origen de un solo archivo. Ejecuta el script de fusión para aplanar tu conector en un único archivo Python independiente. La canalización usa este archivo en tiempo de ejecución en lugar del repositorio completo.

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

    El script escribe este archivo en dist/<your-source>/.

  • Wheels de Python (.whl) para las dependencias del conector. El marco del conector y cualquier biblioteca de terceros que importe tu conector deben estar disponibles para el canal como wheels almacenados en un volumen de Unity Catalog. Si faltan estas ruedas, la tubería puede fallar durante el descubrimiento de la fuente. Puedes subirlos tú mismo y referenciarlos en la interfaz, o dejar que la CLI del Community Connector los construya y suba por ti. Consulta Desplegar con la CLI de Community Connector.

Despliega tu conector de una de dos maneras:

  • La interfaz de usuario de Azure Databricks es la ruta de apuntar y hacer clic. Úsala para una implementación puntual cuando proporciones tú mismo los wheels del conector en el campo Dependencias de la biblioteca. Consulta Desplegar en la interfaz de Databricks.
  • La community-connector CLI es la vía automatizable mediante scripts. Úsalo cuando estés desarrollando localmente y quieras que las ruedas del conector se construyan y suban para ti, o cuando quieras un despliegue repetible que puedas automatizar. Consulte Implementar con la CLI de conectores de la comunidad.

Despliegue en la interfaz de Databricks

Despliega tu conector en la interfaz de Azure Databricks en dos fases: añade el conector y luego crea la canalización.

Añadir el conector personalizado

Primero, añade tu conector para que aparezca como un mosaico en la página de Añadir datos :

  1. En la barra lateral de tu espacio de trabajo de Azure Databricks, haz clic en +Nuevo>Añadir o subir datos y, en conectores de la Comunidad, añade un conector personalizado.
  2. En Nombre de origen, escriba el nombre del conector. Esto debe coincidir con el nombre del directorio que contiene el código fuente de tu conector (sources/<source-name>).
  3. En Nombre para mostrar, introduce un nombre descriptivo para el conector. Si dejas esto en blanco, por defecto aparece el nombre de la fuente.
  4. En Dependencias de la biblioteca, añada los archivos wheel de Python (.whl) que necesita su conector desde un volumen de Unity Catalog. Consulta Crear el artefacto de implementación.
  5. Para Especificación de conexión, pega la especificación de conexión del conector en formato YAML, de modo que coincida con su archivo connector_spec.yaml.
  6. Haz clic en Guardar. El conector aparece como un mosaico Personalizado en Conectores de la comunidad.

Creación de la canalización de ingesta

Luego crea la tubería que ingiere datos de tu fuente:

  1. Seleccione el mosaico de su conector para abrir el asistente Importar datos.
  2. En el paso de Conexión , haz clic + Crear conexión o selecciona una conexión existente, introduce los datos de la conexión de tu fuente y luego haz clic en Siguiente.
  3. En el paso Configuración de la ingestión, introduce un nombre de pipeline, establece la ubicación del registro de eventos (catálogo y esquema), elige un tipo de proceso y, a continuación, haz clic en Crear pipeline y continuar.
  4. En el paso Fuente, selecciona las tablas que se van a importar.
  5. En el paso Destino, elige el catálogo y el esquema donde se escriben las tablas ingeridas.
  6. En el paso de Horarios y notificaciones, establece un horario opcional y notificaciones, y luego termina.
  7. Ejecuta la canalización manualmente o según lo programado.

Para configurar aún más el pipeline, puedes editar ingest.py en el editor de pipeline. Consulta Opciones de configuración de la canalización.

Opciones de configuración de canalización

Puede configurar las siguientes opciones en ingest.py:

Opción Description
connection_name Required. Nombre de la conexión que almacena las credenciales de autenticación para el origen.
objects Required. Lista de tablas que se van a ingerir. Cada entrada tiene el formato {"table": {"source_table": "..."}}. También puede especificar un opcional destination_table dentro del table objeto .
destination_catalog Catálogo donde se escriben las tablas ingeridas. El valor predeterminado es el catálogo establecido durante la creación de la canalización.
destination_schema Esquema donde se escriben las tablas ingeridas. Se asigna por defecto al esquema establecido durante la creación de la canalización.
scd_type Estrategia de dimensión de variación lenta: SCD_TYPE_1, SCD_TYPE_2o APPEND_ONLY. Tiene como valor predeterminado SCD_TYPE_1.
primary_keys Invalide las claves principales predeterminadas para una tabla. Proporcione una lista de nombres de columna.

Implementa con la CLI de Community Connector

El comando publish de la CLI compila el marco y los archivos wheel del conector a partir de tu código fuente local, los sube a un volumen de Unity Catalog, registra sus rutas en el manifiesto del conector y publica el conector como un mosaico Personalizado en la página Añadir datos. Indícale la especificación de tu conector local para que no busque el conector en el repositorio de origen:

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

Para reutilizar los wheels que ya has creado y saltarte el paso de creación, pásalos con --package. Para reemplazar un conector que publicaste antes, añade --overwrite. Para ver la lista completa de opciones, incluidas --package, --volume-path, --catalog y --schema, consulte la publishreferencia de comandos.

Para ejecutar todo el flujo de trabajo desde la línea de comandos, incluyendo la creación de una conexión, la creación y actualización de la pipeline de ingestión, la publicación y la despublicación, consulta la referencia de la CLI del Community Connector.

Aporta tu conector a la comunidad

Tu conector se ejecuta en tu espacio de trabajo, independientemente de si lo compartes o no. Si quieres compartirlo para que otros usuarios puedan descubrirlo y utilizarlo, abre una solicitud pull en el repositorio de Lakeflow Community Connectors . Los conectores contribuidos se convierten en conectores comunitarios, que la comunidad mantiene y que no están respaldados por SLA de Databricks.