Sustitución de instantáneas parciales con flujos REPLACE USING

Importante

Esta característica se encuentra en su versión beta.

Un flujo REPLACE USING mantiene una tabla de destino sincronizada con un origen de datos en streaming: sustituye todas las filas que coinciden con las columnas clave especificadas y deja intactos todos los demás datos.

Una columna ordena las actualizaciones para que el resultado sea correcto incluso cuando las SEQUENCE BY actualizaciones llegan fuera de orden. Para cada clave, prevalece la secuencia más alta, y una fila con una secuencia inferior nunca sobrescribe a otra con una secuencia superior que ya se encuentre en el destino. Las filas que comparten la misma clave y la misma secuencia se añaden en lugar de reemplazarlas.

Cómo funciona REPLACE USING

Consideremos una tabla de eventos que contiene eventos de clics y conversión para dos regiones, secuenciados por seq:

region_id tipo_de_dispositivo event_type sec
1 iOS click 1
1 Android conversión 1
2 iOS click 1
2 escritorio click 1

Un REPLACE USING (region_id) SEQUENCE BY seq flujo recibe estas actualizaciones para las regiones 1 y 3. La Región 2 no tiene actualizaciones:

region_id tipo_de_dispositivo tipo_de_evento sec
1 iOS click 2
1 Android conversión 2
1 escritorio click 2
3 iOS click 1
3 escritorio click 2

El objetivo se convierte en:

region_id tipo_de_dispositivo event_type sec Resultado
1 iOS click 2 Reemplazado, porque seq 2 es mayor que seq 1
1 Android conversión 2 Reemplazado, porque la secuencia 2 es mayor que la secuencia 1
1 escritorio click 2 Reemplazado, porque seq 2 es mayor que seq 1
2 iOS click 1 Intacto, porque la clave no está presente en esta actualización
2 escritorio click 1 Intacto, porque la clave no está presente en esta actualización
3 escritorio click 2 Añadido. La fila seq 1 para la región 3 no se añade, porque solo se aplica la secuencia más alta para una clave.

Requisitos

Los flujos REPLACE USING tienen los siguientes requisitos:

  • Los flujos REPLACE USING se ejecutan en Databricks Runtime 18.2 y versiones posteriores, en entornos de proceso clásicos o sin servidor. Databricks recomienda Unity Catalog.
  • El origen debe ser un origen de streaming. REPLACE USING rechaza una fuente que no sea de streaming.
  • Debes especificar al menos una columna clave y exactamente una SEQUENCE BY columna.

Cuándo utilizar flujos REPLACE USING

Las tuberías de flujo lacustre ofrecen tres flujos que sobrescriben las filas existentes. Elige en función de cómo es tu fuente y cómo identifica las filas a reemplazar:

  • Utilice REPLACE USING cuando su fuente sea una serie de instantáneas parciales indexadas por columna. REPLACE USING sobrescribe únicamente los datos que tienen una coincidencia en los datos entrantes, dejando intactos todos los demás datos. No requiere una clave primaria.
  • Utiliza AUTO CDC cuando tu fuente sea una fuente de captura de datos de cambio (CDC) con operaciones explícitas de inserción, actualización y eliminación , o cuando necesites un historial de dimensión de cambio lento (SCD) Tipo 2 . AUTO CDC también requiere una clave primaria verdadera. Consulte Las API AUTO CDC: simplifican la captura de datos modificados con canalizaciones.
  • Utiliza REPLACE WHERE cuando la fuente sea una copia instantánea y quieras volver a calcular y sobrescribir un rango de la tabla de destino seleccionado mediante un predicado, por ejemplo, los últimos 7 días, en una operación por lotes. No requiere una clave primaria. Consulte Procesamiento por lotes con flujos de tipo REPLACEWHERE.

Crear un flujo REPLACE USING

Defina flujos REPLACE USING en SQL o en Python.

SQL

Use la cláusula FLOW REPLACE USING en línea con CREATE STREAMING TABLE:

CREATE STREAMING TABLE payments_current
FLOW REPLACE USING (payment_id) SEQUENCE BY payment_date BY NAME
SELECT payment_id, booking_id, status, payment_date
FROM STREAM(samples.wanderbricks.payments);

Como alternativa, use la sintaxis de formato CREATE FLOW largo:

CREATE STREAMING TABLE payments_current;

CREATE FLOW payments_flow AS
INSERT INTO payments_current BY NAME
REPLACE USING (payment_id) SEQUENCE BY payment_date
SELECT payment_id, booking_id, status, payment_date
FROM STREAM(samples.wanderbricks.payments);

Note

BY NAME es necesario en SQL. Hace coincidir las columnas por nombre en lugar de por posición.

Python

Declaremos la tabla y el flujo junto con @dp.table:

from pyspark import pipelines as dp

@dp.table(name="payments_current", replace_using=["payment_id"], sequence_by="payment_date")
def payments_current():
  return spark.readStream.table("samples.wanderbricks.payments")

Como alternativa, selecciona una tabla de streaming existente mediante @dp.replace_flow:

from pyspark import pipelines as dp

dp.create_streaming_table("payments_current")

@dp.replace_flow(target="payments_current", replace_using=["payment_id"], sequence_by="payment_date")
def payments_flow():
  return spark.readStream.table("samples.wanderbricks.payments")

replace_using es una lista de columnas clave. sequence_by es un nombre de columna o una Column expresión, y es necesario siempre que replace_using se establezca.

Secuenciación y datos fuera de orden

La SEQUENCE BY columna hace que el resultado sea independiente del orden en que llegan las actualizaciones. Una fila se aplica a una clave solo si su secuencia es mayor que la secuencia ya almacenada para esa clave, por lo que se ignora una fila tardía o reproducida que sea anterior al valor actual. Las teclas que no están presentes en una actualización quedan intactas.

Sigue estas prácticas para que el reemplazo se comporte de forma previsible:

Práctica Motivo
Utilice una secuencia que aumente estrictamente por versión de clave, como una marca de tiempo, un número de versión o un desplazamiento de registro. Se mantienen dos filas con la misma clave y la misma secuencia, lo que da lugar a filas duplicadas para esa clave.
Usa una secuencia no nula. Una secuencia nula puede dar lugar a un comportamiento indefinido.

Expectations

Los flujos REPLACE USING admiten expectativas. warn y fail se comportan como en otros flujos: warn siguen violando filas y registran la infracción, y fail detienen la actualización. Consulte Administración de la calidad de los datos con las expectativas de canalización.

Una dropexpectativa trata una fila que incumple las reglas como si la fuente nunca la hubiera generado. La fila eliminada no reemplaza, elimina ni modifica las claves correspondientes en la tabla de destino:

  • La eliminación se produce antes de la deduplicación, por lo que el flujo conserva la última versión válida de la clave.
  • Si se elimina cada fila entrante de una clave, las filas existentes de la clave quedan intactas.
  • Como una fila eliminada no establece un piso de secuencia, una actualización válida posterior sigue teniendo lugar incluso si su secuencia es menor que la de la fila eliminada.

Limitaciones

Los flujos REPLACE USING tienen las siguientes limitaciones:

  • REPLACE USING admite un único flujo por tabla de destino. No se admite combinar REPLACE USING con otro tipo de flujo en el mismo destino.
  • La tabla de destino debe crearse dentro de la canalización.
  • El origen debe ser un origen de streaming.
  • Debes especificar al menos una columna clave y una SEQUENCE BY columna. Las columnas clave no pueden repetirse, y el tipo de cada columna clave debe ser ordenable. Los tipos atómicos, como los enteros, las cadenas y las fechas, pueden utilizarse como claves, mientras que MAP y VARIANT no pueden utilizarse como tales.
  • Para las tablas de streaming independientes, consulta Aplicar la sustitución de instantáneas parciales con flujos REPLACE USING para conocer las diferencias de sintaxis.

Examples

Los siguientes ejemplos leen desde samples.wanderbricks.booking_updates, una tabla de ejemplo de cambios en el estado de las reservas que está disponible en todos los área de trabajo habilitados para Unity Catalog. Cada reserva aparece una vez por cada cambio, por lo que booking_id se repite con un nuevo booking_update_id. Consulta el conjunto de datos de Wanderbricks.

Ejemplo 1: Conserva el registro más reciente de cada tecla

Conserva solo el estado actual de cada reserva. El flujo se basa en booking_id y se ordena por booking_update_id, por lo que la actualización más reciente de una reserva sustituye a las anteriores. Usa AUTO CDC en su lugar cuando tu fuente es un feed de cambios con operaciones explícitas de insertar, actualizar y eliminar.

SQL

CREATE OR REFRESH STREAMING TABLE bookings_current
FLOW REPLACE USING (booking_id) SEQUENCE BY booking_update_id BY NAME
SELECT booking_id, status, total_amount, booking_update_id
FROM STREAM(samples.wanderbricks.booking_updates);

Python

from pyspark import pipelines as dp

@dp.table(
  name="bookings_current",
  replace_using=["booking_id"],
  sequence_by="booking_update_id"
)
def bookings_current():
  return spark.readStream.table("samples.wanderbricks.booking_updates")

Este ejemplo ordena por booking_update_id en lugar de la marca de tiempo updated_at, porque varias actualizaciones de la misma reserva pueden compartirla. Las filas que empatan en la secuencia se añaden en lugar de sustituirse, lo que dejaría más de una fila para esas reservas.

Ejemplo 2: Clave en más de una columna

Cuando un registro se identifica por una combinación de columnas, anúntalos todos en REPLACE USING. Aquí cada reserva se identifica por (property_id, booking_id), por lo que el flujo mantiene el estado actual de cada reserva por propiedad. Si una columna clave puede ser nula, REPLACE USING hace coincidir nulo con nulo en lugar de omitir la fila.

SQL

CREATE OR REFRESH STREAMING TABLE bookings_by_property
FLOW REPLACE USING (property_id, booking_id) SEQUENCE BY booking_update_id BY NAME
SELECT property_id, booking_id, status, total_amount, booking_update_id
FROM STREAM(samples.wanderbricks.booking_updates);

Python

from pyspark import pipelines as dp

@dp.table(
  name="bookings_by_property",
  replace_using=["property_id", "booking_id"],
  sequence_by="booking_update_id"
)
def bookings_by_property():
  return spark.readStream.table("samples.wanderbricks.booking_updates")

Ejemplo 3: Eliminar registros inválidos con una expectativa

Añade una condición para excluir las filas erróneas del destino. Una fila caída se trata como si la fuente nunca la hubiera producido: no reemplaza ni elimina la clave correspondiente, y el flujo vuelve a la última fila válida para esa clave. Este flujo descarta las actualizaciones que no tienen un total_amount positivo.

from pyspark import pipelines as dp

@dp.table(
  name="bookings_validated",
  replace_using=["booking_id"],
  sequence_by="booking_update_id"
)
@dp.expect_or_drop("positive_amount", "total_amount > 0")
def bookings_validated():
  return spark.readStream.table("samples.wanderbricks.booking_updates")