Ingesta de datos de SQL Server

Aprenda a ingerir datos de SQL Server en Azure Databricks mediante Lakeflow Connect.

El conector de SQL Server admite azure SQL Database, Azure SQL Managed Instance y bases de datos SQL de Amazon RDS. Esto incluye SQL Server que se ejecuta en máquinas virtuales (VM) de Azure y Amazon EC2. El conector también admite SQL Server local mediante redes de Azure ExpressRoute y AWS Direct Connect.

Requisitos

  • Para crear una puerta de enlace de ingesta y una canalización de ingesta, primero debe cumplir los siguientes requisitos:

    • El área de trabajo está habilitada para Unity Catalog.

    • El proceso sin servidor está habilitado para el área de trabajo. Consulte Requisitos de proceso sin servidor.

    • Si tiene previsto crear una conexión: tiene CREATE CONNECTION privilegios en el metastore. Consulte Administración de privilegios en Unity Catalog.

      Si el conector admite la creación de canalizaciones basadas en la interfaz de usuario, puede crear la conexión y la canalización al mismo tiempo completando los pasos de esta página. Sin embargo, si usa la creación de canalizaciones basadas en API, debe crear la conexión en el Explorador de catálogos antes de completar los pasos de esta página. Consulte Conexión a orígenes de ingesta administrados.

    • Si planea utilizar una conexión existente: Tiene privilegios USE CONNECTION o ALL PRIVILEGES en la conexión.

    • Tiene USE CATALOG privilegios en el catálogo de destino.

    • Tiene privilegios USE SCHEMA, CREATE TABLE y CREATE VOLUME en un esquema existente o privilegios CREATE SCHEMA en el catálogo de destino.

    • Tiene acceso a una instancia principal de SQL Server. Las características de captura de datos modificados y seguimiento de cambios no se admiten en réplicas de lectura ni en instancias secundarias.

    • Permisos sin restricciones para crear clústeres o una directiva personalizada (solo API). Una directiva personalizada para la puerta de enlace debe cumplir los siguientes requisitos:

      • Familia: Proceso de trabajos

      • Invalidaciones de familia de directivas:

        {
          "cluster_type": {
            "type": "fixed",
            "value": "dlt"
          },
          "num_workers": {
            "type": "unlimited",
            "defaultValue": 1,
            "isOptional": true
          },
          "runtime_engine": {
            "type": "fixed",
            "value": "STANDARD",
            "hidden": true
          }
        }
        
      • Databricks recomienda especificar los nodos de trabajo más pequeños posibles para las puertas de enlace de ingesta porque no afectan al rendimiento de la puerta de enlace. La siguiente política de cómputo permite a Azure Databricks escalar la puerta de enlace de ingestión para cumplir con los requisitos de su carga de trabajo. El requisito mínimo para el nodo controlador es de 8 núcleos para permitir una extracción eficiente y eficiente de datos de tu base de datos fuente.

        {
          "driver_node_type_id": {
            "type": "fixed",
            "value": "Standard_E64d_v4"
          },
          "node_type_id": {
            "type": "fixed",
            "value": "Standard_F4s"
          }
        }
        

      Para obtener más información sobre las directivas de clúster, consulte Seleccionar una directiva de procesamiento.

  • Para ingerir desde SQL Server, primero debe completar los pasos descritos en Configuración de Microsoft SQL Server para la ingesta en Azure Databricks.

Creación de una puerta de enlace y una canalización de ingesta

Advertencia

No detenga manualmente la pasarela de ingestión. La pasarela debe ejecutarse de forma continua para capturar los cambios antes de que los registros de cambios se trunquen en la base de datos de origen. Si se detiene la pasarela, los cambios pueden perderse debido a la retención de registros, lo que requeriría una actualización completa de todas las tablas afectadas. Detener y reiniciar la pasarela también provoca un nuevo aprovisionamiento de la máquina virtual, lo que aumenta el tiempo de arranque. Si necesita solucionar problemas del gateway, consulte Solucionar problemas de la ingesta de SQL Server o póngase en contacto con el soporte de Databricks.

Interfaz de usuario de Databricks

  1. En la barra lateral del área de trabajo de Azure Databricks, haga clic en Ingesta de datos.

  2. En la página Agregar datos , en Conectores de Databricks, haga clic en SQL Server.

  3. En la página Connection del Asistente para ingesta, seleccione la conexión que almacena las credenciales de acceso de SQL Server. Si tiene el privilegio CREATE CONNECTION en el metastore, puede hacer clic en Plus icon. Create connection para crear una conexión nueva con los detalles de autenticación que se indican en Create a SQL Server connection.

  4. Haga clic en Next.

  5. En la página Configuración de Ingesta, escriba un nombre único para la canalización de ingesta. Esta tubería mueve los datos del lugar de almacenaje provisional al destino.

  6. Seleccione un catálogo y un esquema en el que escribir registros de eventos. El registro de eventos contiene registros de auditoría, comprobaciones de calidad de datos, progreso de la canalización y errores. Si tiene los privilegios USE CATALOG y CREATE SCHEMA sobre el catálogo, puede hacer clic en Icono Plus. Crear esquema en el menú desplegable para crear un nuevo esquema.

  7. (Opcional) Establezca Actualización completa automática para todas las tablas en Activado. Cuando la actualización automática está activada, la canalización intenta corregir automáticamente problemas como eventos de limpieza de registros y ciertos tipos de evolución del esquema actualizando completamente la tabla afectada. Si el seguimiento del historial está habilitado, una actualización completa borra ese historial.

  8. Escriba un nombre único para la puerta de enlace de ingesta. La pasarela es un proceso que extrae los cambios del origen y los prepara para que los cargue el proceso de ingesta.

  9. Seleccione un catálogo y un esquema para la ubicación de almacenamiento provisional. En esta ubicación se crea un volumen para almacenar provisionalmente los datos extraídos. Si tiene los privilegios USE CATALOG y CREATE SCHEMA sobre el catálogo, puede hacer clic en Icono Plus. Crear esquema en el menú desplegable para crear un nuevo esquema.

  10. Haga clic en Crear canalización y continúe.

  11. En la página Origen , seleccione las tablas que se van a ingerir. Si selecciona tablas específicas, puede configurar las opciones de tabla:

    a) (Opcional) En la pestaña Configuración , especifique un nombre de destino para cada tabla ingerida. Esto resulta útil para diferenciar entre las tablas de destino al ingerir un objeto en el mismo esquema varias veces. Consulte Nombre de una tabla de destino.

    a) (Opcional) Cambie la configuración de seguimiento del historial predeterminada. Consulte Habilitación del seguimiento del historial (tipo 2 de SCD).

  12. Haga clic en Siguiente y, a continuación, haga clic en Guardar y continuar.

  13. En la página Destino , seleccione un catálogo y un esquema en el que cargar datos. Si tiene los privilegios USE CATALOG y CREATE SCHEMA sobre el catálogo, puede hacer clic en Icono Plus. Crear esquema en el menú desplegable para crear un nuevo esquema.

  14. Haga clic en Guardar y continuar.

  15. En la página Configuración de la base de datos , haga clic en Validar para confirmar que el origen está configurado correctamente para la ingesta de Azure Databricks. Se devuelven las configuraciones que faltan. Para conocer los pasos para resolverlo, haga clic en Completar configuración. A continuación, haga clic en Siguiente. Como alternativa, haga clic en Omitir validación.

  16. (Opcional) En la página Programaciones y notificaciones , haga clic en el icono Más. Crear programación. Establezca la frecuencia para actualizar las tablas de destino.

  17. (Opcional) Haga clic en el icono Más. Agregue una notificación para establecer notificaciones por correo electrónico para que la operación de canalización se complete correctamente o no y, a continuación, haga clic en Guardar y ejecutar canalización.

Agrupaciones de automatización declarativa

Antes de realizar la ingesta mediante paquetes de automatización declarativa, debe tener acceso a una conexión existente. Para obtener instrucciones, consulte Create a SQL Server connection.

El catálogo y el esquema de almacenamiento provisional pueden ser los mismos que el catálogo y el esquema de destino. El catálogo de ensayo no puede ser un catálogo externo. Indique la ubicación de almacenamiento provisional en la sección gateway_definition del archivo YAML de la canalización de paquetes.

La puerta de enlace de ingesta extrae la instantánea y cambia los datos de la base de datos de origen y los almacena en el volumen de almacenamiento provisional de Unity Catalog. Debe usar la puerta de enlace como una canalización continua. Esto ayuda a adaptarse a las directivas de retención del registro de cambios que tenga en la base de datos de origen.

La canalización de ingesta aplica la instantánea y cambia los datos del volumen de almacenamiento provisional en tablas de streaming de destino.

Las agrupaciones pueden contener definiciones de YAML de trabajos y tareas, se administran mediante la CLI de Databricks y se pueden compartir y ejecutar en diferentes áreas de trabajo de destino (como desarrollo, almacenamiento provisional y producción). Para obtener más información, consulte ¿Qué son los conjuntos de automatización declarativos?.

  1. Cree una agrupación mediante la CLI de Databricks:

    databricks bundle init
    
  2. Añada la configuración de su canal y de su trabajo al paquete. Consulte Ejemplos para obtener un ejemplo completo con todas las opciones disponibles.

  3. Implemente la canalización mediante la CLI de Databricks:

    databricks bundle deploy
    

Notebook de Databricks

Actualice la celda Configuration del cuaderno siguiente con la conexión de origen, el catálogo de destino, el esquema de destino y las tablas que se van a importar desde el origen.

Obtención del cuaderno

Terraform

Puede usar Terraform para implementar y administrar canalizaciones de ingesta de SQL Server. Para obtener un marco de ejemplo completo, incluidas las configuraciones de Terraform para crear puertas de enlace y canalizaciones de ingesta, consulte el repositorio de ejemplos de Terraform de Lakeflow Connect en GitHub.

Comprobación de la ingesta de datos correcta

La vista en lista en la página de detalles de la canalización muestra el número de registros procesados conforme se incorporan los datos. Estos números se actualizan automáticamente.

Comprobación de la replicación

Las columnas Upserted records y Deleted records no se muestran de forma predeterminada. Puede habilitarlas haciendo clic en el botón de configuración de columnas Icono de configuración de columnas y seleccionándolas.

Ejemplos

Use estos ejemplos para configurar la canalización.

Configuración de canalizaciones

Agrupaciones de automatización declarativa

El siguiente paquete define un canal de puerta de enlace, un canal de ingesta y un trabajo programado. Las opciones de salida comentadas muestran toda la configuración disponible. Actualice las variables secciones y targets con los detalles de origen y destino.

bundle:
  name: lakeflow-connect-sqlserver

# Variables parameterize the bundle for different environments and sources.
# Set values here, override per-target, or pass with: databricks bundle deploy -var="key=value"
variables:
  # The name of the Unity Catalog connection to your SQL Server instance.
  # This connection must already exist and be of type SQLSERVER.
  connection_name:
    description: 'Unity Catalog connection name for the SQL Server source'
  # The SQL Server database name to ingest from.
  # In Lakeflow Connect, this maps to source_catalog in the table/schema spec.
  source_database:
    description: 'SQL Server database name (maps to source_catalog in table specs)'
  # The SQL Server schema to ingest from (for example, "dbo", "sales").
  source_schema:
    description: 'SQL Server schema name to ingest from'
  # The Unity Catalog catalog where ingested Delta tables are created.
  dest_catalog:
    description: 'Destination Unity Catalog catalog for ingested tables'
  # The Unity Catalog schema where ingested Delta tables are created.
  dest_schema:
    description: 'Destination Unity Catalog schema for ingested tables'
  # The Unity Catalog catalog for the gateway's internal staging volume.
  # Can be the same as dest_catalog. Must not be a foreign catalog.
  staging_catalog:
    description: 'Catalog for gateway staging volume'
  # The Unity Catalog schema for the gateway's internal staging volume.
  staging_schema:
    description: 'Schema for gateway staging volume'

resources:
  pipelines:
    # --- Gateway pipeline ---
    # Extracts change data from SQL Server and stages it in a Unity Catalog
    # volume. Must run continuously to capture changes before change logs are
    # truncated in the source database.
    gw_pipeline:
      name: 'lfc-sqlserver-gateway-${bundle.target}'
      # Gateway pipelines must be continuous.
      continuous: true
      # "CURRENT" (stable) or "PREVIEW" (early access).
      channel: 'CURRENT'
      # (Optional) Associate with a budget policy for cost tracking.
      # budget_policy_id: "<policy-uuid>"
      # The gateway runs on classic compute. Cluster settings are managed
      # automatically. You can optionally customize the cluster:
      # clusters:
      #   - label: "default"
      #     autoscale:
      #       min_workers: 1
      #       max_workers: 4
      #     # node_type_id: "i3.xlarge"
      #     # Restrict the cluster to an approved cluster policy.
      #     # policy_id: "<cluster-policy-id>"
      catalog: ${var.staging_catalog}
      schema: ${var.staging_schema}
      gateway_definition:
        # (Required) Unity Catalog connection name (type SQLSERVER).
        connection_name: ${var.connection_name}
        # (Required) Catalog and schema for the staging volume.
        gateway_storage_catalog: ${var.staging_catalog}
        gateway_storage_schema: ${var.staging_schema}
        # (Optional) Custom staging volume name. If not set, the system
        # auto-generates: __databricks_ingestion_gateway_staging_data-<pipeline_id>
        # gateway_storage_name: "my_custom_staging_volume"

    # --- Ingestion pipeline ---
    # Reads staged data from the gateway and applies it to Delta tables.
    mi_pipeline:
      name: 'lfc-sqlserver-ingestion-${bundle.target}'
      # Continuous mode is not supported for the ingestion pipeline.
      # Use a scheduled job to trigger runs.
      continuous: false
      channel: 'CURRENT'
      # (Optional) Associate with a budget policy for cost tracking.
      # budget_policy_id: "<policy-uuid>"
      # The ingestion pipeline runs on serverless compute only.
      serverless: true
      # (Optional) Development mode for faster iteration (no retries).
      # development: true
      catalog: ${var.dest_catalog}
      schema: ${var.dest_schema}
      # (Optional) Email notifications for pipeline events.
      # notifications:
      #   - email_recipients:
      #       - "team@example.com"
      #     alerts:
      #       - "on-update-failure"
      #       - "on-update-fatal-failure"
      #       - "on-flow-failure"
      # (Optional) Run as a service principal for production.
      # run_as:
      #   service_principal_name: "my-service-principal"
      ingestion_definition:
        # (Required) References the gateway pipeline. The connection is
        # inherited from the gateway. Do not specify connection_name here.
        ingestion_gateway_id: ${resources.pipelines.gw_pipeline.id}
        # Pipeline-level table configuration defaults. These apply to all
        # tables unless overridden at the schema or table level.
        table_configuration:
          # SCD Type: How changes are applied to destination tables.
          #   SCD_TYPE_1: Overwrites rows with latest values (default).
          #   SCD_TYPE_2: Preserves history with __START_AT/__END_AT columns.
          #               Requires CDC on source. CT does not support SCD_TYPE_2.
          scd_type: 'SCD_TYPE_1'
          # (Optional) Auto full refresh policy. Triggers a snapshot when the
          # pipeline detects issues resolvable by re-reading all source data
          # (for example, CT/CDC retention window expired).
          # auto_full_refresh_policy:
          #   enabled: true
          #   min_interval_hours: 24
        # (Optional) Schedule automatic full refreshes.
        # full_refresh_window:
        #   start_hour: 2
        #   days_of_week:
        #     - "SUNDAY"
        #   time_zone_id: "America/Los_Angeles"
        objects:
          # Option 1: Schema-level ingestion. Ingests all tables from a source
          # schema. New tables added to the schema are picked up automatically.
          - schema:
              source_catalog: ${var.source_database}
              source_schema: ${var.source_schema}
              destination_catalog: ${var.dest_catalog}
              destination_schema: ${var.dest_schema}
              # (Optional) Override table_configuration for this schema.
              # table_configuration:
              #   scd_type: "SCD_TYPE_2"
          # Option 2: Table-level ingestion. Provides granular control.
          # Replace or combine with the schema-level spec.
          # - table:
          #     source_catalog: ${var.source_database}
          #     source_schema: ${var.source_schema}
          #     source_table: "customers"
          #     destination_catalog: ${var.dest_catalog}
          #     destination_schema: ${var.dest_schema}
          #     # (Optional) Rename the table at the destination.
          #     # destination_table: "customers_v2"
          #     table_configuration:
          #       scd_type: "SCD_TYPE_1"
          #       # Include only specific columns (mutually exclusive with exclude_columns).
          #       # include_columns:
          #       #   - "customer_id"
          #       #   - "first_name"
          #       #   - "email"
          #       # Exclude specific columns. All other columns are included.
          #       # exclude_columns:
          #       #   - "internal_notes"
          #       # Override the primary key used for change detection.
          #       # primary_keys:
          #       #   - "customer_id"
          #       # Logical ordering columns for change resolution.
          #       # sequence_by:
          #       #   - "updated_at"
          #       # Auto full refresh for this table.
          #       # auto_full_refresh_policy:
          #       #   enabled: true
          #       #   min_interval_hours: 48
      # (Optional) Grant additional users or groups access.
      # permissions:
      #   - user_name: "analyst@example.com"
      #     level: "CAN_VIEW"
      #   - group_name: "data-engineers"
      #     level: "CAN_RUN"

  # --- Scheduled job ---
  # Triggers the ingestion pipeline on a schedule.
  jobs:
    mi_schedule:
      name: 'lfc-sqlserver-ingestion-schedule-${bundle.target}'
      # Quartz cron syntax: "seconds minutes hours day month day-of-week"
      # Examples: "0 0 * * * ?" (hourly), "0 0 */4 * * ?" (every 4 hours)
      schedule:
        quartz_cron_expression: '0 */30 * * * ?'
        timezone_id: 'UTC'
      tasks:
        - task_key: 'run_ingestion'
          pipeline_task:
            pipeline_id: ${resources.pipelines.mi_pipeline.id}
      # email_notifications:
      #   on_failure:
      #     - "team@example.com"

# Deploy to different workspaces with: databricks bundle deploy -t <target>
targets:
  dev:
    default: true
    workspace:
      host: https://<workspace-url>.cloud.databricks.com
    variables:
      connection_name: '<sqlserver-connection>'
      source_database: '<database-name>'
      source_schema: 'dbo'
      dest_catalog: '<dest-catalog>'
      dest_schema: '<dest-schema>'
      staging_catalog: '<staging-catalog>'
      staging_schema: '<staging-schema>'

Notebook de Databricks

A continuación se muestra una sección de ejemplo Configuration de una especificación de canalización:

# The name of the UC connection with the credentials to access the source database
connection_name = "my_connection"

# The name of the UC catalog and schema to store the replicated tables
target_catalog_name = "main"
target_schema_name = "lakeflow_sqlserver_connector_cdc"

# The name of the UC catalog and schema to store the staging volume with intermediate
# CDC and snapshot data. Use the destination catalog/schema by default.
stg_catalog_name = target_catalog_name
stg_schema_name = target_schema_name

# The name of the Gateway pipeline to create
gateway_pipeline_name = "cdc_gateway"

# The name of the Ingestion pipeline to create
ingestion_pipeline_name = "cdc_ingestion"

# Construct the full list of tables to replicate.
# IMPORTANT: The letter case of catalog, schema, and table names must match exactly
# the case used in the source database system tables.
tables_to_replicate = replicate_full_db_schema("MY_DB", ["MY_DB_SCHEMA"])
# Append tables from additional schemas as needed:
#  + replicate_tables_from_db_schema("MY_DB", "MY_SCHEMA_2", ["table3", "table4"])

Patrones comunes

Para conocer las configuraciones avanzadas de canalización, consulte Patrones comunes para canalizaciones de ingesta administradas.

Pasos siguientes

Inicie, programe y establezca alertas en su flujo de trabajo. Consulte Tareas comunes de mantenimiento de canalización.

Recursos adicionales