Nota:
El acceso a esta página requiere autorización. Puede intentar iniciar sesión o cambiar directorios.
El acceso a esta página requiere autorización. Puede intentar cambiar los directorios.
Para obtener procedimientos recomendados completos sobre la configuración de Auto Loader, incluida la selección del modo de detección de archivos, la administración de esquemas y el control de la calidad de los datos, consulte Procedimientos recomendados del cargador automático.
Databricks recomienda usar Auto Loader en canalizaciones de Lakeflow para la ingesta incremental de datos. Las canalizaciones de Lakeflow amplían la funcionalidad en Apache Spark Structured Streaming y permiten escribir solo unas pocas líneas de Python declarativas o SQL para implementar una canalización de datos de calidad de producción con:
- Infraestructura informática con escalado automático para ahorrar costes:Optimiza la utilización del clúster de pipelines de Lakeflow con escalado automático
- Comprobaciones de calidad de datos con expectativas:administración de la calidad de los datos con expectativas de canalización
- Control automático de la evolución del esquema:Configurar la inferencia de esquemas y la evolución en Auto Loader
- Supervisión a través de métricas en el registro de eventos:Registro de eventos de canalización
Databricks también recomienda seguir los procedimientos recomendados de streaming para ejecutar Auto Loader en producción. Consulte Consideraciones de producción para Structured Streaming.
Note
Las canalizaciones de Lakeflow son la forma recomendada de ejecutar Auto Loader para la mayoría de los procesos de ingesta en producción. Si la carga de trabajo no tiene requisitos de baja latencia y su prioridad minimiza el costo de proceso, en su lugar puede programar Auto Loader como un trabajo por lotes desencadenado que use Trigger.AvailableNow. Consulte Consideraciones sobre los costos.
Supervisión de cargador automático
En las secciones siguientes se describe cómo supervisar el cargador automático en producción, incluidas las métricas, los registros, las alertas y los flujos de trabajo de solución de problemas comunes. Para obtener una referencia completa sobre patrones de paneles, análisis de latencia y detección de la deriva del esquema, consulte Supervisar y observar Auto Loader.
Consulta de archivos descubiertos por el cargador automático
El cargador automático proporciona una API de SQL para inspeccionar el estado de una secuencia. Con la función cloud_files_state, puede encontrar metadatos sobre los archivos detectados por una secuencia del cargador automático. Consulte cloud_files_state y proporcione la ubicación del punto de control asociada a un flujo de Auto Loader.
Note
La opción cloud_files_state está disponible en Databricks Runtime 11.3 LTS y versiones posteriores.
SELECT * FROM cloud_files_state('path/to/checkpoint');
Escuchar actualizaciones de secuencias
Para monitorear más a fondo los flujos de Auto Loader, Databricks recomienda usar la interfaz de Streaming Query Listener de Apache Spark. Vea Supervisión de consultas de Structured Streaming en Azure Databricks.
El cargador automático notifica métricas al agente de escucha de consultas de streaming en cada lote. Puede ver cuántos archivos existen en el trabajo pendiente y el tamaño del trabajo pendiente en la numFilesOutstanding y numBytesOutstandinglas métricas en la pestaña Datos sin procesar en el panel de progreso de la consulta de streaming:
{
"sources": [
{
"description": "CloudFilesSource[/path/to/source]",
"metrics": {
"numFilesOutstanding": "238",
"numBytesOutstanding": "163939124006"
}
}
]
}
Al usar el modo de notificación de archivos en Databricks Runtime 10.4 LTS y versiones posteriores, las métricas también incluyen el número aproximado de eventos de archivo en la cola en la nube: approximateQueueSize para AWS y Azure.
Consideraciones sobre los costos
Al ejecutar Auto Loader, los principales orígenes de costo son los recursos de proceso y la detección de archivos.
Si su carga de trabajo no tiene requisitos de baja latencia, puede reducir los costes de computación usando Lakeflow Jobs para programar Auto Loader en trabajos por lotes mediante Trigger.AvailableNow en lugar de ejecutarlo continuamente. Vea Configurar intervalos del desencadenador de Structured Streaming. Estos trabajos por lotes se pueden desencadenar mediante desencadenadores de llegada de archivos para reducir aún más la latencia entre la llegada y el procesamiento de archivos.
Los costos de detección de archivos pueden presentarse en forma de operaciones de LIST en las cuentas de almacenamiento en el modo de lista de directorios y las solicitudes de API en el servicio de suscripción y en cola en modo de notificación de archivos. Los desencadenadores continuos, como Trigger.ProcessingTime , por ejemplo, son especialmente costosos en el modo de lista de directorios, ya que Auto Loader enumera continuamente todo el directorio para buscar nuevos archivos. Si la carga de trabajo requiere desencadenadores continuos, Databricks recomienda elegir un modo de detección de archivos en función de los requisitos de latencia:
- Baja latencia y simplicidad: Use Auto Loader con eventos de archivos. Los eventos de archivo solo requieren una cola por cubo y usan la detección incremental en ejecuciones posteriores. Para obtener más información, consulte Auto Loader: introducción a los eventos de archivo.
- Aplicaciones muy sensibles a la latencia: use el modo de notificación de archivos clásico. El modo clásico lee directamente desde la cola en la nube sin el paso adicional por la caché que introducen los eventos de archivo. En este modo, puede etiquetar los recursos creados por Auto Loader para realizar un seguimiento de los costos mediante etiquetas de recursos. Para obtener más información, consulte Notificación de archivos.
Retención de datos de origen
Note
Disponible en Databricks Runtime 16.4 LTS y versiones posteriores.
A medida que los archivos se acumulan en el directorio de origen, los costos de almacenamiento aumentan y la detección de archivos se ralentiza, especialmente en el modo de lista de directorios. Auto Loader proporciona la opción de administrar automáticamente la cloudFiles.cleanSource retención de archivos mediante el archivado o la eliminación de archivos una vez procesados.
Archivado de archivos en el directorio de origen para reducir los costos
Warning
- Al establecer
cloudFiles.cleanSourcese eliminan o mueven archivos en el directorio de origen. - Si usa
foreachBatchpara el procesamiento de datos, sus archivos se convierten en candidatos para moverse o eliminarse tan pronto como su operaciónforeachBatchse completa con éxito, incluso si la operación solo consume un subconjunto de los archivos del lote.
Databricks recomienda usar Auto Loader con eventos de archivo para reducir los costos de detección. Esto también reduce los costos de proceso porque la detección es incremental.
Si no puede usar eventos de archivo y debe usar la lista de directorios para detectar archivos, puede usar la cloudFiles.cleanSource opción de archivar o eliminar archivos automáticamente después de que Auto Loader los procese para reducir los costos de detección. Dado que Auto Loader limpia los archivos del directorio de origen después del procesamiento, es necesario que se muestren menos archivos durante la detección.
Al usar cloudFiles.cleanSource con la MOVE opción , tenga en cuenta los siguientes requisitos:
- Tanto el directorio de origen como el directorio de destino del traslado deben encontrarse en la misma ubicación externa, volumen o montaje DBFS. No se admiten movimientos entre cubos y contenedores, lo que resulta en un error.
- El destino de movimiento puede ser una ruta de acceso de volumen (por ejemplo,
/Volumes/my_catalog/my_schema/my_volume/archive/). - Si el directorio de origen y destino se encuentran en la misma ubicación externa, no deben tener directorios del mismo nivel que contengan almacenamiento administrado (por ejemplo, un volumen administrado o catálogo). En estos casos, el cargador automático no puede obtener los permisos necesarios para escribir en el directorio de destino.
Databricks recomienda usar esta opción cuando:
- El directorio de origen acumula un gran número de archivos a lo largo del tiempo.
- Debe conservar los archivos procesados para el cumplimiento normativo o la auditoría (establecer
cloudFiles.cleanSourceaMOVE). - Quiere reducir los costos de almacenamiento eliminando los archivos tras la ingestión (configurar
cloudFiles.cleanSourceenDELETE). Al usar el modoDELETE, Databricks recomienda habilitar el versionado en el cubo para que Auto Loader actúe como eliminaciones lógicas y estén disponibles en caso de una configuración incorrecta. Además, Databricks recomienda configurar directivas de ciclo de vida de la nube para purgar versiones antiguas y eliminadas temporalmente después de un período de gracia especificado (por ejemplo, 60 o 90 días) en función de los requisitos de recuperación.
Para la referencia completa sobre cleanSource opciones y sus valores predeterminados, consulta Limpiar archivos procesados con Auto Loader.
Mover archivos procesados a una ruta de almacenamiento en frío
En el ejemplo siguiente se configura el cargador automático para mover los archivos procesados a un directorio de archivo dentro del mismo cubo después de 14 días. Puede aplicar una directiva de ciclo de vida de la nube en la ruta de acceso de archivo para realizar la transición de archivos a niveles de almacenamiento más baratos (por ejemplo, AWS S3 Glacier, Azure Cool/Archive o GCS Coldline/Archive).
Python
# Step 1: Configure Auto Loader to move processed files to an archive path.
checkpoint = "/Volumes/my_catalog/my_schema/my_volume/checkpoints/ingest_stream"
archive_path = "s3://my-bucket/archive/landing/"
df = (spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.cleanSource", "MOVE")
.option("cloudFiles.cleanSource.moveDestination", archive_path)
.option("cloudFiles.cleanSource.retentionDuration", "14 days")
.option("cloudFiles.schemaLocation", checkpoint)
.load("s3://my-bucket/landing/")
)
# Step 2: Write to a Delta table.
(df.writeStream
.option("checkpointLocation", checkpoint)
.trigger(availableNow=True)
.toTable("my_catalog.my_schema.raw_events")
)
# Step 3 (outside Databricks): Set up a cloud lifecycle policy on the
# archive path to transition files to cold storage after a grace period.
# For example, in AWS you can configure an S3 Lifecycle rule to move
# objects under s3://my-bucket/archive/landing/ to S3 Glacier after
# 30 days.
SQL
-- Step 1: Configure Auto Loader to move processed files to an archive path
-- using a Lakeflow Declarative Pipeline.
CREATE OR REFRESH STREAMING TABLE raw_events
AS SELECT * FROM STREAM read_files(
's3://my-bucket/landing/',
format => 'json',
cleanSource => 'MOVE',
cleanSourceMoveDestination => 's3://my-bucket/archive/landing/',
cleanSourceRetentionDuration => '14 days'
);
-- Step 2 (outside Databricks): Set up a cloud lifecycle policy on the
-- archive path to transition files to cold storage.
-- For example, in AWS configure an S3 Lifecycle rule to move objects
-- under s3://my-bucket/archive/landing/ to S3 Glacier after 30 days.
Uso de Trigger.AvailableNow y limitación de velocidad
Note
Disponible en Databricks Runtime 10.4 LTS y versiones posteriores.
El cargador automático se puede programar para que se ejecute en Lakeflow Jobs como un trabajo por lotes usando Trigger.AvailableNow. El AvailableNow desencadenador indica al cargador automático que procese todos los archivos que llegaron antes de la hora de inicio de la consulta. Los nuevos archivos que llegan después de que se inicie el flujo se omiten hasta el siguiente desencadenador.
Con Trigger.AvailableNow, la detección de archivos se lleva a cabo de forma asincrónica con el procesamiento de datos y los datos se pueden procesar en varios micro lotes con limitación de velocidad. El cargador automático procesa de manera predeterminada un máximo de 1000 archivos por micro lote. Puede configurar cloudFiles.maxFilesPerTrigger y cloudFiles.maxBytesPerTrigger para configurar cuántos archivos o cuántos bytes se deben procesar en un micro lote. El límite de archivos es un límite máximo, pero el límite de bytes es un límite flexible, lo que significa que se pueden procesar más bytes que el maxBytesPerTrigger proporcionado. Cuando ambas opciones se proporcionan juntas, el cargador automático procesa tantos archivos como sean necesarios para alcanzar uno de los límites.
Ubicación del punto de control
La ubicación del punto de control se usa para almacenar el estado y la información de progreso de la secuencia. Databricks recomienda establecer la ubicación del punto de control en una ubicación sin una política de ciclo de vida de objetos en la nube. Si los archivos de la ubicación del punto de control se limpian según la directiva, el estado de la secuencia está dañado. Si esto ocurre, debe reiniciar la secuencia desde el principio.
Seguimiento de eventos relacionados con archivos
El cargador automático realiza un seguimiento de los archivos detectados en la ubicación del punto de comprobación mediante RocksDB para proporcionar garantías de ingesta exactamente una vez. Para flujos de ingesta de gran volumen o de larga duración, Databricks recomienda actualizar a Databricks Runtime 15.4 LTS o superior. En estas versiones, Auto Loader no espera a que se descargue todo el estado de RocksDB antes de que se inicie la secuencia, lo que puede acelerar el tiempo de inicio de la secuencia.
Si desea evitar que los estados del archivo crezcan sin límites, también puede considerar la posibilidad de usar la cloudFiles.maxFileAge opción de expirar los eventos de archivo anteriores a una determinada edad. El valor mínimo que se puede establecer para cloudFiles.maxFileAge es "14 days". Las eliminaciones en RocksDB aparecen como entradas de marcador de exclusión. Por lo tanto, es posible que vea un aumento temporal en el uso del almacenamiento a medida que los eventos caducan antes de que comience a estabilizarse.
Warning
cloudFiles.maxFileAge se ofrece como mecanismo de control de costes para conjuntos de datos de gran volumen. Un ajuste demasiado agresivo de cloudFiles.maxFileAge puede causar problemas en la calidad de los datos, como la ingesta duplicada o la falta de archivos. Por lo tanto, Databricks recomienda una configuración prudente para cloudFiles.maxFileAge, como 90 días, que es similar a lo que recomiendan soluciones comparables de ingesta de datos.
Intentar optimizar la opción cloudFiles.maxFileAge puede hacer que Auto Loader ignore los archivos sin procesar o que los archivos ya procesados expiren y, a continuación, se vuelvan a procesar, lo que provocará datos duplicados. A continuación se muestran algunos aspectos que deben tenerse en cuenta al elegir un cloudFiles.maxFileAge:
- Si la secuencia se reinicia después de mucho tiempo, se omiten los eventos de notificación de archivos que se extraen de la cola que son anteriores a
cloudFiles.maxFileAge. Del mismo modo, si usa la lista de directorios, se omiten los archivos que pueden haber aparecido durante el tiempo de inactividad y que son anteriores acloudFiles.maxFileAge. - Si usa el modo de lista de directorios y usa
cloudFiles.maxFileAge, por ejemplo, establecida en"1 month", detiene la secuencia y reinicia la secuencia concloudFiles.maxFileAgeestablecida en"2 months", los archivos más antiguos que 1 mes, pero más recientes que 2 meses, se vuelven a procesar.
Si establece esta opción la primera vez que inicia la secuencia, no ingerirá datos anteriores a cloudFiles.maxFileAge, por lo tanto, si desea ingerir datos antiguos, no debe establecer esta opción a medida que inicie la secuencia por primera vez. Sin embargo, debe establecer esta opción en ejecuciones posteriores.
Activar reposiciones periódicas usando cloudFiles.backfillInterval
En raras ocasiones, es posible que se pierdan o tarde los archivos cuando solo dependa de los sistemas de notificación, como cuando se alcancen los límites de retención de mensajes de notificación. Si tiene requisitos estrictos sobre la integridad de los datos y el Acuerdo de Nivel de Servicio, considere la posibilidad de establecer la configuración cloudFiles.backfillInterval para desencadenar rerrellenos asincrónicos en un intervalo especificado. Por ejemplo, establézcalo en un día para reposición diaria o una semana para reposición semanal. El desencadenamiento de reposición normal no provoca duplicados.
Al usar eventos de archivo, ejecute el flujo al menos una vez cada 7 días.
Al usar eventos de archivo, ejecute las secuencias del cargador automático al menos una vez cada 7 días para evitar una lista completa de directorios. Ejecutar las secuencias de Auto Loader con frecuencia garantizará que la detección de archivos sea incremental.
Para obtener procedimientos recomendados completos de eventos de archivos administrados, consulte Procedimientos recomendados para el cargador automático con eventos de archivo.