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.
Los reintentos y las nuevas ejecuciones son inevitables en cualquier canalización real, por lo que esta página explica las garantías de procesamiento que ofrecen las canalizaciones de Lakeflow y cómo hacer que las partes que escribes puedan volver a ejecutarse de forma segura.
Overview
Dos propiedades relacionadas determinan si es seguro volver a ejecutar una pipeline:
- Idempotencia significa que una tubería produce el mismo resultado sin importar cuántas veces la ejecutes sobre la misma entrada. Volver a ejecutar tras un fallo, rellenar un rango de fechas dos veces o reactivar manualmente un trabajo nunca crea filas duplicadas ni corrompe el estado.
- La garantía de procesamiento describe cuántas veces cada registro afecta al resultado. Al menos una vez el procesamiento garantiza que todos los registros se procesen, pero un fallo y un nuevo intento pueden procesar algunos registros más de una vez, lo que puede provocar duplicados. El procesamiento exacto una vez garantiza que cada registro afecta al resultado como si se hubiera procesado exactamente una vez, incluso entre intentos, sin duplicados ni huecos.
Las canalizaciones de Lakeflow son idempotentes de forma predeterminada en los componentes que gestionan y ofrecen procesamiento exactamente una vez en sus propias tablas administradas. Lo importante es entender dónde dejan de ser automáticas esas garantías, para que puedas añadir las salvaguardas adecuadas en los bordes de tu pipeline.
Cómo funciona
Las canalizaciones de Lakeflow proporcionan procesamiento exactamente una vez e idempotencia para los flujos que gestionan, y también le ofrecen herramientas para que la lógica que escriba siga siendo idempotente.
Procesamiento exacto de una vez para tablas administradas
En las tablas administradas, se obtiene procesamiento exactamente una vez de forma predeterminada. Las tablas de streaming utilizan puntos de control de Structured Streaming combinados con las escrituras transaccionales de Delta Lake: cada microlote confirma conjuntamente sus offsets de origen y sus datos de salida, de modo que un lote reintentado tras un fallo o bien se completa correctamente por entero, o bien se revierte por completo y se vuelve a intentar, sin que llegue a aplicarse parcialmente dos veces. Esto se aplica a la ingesta de archivos de Auto Loader, las lecturas de Kafka, Kinesis y Azure Event Hubs y las operaciones de inserción o actualización AUTO CDC, sin que tenga que escribir nada de código.
Si al menos una fuente envía el mismo registro varias veces, la tubería los procesa como registros únicos y los escribe todos en tu tabla. Eliminar esos duplicados es tu responsabilidad. Consulte Eliminar duplicados en orígenes de al menos una vez.
La idempotencia para las lecturas se deduce de esos mismos puntos de control. Los puntos de control de Auto Loader y de las tablas de streaming garantizan que cada archivo de origen o desplazamiento se procese una sola vez a efectos del seguimiento del estado, de modo que, tras un fallo, el reprocesamiento de una actualización de la canalización se reanuda a partir del punto de control en lugar de reprocesar u omitir datos. Esto se consigue usando tablas en streaming con spark.readStream en lugar de bucles de procesamiento por lotes implementados manualmente. Consulte Tablas de streaming.
Usa AUTO CDC en lugar de MERGE escrito manualmente
AUTO CDC INTO es inherentemente idempotente respecto a su keys y sequence_by. Aplicar el mismo registro de cambio dos veces, o aplicar registros fuera de orden, produce el mismo estado final, porque la tubería utiliza la columna de secuencia para decidir si una fila entrante es realmente más nueva que lo almacenado:
CREATE FLOW customers_cdc_flow AS AUTO CDC INTO customers_silver
FROM stream(customers_cdc_bronze)
KEYS (customer_id)
SEQUENCE BY sequence_num
STORED AS SCD TYPE 1;
Si escribe su propia lógica de inserción o actualización fuera de AUTO CDC (algo raro, pero a veces necesario para condiciones de fusión complejas), basándose en una clave de negocio estable y asegurándose de que pueda aplicarse dos veces sin problemas; por ejemplo, MERGE ... WHEN MATCHED sobre order_id en lugar de INSERT a ciegas. Para más información, consulte Las APIs de AUTO CDC: Simplificar la captura de datos de cambios con pipelines.
Mantén tus propias transformaciones idempotentes
Para mantener la lógica idempotente al volver a ejecutar operaciones de escritura, sigue estas dos directrices:
- Evitar transformaciones no deterministas en vistas materializadas. Como una vista materializada puede recomputar total o incrementalmente, evita funciones cuya salida depende de cuándo se ejecutan en lugar de cuál es la entrada. Por ejemplo, no uses
current_timestamp()para calcular un valor de negocio que debería permanecer fijo una vez escrito; toma la marca de tiempo del evento fuente o pásala como parámetro para que el recálculo produzca una salida idéntica. - Diseñe actualizaciones completas para estar seguro. Una actualización completa elimina y vuelve a calcular una tabla desde cero, lo cual solo es seguro si todas las fuentes de origen pueden seguir proporcionando el historial completo. Si una fuente aguas arriba solo expone una ventana móvil de cambios, una actualización completa de una tabla
AUTO CDCaguas abajo puede provocar la pérdida silenciosa del historial, así que diseña la retención de la fuente y del tema teniendo esto en cuenta.
Obtener semántica de exactamente una vez en los extremos
Donde exactamente una vez deja de ser automático es en los límites de lo que la canalización controla directamente, como las escrituras en sistemas externos. Cuando se expande a un sistema externo, haga que la escritura sea idempotente, por ejemplo inserción o actualización por tecla en el lado receptor, ya que un microlote reprobado podría escribir el mismo lote dos veces. El siguiente receptor escribe cada partición del lote desde los ejecutores y utiliza una clave de idempotencia para que un lote reintentado no escriba doblemente:
from pyspark import pipelines as dp
@dp.foreach_batch_sink(name="orders_to_external_api")
def write_orders_to_api(batch_df, batch_id):
def write_partition(rows):
# Open one client per partition.
for row in rows:
# Use an idempotency key (order_id) so a retried batch doesn't double-write.
upsert_to_external_system(key=row.order_id, payload=row.asDict())
batch_df.select("order_id", "amount").foreachPartition(write_partition)
Para obtener más información sobre cómo escribir en sistemas externos, consulte Sinks in Lakeflow pipelines.
Eliminar duplicados en orígenes de al menos una vez
Cuando un origen puede entregar un registro más de una vez, elimina los duplicados aguas abajo. Combina una marca de agua temporal con dropDuplicatesWithinWatermark, que es compatible con marcas de agua temporales y no requiere un estado no acotado para detectar duplicados. Elimine duplicados según las columnas que identifican de forma unívoca cada evento. La identidad puede abarcar varias columnas cuando ninguna columna es única por sí sola. En el siguiente ejemplo, un número de secuencia de clics es único solo dentro de su sesión, por lo que las dos columnas juntas identifican el evento:
from pyspark import pipelines as dp
@dp.table(name="clicks_deduped")
def clicks_deduped():
return (
spark.readStream.table("clicks_bronze")
.withWatermark("click_ts", "5 minutes")
.dropDuplicatesWithinWatermark(["session_id", "click_seq_num"])
)
Elija esas columnas del contrato de unicidad del origen, no de lo que se vea distinto en los datos de muestra. Las columnas que pueden repetirse de forma legítima descartan eventos reales cuando se las trata como identificador. Que un usuario haga clic dos veces en el mismo anuncio es un ejemplo habitual: al deduplicar por usuario y anuncio, se descarta silenciosamente el segundo clic.
La semántica de inserción o actualización basada en claves de AUTO CDC también elimina duplicados de forma natural, por lo que encaminar datos con entrega al menos una vez a través de un flujo AUTO CDC basado en una clave de negocio estable es otra forma de converger hacia un estado exacto de una vez.
Limitations
El procesamiento exacto de una vez se aplica a flujos gestionados de Delta a Delta. Trate los siguientes extremos como de al menos una vez y añada lógica explícita de deduplicación o escritura idempotente en esos puntos:
-
foreach_batch_sinky operaciones de escritura externas personalizadas. Spark garantiza que se intente un lote al menos una vez, pero si se intenta de nuevo tras una escritura parcial, algunas filas pueden quedarse visibles dos veces en el sistema externo. Haga que la escritura externa sea idempotente, por ejemplo insertando o actualizando según una clave natural o escribiendo un identificador de lote con el que el receptor pueda eliminar duplicados. - Kafka como un receptor. Los temas de Kafka no soportan escrituras transaccionales de exactamente una vez como lo hace Delta, así que una escritura microlote reintentada en Kafka puede producir mensajes duplicados. Si los consumidores aguas abajo son sensibles a los duplicados, deduplique en el lado del consumidor, por ejemplo por ID de evento.
- Fuentes de datos personalizadas en Python utilizadas como fuentes. Que las lecturas sean de exactamente una vez depende de si la implementación del origen notifica correctamente los offsets y reanuda la lectura a partir de ellos. Si no rastrea los offsets, trátelo como de al menos una vez y deduplique aguas abajo con
dropDuplicatessobre un ID de evento o basándose en la semántica de inserción o actualización basada en claves deAUTO CDC.
Como regla general, si toda su canalización es Delta-a-Delta (tablas en streaming y vistas materializadas leyendo y escribiendo tablas Delta a través de flujos gestionados), ya tiene exactamente una vez. En el momento en que añade un foreach_batch_sink, un receptor que no es Delta o un origen personalizado no verificado, trate ese extremo específico como de al menos una vez y añada lógica de escritura de idempotentes o deduplicación en esos puntos.