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.
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 una superior ya existente 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 | seq |
|---|---|---|---|
| 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 | seq |
|---|---|---|---|
| 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 | seq | 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 de REEMPLAZAR USANDO tienen los siguientes requisitos:
- REEMPLAZAR CON flujos que se ejecutan en Databricks Runtime 18.2 o versiones posteriores, con computación clásica o sin servidor. Databricks recomienda Unity Catalog.
- El origen debe ser un origen de streaming. REPLACE USING rechaza una fuente no continua.
- Debes especificar al menos una columna clave y exactamente una
SEQUENCE BYcolumna.
Cuándo usar REEMPLAZAR USANDO flujos
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:
- Use REPLACE USING cuando la fuente sea una serie de instantáneas parciales indexadas por columna. REEMPLAZAR USANDO sobrescribe solo los datos que coinciden con los datos entrantes, dejando todos los demás datos intactos. 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 DE AUTO CDC: Simplificación de 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 REPLACEWHERE.
Crea un flujo de REEMPLAZAR USANDO
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 |
|---|---|
| Utiliza una secuencia que aumente estrictamente por versión de clave, como una marca de tiempo, número de versión o desplazamiento logarítmico. | 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
REEMPLAZAR USANDO flujos que apoyan las 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 comprobación drop trata una fila que infringe la regla como si la fuente no la hubiera producido nunca. La fila eliminada no reemplaza, elimina ni modifica las claves correspondientes en la tabla de destino:
- El descarte ocurre antes de la deduplicación, así que el flujo mantiene la última versión válida para 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
REEMPLAZAR USANDO flujos tiene 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 BYcolumna. 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 queMAPyVARIANTno pueden utilizarse como tales. - Para tablas de transmisión independientes, consulta Aplicar el reemplazo parcial de instantáneas con flujos REPLACE USING para ver 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 espacios de trabajo con Unity Catalog habilitado. Cada reserva aparece una vez por cambio, de modo que booking_id vuelve a aparecer 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 usa booking_id como clave y ordena por booking_update_id, así que la actualización más reciente de una reserva reemplaza 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 reemplazarse, 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 null, REPLACE USING hace coincidir null con null 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 la expectativa de evitar que las malas filas lleguen al objetivo. 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")