Desde el cumplimiento bitemporal hasta las actualizaciones parciales de registros: ofreciendo una captura de datos de cambios sólida y lista para auditorías sin código personalizado
La captura de datos modificados (CDC) es una de las tareas más comunes que los ingenieros de datos crean en Spark, y una de las más tediosas de implementar correctamente a mano. En nuestra publicación anterior, Deje de programar a mano las canalizaciones de captura de datos modificados, presentamos cómo AUTO CDC en Apache™ Spark Declarative Pipelines (SDP) automatiza SCD Tipo 1, SCD Tipo 2 y Snapshot CDC al reemplazar cientos de líneas de lógica MERGE frágil con unas pocas declaraciones simples.
A medida que los requisitos de las canalizaciones evolucionan, los ingenieros se enfrentan a situaciones que los patrones CDC estándar tienen dificultades para resolver:
Hoy, llevamos AUTO CDC al siguiente nivel para resolver exactamente estos desafíos del mundo real, y expandimos estas capacidades a Apache Spark 4.2 de código abierto.
Las tablas SCD Tipo 2 estándar pueden indicarle cuándo cambió un hecho en el mundo real, pero no pueden decirle lo que su sistema creía en un momento dado.
Bajo la Regla 17a-4 de la SEC y las normas de conservación de registros de FINRA, las empresas deben poder reconstruir los registros tal como existían en un momento dado; solo la campaña de control de conservación de registros de la SEC ha generado más de 2000 millones de dólares en multas a más de 100 empresas desde 2021. La parte difícil rara vez es almacenar el valor actual. Lo difícil es responder, meses después, qué decían los datos de referencia en la fecha de informe y qué creían nuestros sistemas en ese momento.
SCD Tipo 2 estándar realiza el seguimiento de una línea de tiempo: cuándo cambió un hecho. AUTO CDC bitemporal realiza el seguimiento de dos, de forma independiente:
Cada tabla de destino obtiene cuatro columnas administradas por el sistema: __START_AT y __END_AT para el tiempo de negocio, __SYSTEM_START_AT y __SYSTEM_END_AT para el tiempo del sistema. Un solo hecho lógico puede tener varias filas físicas, una por cada combinación de versión de negocio y versión de sistema, que es lo que hace posible la reconstrucción en un momento dado a lo largo de cualquiera de los ejes. La garantía de comportamiento clave: los eventos pueden llegar en cualquier orden en cualquiera de las líneas de tiempo.
Cuando aparece una corrección con un tiempo de negocio o de sistema anterior a algo que ya se ha procesado, el motor reescribe el historial afectado en lugar de simplemente añadirlo al final. Sin lógica escrita a mano, solo declare las dos columnas de secuenciación y el motor mantendrá ambos intervalos. Esto funciona igual de bien para tablas de dimensiones, como maestros de símbolos, y para tablas de hechos, como el historial de transacciones o lecturas de sensores, que necesitan una auditabilidad estricta. Así es como se ve frente a los datos de referencia de FINRA CAT:
Tenga en cuenta que la cláusula SQL exacta es STORED AS BITEMPORAL, no STORED AS SCD TYPE BITEMPORAL, y requiere tanto SEQUENCE BY como SYSTEM SEQUENCE BY. Supongamos que el indicador de declarable de Acme cambia el 1 de enero (tiempo de negocio), pero la fuente no lo recibe hasta el 5 de enero (tiempo del sistema). Luego, el 8 de enero llega una corrección con fecha anterior que dice que el cambio real fue el 1 de enero pero con un valor diferente. AUTO CDC bitemporal puede responder a ambas preguntas:
El 3 de enero, la primera consulta no devuelve nada, la respuesta correcta y auditable de lo que mostraba el sistema en ese momento. La segunda consulta, ejecutada hoy, refleja la verdad corregida. Dos relojes, dos respuestas, ambas correctas. Las columnas de secuenciación deben ser de tipos ordenables, sin valores de secuenciaci ón NULL. La función se ejecuta en SDP serverless o en las ediciones de producto Pro/Advanced, y actualmente está en versión Beta, así que fije la canalización al canal: PREVIEW.
Cuando un modelo se entrena con datos de referencia o de características, la reproducibilidad significa poder reconstruir el conjunto de datos exacto que utilizó el modelo, meses después, durante una revisión o una auditoría. El instinto es recurrir al viaje en el tiempo de Delta Lake, pero esa es una propiedad del historial de archivos de la tabla, no un registro permanente. VACUUM elimina permanentemente los archivos de datos a los que ya no hacen referencia las versiones recientes; una vez pasada la ventana de retención predeterminada de 7 días, un TIMESTAMP AS OF registrado en el momento del entrenamiento puede dejar de resolverse silenciosamente. Una tabla bitemporal almacena ese historial como datos, no como versiones de archivos. VACUUM y OPTIMIZE compactan archivos pero nunca tocan el historial lógico, por lo que cada versión anterior de negocio o sistema sigue siendo una fila consultable. Hay dos formas de obtener reproducibilidad a partir de esto: registrar dos instantes de referencia (tiempo de negocio y de sistema) como parámetros de MLflow y fijar la consulta de entrenamiento a ese estado de creencia:
O, si la tabla expone una vista actual, registre un único instante del sistema en el momento del entrenamiento y reconstruya más tarde con una consulta de tiempo del sistema en esa marca de tiempo:
De cualquier manera, el contrato de reproducibilidad consiste en un par de marcas de tiempo en la ejecución de MLflow y, debido a que el historial bitemporal se almacena como filas, ese contrato se mantiene incluso después de que VACUUM haya limpiado los archivos subyacentes.
No todas las fuentes de captura de datos modificados (CDC) emiten filas completas para las actualizaciones. En su lugar, muchas solo envían los campos que cambiaron, representando todas las demás columnas como NULL. Sin un manejo especial, estos valores NULL pueden sobrescribir involuntariamente los datos existentes en la tabla de destino. Hasta ahora, los clientes tenían que crear una lógica personalizada para solucionar este comportamiento. Con las actualizaciones parciales de AutoCDC, esto ahora se maneja automáticamente.
Las actualizaciones parciales extienden AutoCDC al permitir que los eventos de actualización modifiquen solo un subconjunto de columnas. Para las columnas seleccionadas, los valores NULL en una actualización entrante se interpretan como "no actualizar" en lugar de sobrescribir el valor existente.
Esto es especialmente útil para las fuentes CDC que omiten los valores sin cambios emitiendo NULL. Sin las actualizaciones parciales, estos valores NULL sobrescribirían los datos existentes en la tabla de destino.
Por ejemplo, supongamos que la tabla de destino contiene: (1, 'A', 20)
Un evento de actualización entrante contiene: (1, NULL, 30)
Por defecto, AutoCDC actualizaría la fila a: (1, NULL, 30).
Con las actualizaciones parciales habilitadas, el NULL en name se trata como "dejar el valor existente sin cambios", lo que da como resultado: (1, 'A', 30).
Habilitar las actualizaciones parciales solo requiere agregar un parámetro a su definición de AutoCDC. Puede elegir entre tres formas de especificar qué columnas deben tratarse como actualizaciones parciales:
IGNORE NULL UPDATES ON columnListIGNORE NULL UPDATES ON * EXCEPT (columnList)COLUMNS TO UPDATEPara conocer la sintaxis completa, ejemplos y guías de uso, consulte la documentación de Apply Partial Updates.
Spark Declarative Pipelines es de código abierto, por lo que su tipo de flujo más utilizado también debería serlo. Comenzamos contribuyendo con la API de Python para AUTO CDC Tipo 1 a Apache Spark 4.2.
Lo hemos aportado de la misma manera en que evoluciona el resto de Spark: como una serie de propuestas revisadas y solicitudes de extracción (pull requests), no como una entrega de código única (consulte el SPIP y SPARK-56249).
La corrección con datos fuera de orden viene integrada: una pequeña tabla auxiliar realiza el seguimiento del estado de los eventos que llegan antes de tiempo, como las marcas de borrado (delete tombstones), los microbatches reintentados convergen en lugar de corromper el destino y, dado que se basa en las abstracciones de tablas y streaming de Spark en lugar de en un formato de almacenamiento, se ejecuta tanto en Delta Lake como en Apache Iceberg.
Lo que viene a continuación, en código abierto:
CREATE FLOW ... AS AUTO CDC INTO) en la rama master, que se incluirá en la próxima versión de Apache Spark.NULL sobrescriban los datos de destino.Ya sea que desee implementar el cumplimiento bitemporal, configurar actualizaciones parciales o explorar AutoCDC de código abierto en Apache Spark, consulte los siguientes recursos para comenzar:
(Esta entrada del blog ha sido traducida utilizando herramientas basadas en inteligencia artificial) Publicación original
Suscríbete a nuestro blog y recibe las últimas publicaciones directamente en tu bandeja de entrada.