Ir al contenido principal
Producto

Llevar AUTO CDC al siguiente nivel: resolviendo los casos de uso más difíciles del mundo real

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

por Josh Seidel, Shanelle Roman y Sudhanva Huruli

  • AUTO CDC reemplaza la lógica MERGE escrita a mano para la captura de datos de cambios con un pipeline declarativo
  • Spark Declarative Pipelines ahora es compatible con AUTO CDC bitemporal para realizar el seguimiento del tiempo comercial y del sistema de forma independiente, junto con actualizaciones parciales para gestionar de forma segura los campos faltantes
  • Las capacidades de AUTO CDC se expanden a la versión de código abierto de Apache Spark 4.2 para llevar la captura de datos de cambios fuera de orden estandarizada al ecosistema más amplio

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:

  • Gestionar líneas de tiempo bitemporales fuera de orden
  • Procesar actualizaciones parciales de registros sin corromper los datos existentes
  • Mantener una auditabilidad que sobreviva a las ventanas de retención de almacenamiento

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.

Seguimiento de historial de doble eje con AUTO CDC bitemporal

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:

  • Tiempo de negocio (también conocido como tiempo de evento o de validez): cuándo el hecho fue realmente cierto en el mundo real. Un símbolo de cotización pasó a ser declarable el lunes; un código de país se retiró al final del trimestre.
  • Tiempo del sistema (también conocido como tiempo de transacción o de procesamiento): cuándo el sistema de registro recibió los datos. Es posible que el cambio del lunes no llegue a la canalización hasta el miércoles.

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.

Más allá del viaje en el tiempo: ML reproducible que sobrevive a VACUUM

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.

Las actualizaciones parciales de AutoCDC ya están disponibles de forma general

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:

  1. una lista de columnas que deben ignorar los valores NULL:
    IGNORE NULL UPDATES ON columnList
  2. una lista de columnas que NO deben ignorar los valores NULL:
    IGNORE NULL UPDATES ON * EXCEPT (columnList)
  3. un nombre de columna de origen que puede ser diferente para cada fila:
    COLUMNS TO UPDATE

Para conocer la sintaxis completa, ejemplos y guías de uso, consulte la documentación de Apply Partial Updates.

Seguimos comprometidos con el código abierto

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:

  • Funciones de la próxima versión: Ya hemos fusionado la interfaz SQL (CREATE FLOW ... AS AUTO CDC INTO) en la rama master, que se incluirá en la próxima versión de Apache Spark.
  • Semántica avanzada de canalizaciones: Se está trabajando en el desarrollo de la gestión del historial completo de SCD Tipo 2, entradas de registro de cambios (changelog) nativas y soporte de actualizaciones parciales para evitar que los valores NULL sobrescriban los datos de destino.
  • Confiabilidad y pruebas: Estamos agregando capacidades de aplicar como truncar (apply-as-truncate) a la vez que ampliamos nuestros conjuntos de pruebas automatizadas para datos fuera de orden y reintentos idempotentes.

Primeros pasos

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

Recibe las últimas publicaciones en tu bandeja de entrada

Suscríbete a nuestro blog y recibe las últimas publicaciones directamente en tu bandeja de entrada.