Saltar al contenido

Ingesta managed

Menú

Puedes leer sin crear un espacio. Créalo solo cuando quieras guardar.

Guardar progreso

Auto Loader y Lakeflow Connect

Selecciona entre descubrimiento de archivos, conectores gestionados y alternativas de integración.

Lectura pública
Detalles
Reto observable

Implementa una ingesta incremental con estado y evolución controlados, y justifica Auto Loader o Lakeflow Connect mediante volumen, frescura y operación.

Al terminar podrás
  • Configurar schemaLocation y checkpointLocation
  • Gestionar evolución y rescued data
  • Elegir Lakeflow Connect, Auto Loader o partner connector
Prerrequisitos
m08
Última revisión
25 ago 2026
Nivel
Associate + Professional
Ruta relacionada
core
Dominios blueprint
Data Ingestion and Loading · Lakeflow Connect
Estado
Revisión editorial interna
Fuentes principales
Auto Loader · Lakeflow Connect
Reportar un error
01
Modelo mental

Auto Loader y cloudFiles

Auto Loader descubre archivos incrementalmente mediante cloudFiles.

readStream.format('cloudFiles') conserva estado de archivos y escala sin listar todo repetidamente.

Necesita formato de origen, checkpoint durable y, según el caso, schemaLocation o esquema proporcionado.

PySparkAuto Loader
stream = (spark.readStream.format("cloudFiles")
 .option("cloudFiles.format", "json")
 .option("cloudFiles.schemaLocation", schema_path)
 .load(landing_path))

Usa rutas distintas de schema y checkpoint por workload.

¿Qué opción indica formato de los archivos?

Profundiza

Auto Loader es una fuente incremental de Structured Streaming llamada cloudFiles que transforma un directorio creciente en una secuencia de archivos descubiertos y procesados con estado. No observa filas nuevas dentro de un archivo mutable ni convierte automáticamente duplicados de negocio en eventos únicos. Su unidad de progreso es el archivo identificado durante descubrimiento. Un checkpoint conserva avance de streaming y el schema location mantiene la historia usada para inferencia y evolución; cumplen funciones distintas y ambos pertenecen a una única carga lógica. Auto Loader puede ejecutarse directamente con Structured Streaming o dentro de Spark Declarative Pipelines, el framework actual que sustenta la oferta gestionada Lakeflow pipelines, denominación todavía visible en materiales anteriores.

cloudFiles

Nombre de la fuente de Structured Streaming de Auto Loader que descubre y procesa incrementalmente archivos en almacenamiento cloud.

Distingue el mecanismo de descubrimiento del formato real, como JSON, CSV o Parquet, contenido en cada archivo.
Checkpoint

Directorio durable que conserva offsets, commits y estado necesario para reanudar una consulta de streaming identificable.

Permite continuar después de un fallo sin olvidar el progreso confirmado por esa carga lógica concreta.
Schema location

Ubicación durable donde Auto Loader registra el esquema inferido y su evolución a través de lotes sucesivos.

Separa la historia del contrato estructural del progreso de datos guardado en el checkpoint.
Resumen

Puntos clave

  • cloudFiles es la fuente
  • Checkpoint conserva progreso
  • Esquema debe gobernarse

Evita

  • Compartir checkpoint
  • Borrarlo para reintentar
02
Implementación

Directory listing y file notification

Directory listing simplifica; file notification reduce listados a gran escala.

El modo de descubrimiento se elige por volumen y configuración cloud; ambos necesitan permisos correctos.

rescuedDataColumn conserva campos que no encajan y permite estudiar evolución sin perder payload.

PySparkRescatar cambios
stream = (spark.readStream.format("cloudFiles")
 .option("cloudFiles.format", "json")
 .option("rescuedDataColumn", "_rescued_data")
 .load(landing_path))

Alerta si crece _rescued_data.

¿Para qué sirve _rescued_data?

Profundiza

Directory listing y file notification son estrategias para encontrar archivos, no formatos de entrada ni garantías de entrega final. El listado consulta la jerarquía de object storage y compara resultados con estado; es sencillo y funciona sin infraestructura de eventos, pero su coste de enumeración crece con rutas enormes. Las notificaciones aprovechan eventos cloud y colas para señalar objetos nuevos, reduciendo listados a escala, a cambio de permisos y componentes adicionales. Auto Loader administra el estado de archivos en ambos casos. La elección depende de número total de objetos, tasa de llegada, SLA, restricciones de red y capacidad operativa. Un evento de notificación no sustituye la validación del objeto ni el commit del sink.

Directory listing

Descubrimiento que enumera objetos bajo una ruta y compara sus metadatos con el estado incremental previamente conservado.

Ofrece menor complejidad inicial, pero puede convertirse en el coste dominante cuando la jerarquía acumula millones de archivos.
File notification

Descubrimiento basado en eventos cloud que anuncian la creación de objetos mediante una cola o servicio equivalente.

Reduce enumeraciones masivas, aunque añade configuración de identidad, eventos y operación de la infraestructura asociada.
File event

Señal potencialmente duplicada o desordenada que identifica un objeto candidato, no una confirmación de procesamiento completo.

Obliga a conservar estado y verificar el archivo antes de considerar que los datos llegaron correctamente al destino.
Resumen

Puntos clave

  • Listing requiere menos infraestructura
  • Notifications escalan descubrimiento
  • Rescued data preserva anomalías

Evita

  • Ignorar rescued data
  • Cambiar esquema sin probar checkpoint
03
Operación

Inferencia, hints y evolución de esquema

Lakeflow Connect combina conectores managed/standard y conectores basados en consulta para bases sin CDC.

Managed connectors reducen operación y crean pipelines gobernados para fuentes compatibles. Standard connectors amplían fuentes desde Spark o pipelines. Un query-based connector consulta directamente la base mediante una Unity Catalog connection y escribe streaming tables sin ingestion gateway ni staging volume.

La ingesta incremental basada en consulta se ejecuta según programación y guarda un high-water mark de una única columna cursor monotónica; cada run lee valores mayores que el anterior. Un cursor NULL no se ingiere y un ID sólo sirve si las filas son append-only. No presentes este patrón como CDC de log ni como captura automática de cualquier delete.

YAMLContrato de conector basado en consulta
source: postgres.orders
connector: query_based
cursor: updated_at
schedule: '0 */15 * * * ?'
target_catalog: main

Comprueba monotonicidad, nulos y actualizaciones que no modifican el cursor antes de aceptar el diseño.

¿Qué condición debe cumplir el cursor de un conector basado en consulta?

Profundiza

Lakeflow Connect organiza opciones de ingesta desde conectores muy gestionados hasta interfaces más personalizables. Los managed connectors encapsulan autenticación, lectura incremental, serverless compute y destino gobernado para fuentes soportadas. Los standard connectors exponen capacidades desde SQL, Python o APIs de streaming cuando se necesita controlar transformaciones, opciones o fuentes. La documentación actual recomienda comenzar por la capa más gestionada que cumpla requisitos y descender solo ante una limitación concreta. El framework de transformaciones se denomina Spark Declarative Pipelines; Lakeflow pipelines es la oferta gestionada que lo extiende y aún aparece como nombre abreviado en el blueprint y material heredado. Ningún conector elimina la necesidad de validar semántica CDC y SLA.

Managed connector

Integración configurada que administra lectura incremental, compute y publicación gobernada para una fuente empresarial explícitamente soportada.

Reduce código y operación cuando sus capacidades, región y semántica coinciden con los requisitos reales de ingestión.
Standard connector

Interfaz de ingesta disponible desde código o SQL que concede mayor control sobre opciones y comportamiento del pipeline.

Permite cubrir fuentes o requisitos especiales, pero hace responsable al equipo del estado, pruebas y recuperación.
Connection

Securable de Unity Catalog que representa conectividad y autenticación administrada hacia un sistema de datos externo.

Evita repartir credenciales entre notebooks y permite conceder uso del acceso externo mediante gobierno centralizado.
Resumen

Puntos clave

  • Managed minimiza operación
  • Query-based usa cursor y high-water mark
  • Cursor y semántica dependen de la fuente

Evita

  • Prometer CDC sin verificar conector
  • Elegir un ID como cursor si las filas existentes pueden actualizarse
04
Diagnóstico

Lakeflow Connect standard y managed

La matriz de decisión compara volumen, frecuencia, tipos y gobierno.

COPY INTO encaja en lotes de archivos simples; Auto Loader en archivos incrementales escalables; Connect en SaaS y bases soportadas.

JDBC/REST siguen siendo válidos cuando el origen o contrato no está cubierto, pero trasladan más responsabilidad al equipo.

PythonRegla legible
choice = "autoloader" if source == "files" and continuous else "copy_into"
if managed_connector_available: choice = "lakeflow_connect"

Completa con SLA, CDC y limitaciones.

¿Qué elegir para JSON continuo en S3/ADLS/GCS?

Profundiza

Una matriz de ingesta obliga a comparar soluciones sobre el mismo conjunto de requisitos en vez de elegir por familiaridad. Las filas representan fuentes o casos; las columnas capturan volumen, frecuencia, latencia, mutabilidad, deletes, esquema, autenticación, replay, backfill, disponibilidad regional, coste y operación. Cada opción se puntúa con evidencia y se documenta una condición que invalidaría la decisión. COPY INTO puede ganar para batch SQL simple; Auto Loader para archivos crecientes; un managed connector para una aplicación soportada; CDC para cambios de base; REST personalizado para una API sin integración. La matriz no sustituye una prueba: identifica qué hipótesis debe validar el laboratorio.

Criterio eliminatorio

Requisito obligatorio cuya ausencia descarta una alternativa antes de comparar ventajas secundarias como comodidad o precio.

Evita seleccionar una herramienta popular que no puede representar la semántica indispensable de la fuente.
Benchmark representativo

Prueba con distribución, escala, concurrencia y fallos semejantes a la operación esperada, acompañada de métricas reproducibles.

Convierte la decisión arquitectónica en evidencia y revela límites que una tabla teórica no muestra.
Carga operacional

Tiempo, conocimiento y acciones humanas necesarios para configurar, vigilar, reparar y evolucionar una solución de ingesta.

Forma parte del coste total y favorece servicios gestionados cuando satisfacen el contrato técnico.
Resumen

Puntos clave

  • No existe una herramienta universal
  • Gobierno forma parte de la elección
  • Operación es un coste

Evita

  • Usar streaming para un lote mensual
  • Construir REST para fuente ya soportada
05
Decisión de diseño

Matriz de decisión por volumen, frescura y gobierno

Evolución y estado deben aislarse por pipeline y ambiente.

Checkpoint codifica offsets y estado compatible con la consulta; cambios de fuente, claves o estado pueden exigir nueva ubicación y backfill.

Dev y prod no comparten schemaLocation ni checkpoint. La tabla objetivo se gobierna en UC y la identidad tiene mínimo privilegio.

PySparkEscritura durable
(stream.writeStream
 .option("checkpointLocation", checkpoint_path)
 .trigger(availableNow=True)
 .toTable("main.bronze.orders"))

availableNow procesa lo disponible y se detiene conservando progreso.

¿Qué ocurre al borrar checkpoint?

Profundiza

El checkpoint y el schema location son parte de la identidad de un pipeline, igual que su código y destino. Dev, test, prod, backfill y una nueva lógica no deben compartirlos accidentalmente. El checkpoint afirma qué unidades se confirmaron y puede contener estado de operadores; reutilizarlo con otra consulta puede omitir datos, rechazar cambios incompatibles o mezclar semánticas. El schema location conserva evolución inferida y también debe corresponder a una fuente y contrato. Para reprocesar, se crea una identidad nueva y un destino o estrategia idempotente; borrar un checkpoint productivo no es un botón de retry. Las rutas se nombran, gobiernan y retienen con la misma disciplina que una tabla.

Identidad de pipeline

Conjunto estable de fuente, transformación, destino y estado durable que define una carga incremental concreta.

Impide tratar checkpoints como archivos intercambiables entre ambientes o lógicas con significados diferentes.
Estado de operador

Datos intermedios persistidos para ventanas, agregaciones, deduplicación u otras operaciones que dependen de eventos anteriores.

Hace que ciertos cambios de código sean incompatibles con un checkpoint existente y requieran migración planificada.
Backfill aislado

Reprocesamiento histórico con estado y ámbito propios que integra resultados mediante una escritura controlada e idempotente.

Evita alterar progreso productivo o duplicar datos mientras continúa la ingesta ordinaria.
Resumen

Puntos clave

  • Checkpoint es parte del contrato
  • Entornos aislados
  • Backfill se planifica

Evita

  • Reutilizar checkpoint con otra query
  • Guardar checkpoint en ruta temporal

Vista de lectura · sin ejecución

crdb_to_dbx

crdb_to_dbx · commit 042cb96

crdb_to_dbx/cockroachdb-cdc-tutorial.ipynb

Módulo 09

Contenido del módulo