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.
En esta página se describen los procedimientos recomendados que puede aplicar para configurar Auto Loader para que se ejecuten de forma confiable, rentable y a escala para su caso de uso.
Estos procedimientos recomendados reducen la sobrecarga operativa y evitan problemas comunes que son difíciles de diagnosticar en producción, como: costos de API innecesarios LIST de exámenes de directorio completo, pérdida silenciosa de datos del desfase de esquema y reinicios de canalización causados por una configuración incorrecta del punto de control.
Para obtener más información sobre la configuración de producción, consulte Configuración del cargador automático para cargas de trabajo de producción. Para la supervisión y observabilidad, consulte Supervisión y observación del cargador automático.
Elección del marco de ejecución adecuado
El mejor marco de ejecución para su caso de uso depende de la cantidad de control que necesita sobre la canalización y la sobrecarga operativa que desea administrar. Para la mayoría de los usuarios y canalizaciones de producción, Auto Loader con canalizaciones de Lakeflow es una buena opción. Sin embargo, si necesita el máximo control y la máxima personalización, use Auto Loader con Structured Streaming. Para la configuración más sencilla con una experiencia gestionada, utilice un conector de LakeFlow gestionado cuando esté disponible.
Las canalizaciones de Lakeflow amplían Structured Streaming con escalado automático, comprobaciones de calidad de datos, control de evolución de esquemas y supervisión a través del registro de eventos. Databricks recomienda los procesos de Lakeflow para la mayoría de las cargas de trabajo de ingesta en producción.
Elige el tipo de programación y de desencadenador adecuados
El mejor tipo de programación y desencadenador para su caso de uso depende de los requisitos de latencia y los patrones de llegada de archivos. Para la mayoría de los casos de uso, Databricks recomienda un desencadenador de llegada de archivos con eventos de archivo habilitados. Esto logra una ingesta de baja latencia a bajo costo porque el proceso solo se ejecuta cuando llegan nuevos archivos. Los tres tipos de desencadenantes difieren en cuándo y con qué frecuencia se inicia el pipeline:
- Continuo: el pipeline se ejecuta sin detenerse. Úselo solo cuando una latencia inferior a un segundo sea un requisito indispensable, ya que el cómputo continuo cuesta más. Combínalo con eventos de archivo.
- Desencadenador de llegada de archivos: la canalización se inicia cuando los nuevos archivos llegan a la ubicación de origen. Ideal para patrones de llegada de archivos irregulares o de baja a media latencia. Requiere que los eventos de archivo estén habilitados. Ver Activar trabajos cuando lleguen nuevos archivos.
- Programado: el pipeline se ejecuta según una programación basada en el tiempo (por ejemplo, cada hora). Úselo cuando los requisitos de latencia sean poco estrictos (de minutos a horas). Funciona con listados de directorios, pero los eventos de archivos reducen los costes incluso en modo programado, al evitar escaneos completos del directorio.
Para más información sobre el uso Trigger.AvailableNow de para la programación por lotes, consulte Uso de Trigger.AvailableNow y limitación de velocidad.
Elegir el modo de detección de archivos correcto
Auto Loader admite tres modos de detección de archivos con diferentes inconvenientes en la complejidad de la configuración, la escalabilidad y el costo.
| Mode | Complejidad de la configuración | Scalability | Cost | Cuándo se deben usar |
|---|---|---|---|---|
| Eventos de archivo (recomendado) | Bajo (configuración inicial de permisos) | Millones de archivos por hora | El más bajo | Valor predeterminado para la mayoría de las cargas de trabajo |
| Notificación clásica de archivo | Alta (21+ opciones de configuración en la nube) | Millones de archivos por hora | Media | Cuando los eventos de archivo no están disponibles |
| Lista de directorios | Ninguno | Limitado por tamaño de directorio | Más elevados (LIST costes de la API) | Directorios pequeños, rellenos puntuales o cuando las directivas de seguridad impiden eventos en archivos |
Los eventos de archivos permiten consolidar los recursos de almacenamiento en la nube mediante una suscripción y una cola por ubicación externa, en lugar de una suscripción y una cola por flujo. La diferencia de rendimiento es significativa a escala: la lista de directorios debe examinar todo el directorio de origen en cada desencadenador, por lo que el tiempo de ingesta crece con el tamaño del directorio. Los eventos de archivo entregan notificaciones de archivo nuevas directamente, por lo que el tiempo de ingesta permanece bajo independientemente del número de objetos en el directorio.
Habilitar eventos de archivo
Los eventos de archivo requieren una concesión de permisos de nube único y una ubicación externa configurada para usar el servicio de eventos de archivos administrados. Una vez configurada, todas las streams de Auto Loader que leen desde esa ubicación externa pueden usar eventos de archivos sin necesidad de configuración adicional.
Conceda los permisos de nube necesarios en el lado del proveedor de nube. Los requisitos varían según el proveedor de nube. Consulte Configuración de eventos de archivo para una ubicación externa.
Establezca
cloudFiles.useManagedFileEventsentrueen la consulta de Auto Loader.df = (spark.readStream .format("cloudFiles") .option("cloudFiles.format", "json") .option("cloudFiles.useManagedFileEvents", "true") .load("/path/to/data/dir"))Para ver los pasos de configuración completos, consulte Migración a Auto Loader con eventos de archivo.
Cuando no se pueden usar eventos de archivo
Es posible que no pueda usar eventos de archivo cuando:
- La ubicación externa no está configurada para usar eventos de archivo.
- Las directivas de seguridad de la organización no permiten habilitar eventos de archivo en una ubicación externa compartida.
En estos casos, use el modo de notificación de archivos clásico o el modo de lista de directorios. Para obtener una comparación completa de los modos de detección de archivos, consulte Comparación de los modos de detección de archivos del cargador automático.
Administrar la evolución del esquema
El cargador automático deduce automáticamente el esquema, pero la forma en que configura la evolución del esquema afecta a la integridad de los datos y a la estabilidad de la canalización. Use la tabla siguiente para elegir una estrategia.
| Escenario | Recommendation |
|---|---|
| El esquema se conoce y se ha corregido | Proporcionar un esquema explícito con .schema() |
| El esquema es desconocido y se esperan cambios aditivos |
schemaEvolutionMode: addNewColumns |
| El esquema es desconocido y se esperan cambios de tipo |
schemaEvolutionMode: addNewColumnsWithTypeWidening |
| Se requiere un contrato de esquema estricto |
schemaEvolutionMode: failOnNewColumns |
| Esquema arbitrario o imprevisible | Inadministración como tipo Variant |
Después de haber elegido una estrategia, aplique los procedimientos siguientes para ajustar el comportamiento de la evolución del esquema.
Usar sugerencias de esquema para tipos de campo conocidos
Use la cloudFiles.schemaHints opción para aplicar tipos para los campos que conoce de antemano, a la vez que permite la inferencia de esquemas para otros campos.
df = (spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.schemaHints", "id long, amount double")
.load("/path/to/data/dir"))
Utiliza la ampliación de tipos para los cambios de tipo compatibles
El addNewColumnsWithTypeWidening modo de evolución del esquema amplía automáticamente los tipos compatibles (por ejemplo, int a long) en lugar de enrutar los datos a la _rescued_data columna. Esto evita la necesidad de trabajos de posprocesamiento para administrar promociones de tipos simples.
df = (spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "parquet")
.option("cloudFiles.schemaEvolutionMode", "addNewColumnsWithTypeWidening")
.load("/path/to/data/dir"))
Inadministración como tipo Variant para esquemas impredecibles
Cuando los datos no se ajustan a ningún esquema específico o el esquema cambia continuamente, ingiere los datos como un Variant tipo.
df = (spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("singleVariantColumn", "data")
.load("/path/to/data/dir"))
Variant proporciona esquema en lectura en el momento de la consulta, pero es menos eficiente que consultar columnas estructuradas. Para obtener la mecánica completa de la inferencia de esquemas y la evolución, consulte Configuración de la inferencia de esquemas y la evolución en Auto Loader.
Gestionar datos erróneos y la calidad de los datos
Los procedimientos siguientes le ayudan a detectar, capturar y aislar datos incorrectos antes de que se propague a capas descendentes.
Habilitar _rescued_data y _corrupt_record
Auto Loader proporciona dos columnas para capturar datos que no se pueden analizar de forma limpia.
-
_rescued_datacaptura campos que no coinciden con el esquema actual. Se añade automáticamente mediante Auto Loader. -
_corrupt_recordcaptura filas que no se pueden analizar en absoluto. Habilite mediantecolumnNameOfCorruptRecord:
df = (spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.schemaHints", "_corrupt_record string")
.option("columnNameOfCorruptRecord", "_corrupt_record")
.load("/path/to/data/dir"))
Databricks recomienda columnNameOfCorruptRecord en lugar de badRecordsPath para evitar posibles condiciones de carrera que puedan pasar por alto registros dañados.
Utiliza los valores esperados de los pipelines de Lakeflow para la supervisión
Configura los valores esperados de los pipelines de Lakeflow para verificar que _rescued_data y _corrupt_record sean NULL en condiciones normales. Los valores que no son NULL indican el desfase del esquema o los datos dañados.
import dlt
@dlt.table
@dlt.expect("no rescued data", "_rescued_data IS NULL")
@dlt.expect("no corrupt records", "_corrupt_record IS NULL")
def bronze_table():
return (spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.schemaHints", "_corrupt_record string")
.option("columnNameOfCorruptRecord", "_corrupt_record")
.load("/path/to/data/dir"))
Aislar los datos dañados
Aísle las filas que contengan datos imposibles de analizar en un destino específico para su investigación. Esto evita que los datos dañados se propague a las capas de bajada.
import dlt
@dlt.table
def corrupt_records_sink():
return dlt.read_stream("bronze_table").where("_corrupt_record IS NOT NULL")
@dlt.view
def clean_table():
return dlt.read_stream("bronze_table").where("_corrupt_record IS NULL")
Anotación de datos con metadatos de archivo de origen
Incluya la columna _metadata en las consultas de ingestión de Auto Loader. Como mínimo, capture file_path y file_modification_time. Esto te permite rastrear los problemas de datos hasta archivos fuente específicos y hacer una unión con cloud_files_state() para obtener el ciclo de vida completo de los archivos.
df = (spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.load("/path/to/data/dir")
.select("*", "_metadata.file_path", "_metadata.file_modification_time"))
Para obtener más información, consulte Columna de metadatos de archivo.
Optimización del costo y el rendimiento
Los procedimientos siguientes reducen los tres principales factores de costo para Auto Loader: llamadas API en la nube LIST , proceso inactivo y crecimiento del almacenamiento a largo plazo.
Uso de eventos de archivo para minimizar LIST Costos de API: los eventos de archivo proporcionan detección incremental de archivos, lo que elimina la necesidad de listas de directorios completas en cada ejecución. Esta es la optimización de costos más impactante para Auto Loader.
Utilice desencadenadores por llegada de archivos para el procesamiento controlado por eventos: los desencadenadores por llegada de archivos inician la canalización solo cuando llegan archivos nuevos, por lo que no paga por capacidad de cómputo inactiva. Ver Activar trabajos cuando lleguen nuevos archivos.
Archivar archivos procesados con cloudFiles.cleanSource: use
cloudFiles.cleanSourcepara eliminar o mover archivos procesados automáticamente. Esto reduce tanto los costes de almacenamiento como los costes de listado de directorios en flujos de larga duración. Para más información, consulte Archivado de archivos en el directorio de origen para reducir los costos.- Utilice el modo
deletepara eliminar archivos después de la ingesta. - Use el modo
movepara archivar los archivos en otra ubicación para fines de cumplimiento normativo o auditoría.
df = (spark.readStream .format("cloudFiles") .option("cloudFiles.format", "json") .option("cloudFiles.cleanSource", "delete") .load("/path/to/data/dir"))Advertencia
No habilite
cloudFiles.cleanSourcesi varios flujos del cargador automático u otros clientes leen desde el mismo directorio de origen.- Utilice el modo
Aproveche las mejoras de rendimiento: actualice a la versión más reciente de Databricks Runtime o use un proceso sin servidor para beneficiarse de las mejoras recientes de rendimiento del cargador automático.
Administración de puntos de control
El punto de control almacena el progreso del flujo y el estado del archivo. La configuración incorrecta o la pérdida del punto de control requiere un reinicio completo, por lo que se trata como infraestructura crítica.
- Nunca aplique directivas de ciclo de vida de objetos en la nube a ubicaciones de punto de control. Si se eliminan los archivos de punto de control, el estado de la secuencia está dañado y debe reiniciarse desde cero.
- Use puntos de control independientes para cada flujo y directorio de origen.
- Considere
cloudFiles.maxFileAgepara flujos de larga duración y gran volumen a fin de limitar el crecimiento del estado. Use una configuración conservadora (se recomienda un mínimo de 90 días). Establecer este valor demasiado agresivamente corre el riesgo de volver a procesar los archivos que el cargador automático ya ha ingerido si se encuentran fuera de la ventana.
Para obtener más información, consulte Seguimiento de eventos de archivo.
Use volúmenes para una detección óptima de archivos mediante eventos de archivo
Para mejorar el rendimiento con eventos de archivo, cree un volumen externo para cada ruta de acceso o subdirectorio desde el que se carga el cargador automático. Proporcione rutas de volumen (por ejemplo, /Volumes/catalog/schema/volume) a Auto Loader en lugar de rutas de nube (por ejemplo, s3://bucket/path). Esto optimiza la detección de archivos mediante un patrón de acceso a datos optimizado.
Para obtener más procedimientos recomendados de eventos de archivo, consulte Procedimientos recomendados para el cargador automático con eventos de archivo.