Saltar al contenido

Ingesta batch

Menú

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

Guardar progreso

Ingesta batch, formatos, COPY INTO, JDBC y REST

Elige una vía de entrada reproducible para archivos, bases de datos y APIs.

Lectura pública
Detalles
Reto observable

Carga dos veces el mismo lote gobernado, conserva procedencia y cuarentena, y demuestra que el segundo run no duplica entidades ni archivos.

Al terminar podrás
  • Comparar cargas completas e incrementales
  • Usar COPY INTO de forma idempotente
  • Controlar formatos, compresión y metadatos de origen
Prerrequisitos
m07
Última revisión
25 ago 2026
Nivel
Associate
Ruta relacionada
core
Dominios blueprint
Data Ingestion and Loading · File formats
Estado
Revisión editorial interna
Fuentes principales
Ingestion · COPY INTO
Reportar un error
01
Modelo mental

Matriz de patrones de ingesta

La ingesta se elige por origen, volumen, frescura, cambios y gobierno.

Carga completa simplifica fuentes pequeñas; incremental reduce movimiento pero necesita cursor, archivos descubiertos o CDC.

Antes de seleccionar herramienta define reintento, borrados, esquema y límite de responsabilidad.

YAMLContrato de ingesta
source: orders_api
mode: incremental
cursor: updated_at
replay: true
deletes: tombstone

El cursor debe ser estable y soportar desempates.

¿Qué exige una carga incremental desde API?

Profundiza

Seleccionar ingesta significa traducir propiedades de la fuente a garantías de destino. Un conjunto finito de archivos puede cargarse en batch; una ruta creciente necesita descubrimiento incremental; una base operacional puede requerir snapshot, consultas JDBC o CDC; una API exige paginación y control de límites. Volumen, velocidad, esquema, orden, borrados, autenticación, replay y SLA determinan la herramienta. COPY INTO, Auto Loader y Lakeflow Connect no son sinónimos de pequeño, mediano y grande: ofrecen modelos de estado y operación distintos. La decisión autosuficiente identifica primero qué constituye un dato nuevo y cómo se demuestra que no se perdió ni se procesó dos veces.

Unidad incremental

Elemento cuya identidad permite reconocer progreso, como archivo, offset o versión de cambio.

Define cómo reanudar y evitar reprocesamiento accidental.
Replay

Capacidad de volver a procesar una entrada durable desde un punto conocido.

Es esencial para corregir lógica sin depender de la disponibilidad de la fuente.
CDC

Captura de inserciones, actualizaciones y borrados de una fuente con orden o secuencia asociados.

Conserva cambios que un snapshot periódico podría perder.
Resumen

Puntos clave

  • Full e incremental tienen trade-offs
  • Idempotencia es requisito
  • El origen condiciona el patrón

Evita

  • Incremental sin cursor fiable
  • Full load que borra historia útil
02
Implementación

CSV, JSON, Parquet, Avro, ORC y binary

Formato y compresión afectan esquema, pushdown, tamaño y coste.

Parquet y Delta son columnares; JSON/CSV requieren parseo y contratos; XML, text y binary cubren fuentes no tabulares.

Bronze puede conservar payload y metadatos, pero debe fijar encoding, delimitador y tratamiento de registros corruptos.

PySparkLectura con esquema
raw = (spark.read.schema(order_schema)
  .option("mode", "PERMISSIVE")
  .json("/Volumes/main/landing/orders/incoming"))

Mide la columna de registros corruptos si la configuras.

¿Qué añade Delta a Parquet?

Profundiza

El formato determina qué información estructural puede usar el motor antes de leer cada valor. CSV es texto sin tipos embebidos y exige delimitador, escape, locale y esquema externo. JSON conserva jerarquía, pero su verbosidad y variabilidad complican análisis masivo. Parquet es columnar, tipado y comprimido; permite poda de columnas y pushdown de filtros según metadatos. Delta usa Parquet para datos y añade transaction log, versiones y DML. La compresión reduce bytes a costa de CPU con algoritmos distintos. Elegir no es solo comparar tamaño: se valora interoperabilidad, evolución, patrones de lectura, calidad y necesidad de transacciones.

Formato columnar

Organización física que almacena valores de una misma columna juntos y conserva metadatos por grupos.

Favorece poda, compresión y análisis de subconjuntos de columnas.
Predicate pushdown

Aplicación de filtros en el lector para evitar materializar datos que no pueden cumplirlos.

Reduce I/O cuando el formato y la expresión lo permiten.
Codec

Algoritmo que comprime y descomprime bloques de datos con trade-offs de tamaño y CPU.

Afecta coste de storage, red y tiempo de proceso.
Resumen

Puntos clave

  • Parquet es columnar
  • Delta añade transacciones
  • Semiestructurado exige esquema

Evita

  • Inferir CSV en cada run
  • Confundir Parquet con Delta
03
Operación

COPY INTO e historial de archivos

COPY INTO carga archivos nuevos de object storage de forma reintentable.

Mantiene historial de archivos ya procesados para evitar recarga en ejecuciones normales y permite opciones de formato y transformación limitada.

Es apropiado para ingesta incremental SQL sencilla; Auto Loader escala mejor para descubrimiento continuo y evolución avanzada.

SQLCOPY INTO gobernado
COPY INTO main.bronze.orders
FROM '/Volumes/main/landing/orders/incoming'
FILEFORMAT = JSON
FORMAT_OPTIONS ('inferSchema' = 'false');

Crea previamente la tabla con esquema explícito.

¿COPY INTO deduplica por order_id?

Profundiza

COPY INTO es una carga SQL incremental de archivos que registra qué entradas se procesaron para una tabla y permite reejecutar sin cargar normalmente los mismos archivos. Es apropiado para lotes simples o periódicos desde object storage cuando no se necesita el control streaming de Auto Loader. La idempotencia se basa en identidad de archivos y estado de carga, no en la clave de negocio; si el productor publica el mismo contenido con otro nombre, puede volver a entrar. FILEFORMAT, FORMAT_OPTIONS y COPY_OPTIONS separan cómo leer de cómo cargar. Validación, esquema y errores deben configurarse conscientemente antes de convertir COPY INTO en una promesa de exactamente una fila por evento.

COPY INTO

Comando SQL que carga archivos nuevos en una tabla y mantiene estado de archivos procesados.

Ofrece una ruta batch simple y reintentable para object storage.
FILEFORMAT

Declaración del formato de origen que selecciona el lector de los archivos.

Debe corresponder a la representación física y no sustituye opciones de parseo.
Identidad de archivo

Metadatos usados para distinguir entradas ya procesadas de candidatas nuevas.

Aclara por qué renombrar contenido puede provocar una nueva carga.
Resumen

Puntos clave

  • Rastrea archivos
  • Es SQL declarativo
  • No sustituye CDC de filas

Evita

  • Renombrar contenido y esperar deduplicación de filas
  • Usar force sin entender recarga
04
Diagnóstico

JDBC, ODBC y REST

JDBC/ODBC consultan sistemas tabulares; REST requiere paginación, límites y persistencia durable.

Una lectura JDBC paralela necesita partitionColumn y bounds coherentes; demasiadas conexiones pueden dañar el origen.

Una API REST debe manejar rate limits, retries con backoff, tokens, paginación y checkpoints antes de publicar en UC.

PySparkJDBC particionado
df = (spark.read.format("jdbc")
 .option("url", jdbc_url).option("dbtable", "orders")
 .option("partitionColumn", "order_id").option("lowerBound", 1)
 .option("upperBound", 1000000).option("numPartitions", 8).load())

Los bounds dividen lectura, no filtran filas.

¿numPartitions en JDBC también limita conexiones?

Profundiza

JDBC y ODBC presentan un modelo tabular de consulta; REST presenta recursos y respuestas paginadas bajo un contrato HTTP. En JDBC, la partición de lectura determina paralelismo y puede sobrecargar la base si se eligen demasiadas conexiones. Predicados y columnas deben empujarse cuando sea posible. En REST, cada página, cursor, límite de tasa, retry y error parcial forma parte del estado. Ninguna interfaz debe usarse como almacenamiento intermedio invisible: las respuestas se persisten durably antes o junto con su transformación. La autenticación se resuelve mediante secretos o conexiones gobernadas, y los logs nunca deben exponer tokens o datos sensibles.

Pushdown

Ejecución de proyecciones o filtros en el sistema fuente en lugar de transferir todos los datos.

Reduce red y carga Spark, aunque debe verificarse en el plan.
Cursor

Token opaco emitido por una API para continuar una secuencia paginada desde una posición.

Debe persistirse con el lote confirmado para reanudar sin huecos.
Rate limit

Política del servicio que restringe solicitudes en una ventana y comunica cómo reintentar.

Obliga a diseñar backoff y throughput sin provocar bloqueos.
Resumen

Puntos clave

  • Protege el sistema fuente
  • Persistencia antes de transformar
  • Credenciales fuera del código

Evita

  • Abrir cientos de conexiones
  • Loggear tokens REST
05
Decisión de diseño

Metadatos, errores y cuarentena

Una landing gobernada conserva origen, lote y evidencia para replay.

Volumes ofrecen rutas de archivos bajo Unity Catalog; tablas Bronze organizan registros con metadatos de ingestión.

La idempotencia se demuestra repitiendo el mismo lote y comparando filas, claves y commits, no solo ausencia de error.

SQLMetadatos de origen
SELECT *, _metadata.file_path, _metadata.file_modification_time
FROM read_files('/Volumes/main/landing/orders', format => 'json');

Persiste los metadatos necesarios para investigar duplicados.

¿Qué evidencia prueba reintento seguro?

Profundiza

Landing es una frontera durable entre recepción y procesamiento. Conserva los bytes entregados, identidad del objeto, tiempo, productor, lote y, cuando procede, checksum o metadatos de transporte. No es un vertedero anónimo ni necesariamente una tabla para analistas. Su propósito es demostrar qué llegó y permitir replay cuando cambia la lógica Silver o falla un proceso posterior. Unity Catalog Volumes y external locations gobiernan rutas y evitan credenciales en código. La estructura de carpetas puede ayudar a operación, pero no debe sustituir un manifiesto y metadatos consultables. Retención, cifrado, clasificación y acceso se definen según sensibilidad y capacidad de reextracción.

Landing zone

Área gobernada donde se conserva la representación recibida antes de transformaciones irreversibles.

Desacopla disponibilidad de la fuente y permite replay verificable.
Manifiesto

Registro de objetos esperados y observados con tamaño, checksum, lote y estado.

Permite reconciliar transferencias y detectar archivos ausentes o alterados.
Provenance

Información que relaciona un dato con productor, objeto, instante y proceso de origen.

Sustenta auditoría y localización de errores a lo largo del pipeline.
Resumen

Puntos clave

  • Volumes gobiernan archivos
  • Bronze conserva metadata
  • Replay debe probarse

Evita

  • Guardar landing en ruta personal
  • Declarar idempotencia sin segundo run

Vista de lectura · sin ejecución

File operations and ELT notebooks

databricks-notebooks · commit 51e8e4b

notebooks/file-operations-python.ipynb

Módulo 08

Contenido del módulo