Cómo Databricks Feature Store sirve características con una frescura inferior al segundo
por Ian Ackerman, Nick Joung y Abhay Bothra
Los modelos de machine learning son tan buenos como las señales que reciben. Un caso de uso de detección de fraude debe decidir en milisegundos, desde que un usuario presiona comprar, si permite la transacción. Tomar la decisión correcta depende de ver una transacción sospechosa ocurrida hace solo unos segundos. Combinar el promedio de transacciones de un usuario durante los últimos 30 días junto con el monto total de la transacción de los últimos 10 minutos resalta el fraude potencial. Las agregaciones de largo alcance establecen un perfil de referencia del usuario para determinar qué es lo normal, mientras que los datos más recientes ayudan a detectar cualquier comportamiento anormal justo cuando ocurre. La personalización enfrenta la misma presión: las señales más frescas son las que capturan la intención actual del usuario y fomentan la interacción.
Las canalizaciones de Spark son una forma establecida de procesar datos masivos en el Lakehouse para características de referencia históricas. Ejecutar estos trabajos por lotes de manera programada es algo bien conocido, pero introduce un retraso de minutos a horas. Para las señales de referencia sobre los usuarios, este retraso es un precio aceptable a cambio de una infraestructura más simple. Cuando los modelos requieren señales frescas, esta infraestructura falla; bajar a segundos o milisegundos no es posible en las plataformas de feature store existentes. Para ofrecer el valor de las características frescas, los científicos de datos se ven obligados a implementar una lógica compleja y específica de streaming para manejar estas agregaciones y montar una infraestructura alojada personalizada.
Databricks Feature Store le permite crear una característica una vez y usarla en todas partes: la misma definición impulsa flujos por lotes a gran escala fuera de línea y canalizaciones de características muy frescas en línea. El marco de trabajo elimina la carga de infraestructura, orquestando Spark Real-Time Mode (RTM) para el procesamiento continuo de flujos de datos, Lakebase para el almacenamiento en línea optimizado para streaming y Model Serving para la recuperación a escala. Y una vez creada, esa característica se sirve en milisegundos: latencia p99 de extremo a extremo de 200 ms, desde que un evento llega a Kafka hasta que está disponible en el feature store en línea.

Echemos un vistazo al interior para ver cómo Databricks Feature Store toma una definición de característica independiente de la infraestructura y crea una canalización para calcularla de manera constante en milisegundos. El ruta de extremo a extremo para una característica de streaming se ve así:
Vinculemos esto a nuestra característica de fraude: la suma del monto de las transacciones de un usuario durante los últimos 10 minutos. Cada evento entrante lleva los detalles de la transacción (monto, ubicación, ID de usuario, información del comerciante) y se enruta a una canalización con estado. La canalización consulta una instancia local de RocksDB que contiene el total acumulado de transacciones del usuario, con tiempos de expiración que mantienen la ventana en los últimos 10 minutos. La canalización lee e incrementa el valor localmente, luego escribe el valor de la característica actualizado en Lakebase. Así, cuando llega una consulta al modelo para aprobar una nueva transacción, se dispone de una suma de transacciones actualizada con frescura de menos de un segundo en el feature store. Esta característica de suma se recuperará junto con la referencia de compra histórica del usuario para fundamentar la aprobación. Una suma muy por encima de la referencia histórica es un fuerte indicador para el modelo de un posible fraude.
Cada componente de esta canalización se ha optimizado para que los eventos entrantes se enruten, las agregaciones se calculen y las características se escriban en la tienda en línea lo más rápido posible.

Antes de profundizar en la infraestructura, hablemos de las características de agregación y del cambio de un paradigma de sincronización por lotes a actualizaciones en tiempo real.
Las características de agregación sobre una ventana de tiempo (por ejemplo, recuentos, sumas o promedios) son señales potentes y flexibles para ML en tiempo real. Una característica por lotes a largo plazo establece una referencia histórica para el usuario durante un período de tiempo, lo que permite al modelo adaptarse y comprender el comportamiento de cada usuario. Una característica corta y fresca reacciona rápidamente a situaciones cambiantes para distinguir un nuevo interés del usuario o una actividad fraudulenta. Las ventanas de tiempo definen un rango de tiempo (por ejemplo, 10 minutos), así como la forma en que esos rangos de tiempo deben evolucionar con el tiempo (por ejemplo, superponerse o ser disjuntos).
Databricks Feature Store admite 3 ventanas de tiempo diferentes:
Las ventanas de saltos y deslizantes siguen siendo útiles cuando una característica no cambia con frecuencia: emiten menos actualizaciones, son más baratas de mantener y se adaptan de forma natural a canalizaciones programadas más sencillas. Las ventanas móviles sacrifican esa eficiencia a cambio de una frescura máxima, lo que resulta muy valioso para señales en las que cada nuevo evento debe afectar de inmediato al valor servido al modelo.
Así de sencillo es definir una característica de ventana móvil con la API declarativa de Feature Store:
Pasando a la infraestructura subyacente, la canalización de streaming es lo que hace posible obtener características frescas con un alto rendimiento. Esta canalización lleva los datos desde Kafka hasta el feature store en línea. La canalización de streaming cuenta con la tecnología de Spark Real-Time Mode (RTM), un modo de ejecución fundamentalmente nuevo para Spark Structured Streaming. RTM es la innovación arquitectónica clave que hace posible la frescura en milisegundos.
En el modo tradicional de microlotes (MBM), Spark procesa los datos de streaming en lotes discretos. Cada lote recopila eventos durante un intervalo configurable, los procesa secuencialmente a través de cada etapa, realiza un punto de control y luego comienza el siguiente lote. Esto crea un límite mínimo para la latencia: incluso con un ajuste agresivo, las canalizaciones de MBM para agregaciones con estado suelen funcionar en el orden de segundos a minutos. RTM, por otro lado, ejecuta las etapas de forma concurrente. Los operadores de agregación procesan las filas con avidez en el momento en que están disponibles, sin esperar a que la etapa anterior termine de procesar todas las filas.
Para las agregaciones móviles hay dos etapas importantes. La primera etapa es el procesamiento de datos, la validación del esquema, la fusión de datos y la conversión de tipos. Esto ejecuta la lógica de negocio que convierte los eventos de acción genéricos al formato para la agregación de sus características. La segunda etapa consiste en agregar datos por entidad para calcular las agregaciones de ventana móvil. Cada fila entrante actualiza inmediatamente la agregación en un almacén de estado local de RocksDB y emite el nuevo valor de forma descendente. La expiración de la ventana también ocurre por fila: cuando transcurre la duración de la ventana para un evento determinado, la canalización elimina la contribución de ese evento y emite la agregación corregida a Lakebase. RocksDB se ejecuta localmente en cada ejecutor, lo que permite tamaños de estado que superan la capacidad de memoria del clúster.
La creación de puntos de control es esencial para la tolerancia a fallos en el streaming con estado, ya que permite que la canalización se recupere si falla algún nodo de trabajo individual de la canalización. Pero la creación de puntos de control tiene su costo. En el modo de microlotes, Spark realiza puntos de control en cada límite de lote, y cada punto de control agrega latencia a la canalización porque interactúa con los almacenes de objetos en la nube.
RTM adopta un enfoque diferente: el coste de planificación y creación de puntos de control (checkpointing) se amortiza en intervalos más largos. El coste de creación de puntos de control se distribuye entre todas las filas procesadas en ese intervalo, en lugar de bloquear el pipeline en cada límite de lote (batch). Esto no sacrifica la tolerancia a fallos. Se mantienen las garantías de procesamiento exactamente una vez (exactly-once): en caso de fallo, el pipeline vuelve a reproducir como máximo 5 minutos de datos desde el origen de Kafka. El equilibrio es un aumento moderado en el volumen de reproducción a cambio de una reducción significativa en la latencia de procesamiento en estado estacionario.
Feature Store ejecuta pipelines de RTM sin servidor (serverless) en Lakeflow Spark Delta Pipelines (SDP), lo que elimina por completo la gestión de clústeres y la planificación de capacidad. No tiene que aprovisionar máquinas, ajustar el número de ejecutores ni preocuparse por el mantenimiento del clúster. Cuando las actualizaciones de infraestructura requieren reiniciar el pipeline, SDP coordina la transición: el nuevo clúster serverless se aprovisiona y está completamente listo antes de que se detenga el anterior. Esta coordinación se sincroniza en los intervalos de checkpointing de 5 minutos, lo que minimiza el tiempo de inactividad y evita brechas de reprocesamiento. Esto da como resultado una interrupción casi nula en la frescura de las características durante las ventanas de mantenimiento.
Databricks Feature Store utiliza Lakebase para almacenar los valores de características online para la inferencia. La arquitectura de Lakebase de separación de computación y almacenamiento permite el escalado automático (autoscaling) para gestionar la carga variable en la inferencia de modelos. El Online Feature Store aprovecha esta capacidad para escalar a decenas de miles de lecturas por segundo con decenas de milisegundos de latencia.
Las escrituras en streaming son especialmente complejas, ya que consisten en una gran cantidad de pequeñas inserciones/actualizaciones (upserts) a medida que se emiten valores frescos de ventana deslizante (rolling window) por cada fila de Kafka recibida. En Postgres estándar, este patrón puede generar un gran volumen de registro de escritura anticipada (write-ahead log), ya que Postgres utiliza escrituras de página completa para facilitar la recuperación. Después de cada checkpoint, la primera modificación de una página escribe la imagen completa de la página de 8 KB en el registro de escritura anticipada (WAL), no solo el pequeño cambio lógico. Para las filas de entidades activas (hot entities) que se actualizan con frecuencia, esto hace que la amplificación de WAL sea el cuello de botella para el rendimiento de escritura, la replicación y la sobrecarga de recuperación.
Lakebase ahora aprovecha la separación de computación y almacenamiento distribuido para minimizar la amplificación de escritura en streaming en comparación con Postgres estándar. La arquitectura de Lakebase permite que Postgres escriba registros de cambios pequeños y compactos en lugar de escribir repetidamente instantáneas (snapshots) de página completa de 8 KB en el WAL. La durabilidad sigue estando protegida porque un cuórum de nodos de seguridad (safekeeper) distribuidos confirma esos registros compactos. Las instantáneas de página completa siguen siendo necesarias para la recuperación después de suficientes registros de cambios, pero se generan más tarde en la capa de almacenamiento en lugar de sobrecargar la ruta de escritura. Para Feature Store, el resultado es que RTM puede publicar continuamente valores de características frescos en Lakebase con mucha menos amplificación de WAL y una latencia adicional mínima.
La última etapa del proceso consiste en recuperar características frescas de Lakebase y entregarlas al modelo en el momento de la inferencia. De esto se encarga Databricks Model Serving, una infraestructura de servicio totalmente gestionada y optimizada para cargas de trabajo de alta QPS y baja latencia.
Model Serving está diseñado para las demandas de rendimiento del ML en tiempo real:
Para Feature Store, la integración es perfecta. Cuando se registra un modelo con MLflow, se registran sus dependencias de características. En el momento de la inferencia, Model Serving busca automáticamente las características requeridas en Lakebase, sin código de búsqueda personalizado ni conexiones manuales. El agregado fresco calculado por RTM y almacenado en Lakebase se recupera y se une a la solicitud de inferencia de forma transparente.
Las capacidades de alto rendimiento en tiempo real son solo una parte de lo que un Feature Store puede resolver. Vale la pena considerar brevemente otros dos desafíos:
La generación de datos de entrenamiento puede ser difícil para las características de streaming, ya que las ventanas de retención cortas en los flujos requieren mantener un almacén offline independiente. Databricks Feature Store soluciona esto almacenando una copia offline de los datos de Kafka ingeridos. Para el entrenamiento de modelos, Feature Store calcula los mismos valores de características que calcularían los pipelines de streaming para los valores históricos y realiza uniones (joins) precisas en un punto en el tiempo (point-in-time). Esta misma capacidad se utiliza para rellenar (backfill) características de streaming online para permitir un lanzamiento rápido a producción.
Como se muestra arriba, los Feature Stores orquestan varios componentes de infraestructura complejos. Esa fragmentación puede dificultar la gobernanza, el linaje y la reutilización de características. También ralentiza el desarrollo, ya que los ingenieros deben coordinar los cambios a través de los límites del sistema.
En Databricks, las características son objetos de primer nivel en Unity Catalog: detectables, gobernadas con controles de acceso y rastreadas con un linaje completo. Las transformaciones de características se empaquetan con el modelo, MLflow captura qué características se utilizaron y el linaje de despliegue conecta los modelos con sus dependencias de características. La plataforma es una solución integral para desarrollar, desplegar y gobernar todo su entorno de ML.
El Feature Store de Databricks orquesta componentes clave como Spark RTM, Lakebase y Model Serving para que obtenga la mejor latencia y escala de su clase sin tener que gestionar la infraestructura usted mismo. Cada uno de estos sistemas se ha optimizado minuciosamente para cargas de trabajo de streaming con el fin de hacer realidad una frescura de 200 ms para las características de machine learning.
Consulte la documentación de Streaming Pipeline para saber cómo definir características de streaming. Experimente con las características existentes para ver qué tan fuerte sería la señal que proporcionarían con una frescura a nivel de milisegundos.
Si desea comprender mejor la tecnología subyacente, consulte el blog de Lakebase sobre escrituras más rápidas y el desglose de la arquitectura de RTM.
Si este es el tipo de problemas en los que desea trabajar, ¡estamos contratando!
(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.