Processing guarantees in Lakeflow pipelines
Retries and reruns are inevitable in any real pipeline, so this page explains the processing guarantees Lakeflow pipelines give you and how to keep the parts you write safe to re-run.
Visão geral
Duas propriedades relacionadas determinam se a reexecução de um pipeline é segura:
- Idempotency means a pipeline produces the same result no matter how many times you run it over the same input. Re-running after a failure, backfilling a date range twice, or manually retriggering a job never creates duplicate rows or corrupts state.
- Processing guarantee describes how many times each record affects the result. At-least-once processing guarantees every record is processed, but a failure and retry might process some records more than once, which risks duplicates. Exactly-once processing guarantees every record affects the result as if it were processed precisely one time, even across retries, with no duplicates and no gaps.
Lakeflow pipelines are idempotent by default for the pieces they manage, and give you exactly-once processing within their own managed tables. The important thing to understand is where those guarantees stop being automatic, so you can add the right safeguards at the edges of your pipeline.
Como funciona
Lakeflow pipelines fornecem processamento exatamente uma vez e idempotência para os fluxos que gerenciam, e oferecem ferramentas para manter a lógica que você escreve idempotente também.
Processamento exatamente uma vez para tabelas gerenciadas
Within managed tables, you get exactly-once processing by default. Streaming tables use Structured Streaming checkpoints combined with Delta Lake's transactional writes: each micro-batch commits its source offsets and its output together, so a retried batch after a failure either fully succeeds or is fully rolled back and retried, never partially applied twice. This holds for Auto Loader file ingestion, Kafka, Kinesis, and Azure Event Hubs reads, and AUTO CDC upserts, without any code from you.
If an at-least-once source sends the same record multiple times, the pipeline processes them as unique records, and writes all of them to your table. Removing those duplicates is your responsibility. See Deduplicate at-least-once sources.
A idempotência para leituras decorre desses mesmos pontos de verificação. Os pontos de verificação do Auto Loader e da tabela de transmissão garantem que cada arquivo de origem ou deslocamento seja processado uma vez para fins de acompanhamento de estado, portanto, o reprocessamento de uma atualização de pipeline após uma falha é retomado a partir do ponto de verificação, em vez de reprocessar ou ignorar dados. Você obtém isso usando tabelas de transmissão sobre spark.readStream em vez de loops de lotes manuais. Consulte Tabelas de transmissão.
Use o AUTO CDC em vez de MERGE escrito manualmente
AUTO CDC INTO é inerentemente idempotente com relação aos seus keys e sequence_by. Aplicar o mesmo registro de alteração duas vezes, ou aplicar registros fora de ordem, produz o mesmo estado final, porque o pipeline usa a coluna de sequência para decidir se uma linha recebida é realmente mais recente do que o que está armazenado:
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;
If you write your own upsert logic outside of AUTO CDC (rare, but sometimes necessary for complex merge conditions), key it on a stable business key and make it safe to apply twice, for example a MERGE ... WHEN MATCHED keyed on order_id rather than a blind INSERT. For more information, see The AUTO CDC APIs: Simplify change data capture with pipelines.
Mantenha suas próprias transformações idempotentes
Para manter a lógica idempotente ao reexecutar operações de gravação, siga estas duas diretrizes:
- Avoid non-deterministic transformations in materialized views. Because a materialized view can fully or incrementally recompute, avoid functions whose output depends on when they run rather than what the input is. For example, don't use
current_timestamp()to compute a business value that should stay fixed once written; take the timestamp from the source event or pass it in as a parameter so recomputation produces identical output. - Projete refresh completos para serem seguros. Um refresh completo descarta e recalcula uma tabela do zero, o que só é seguro se cada fonte upstream ainda puder produzir o histórico completo. Se uma fonte upstream expõe apenas uma janela móvel de alterações, um refresh completo de uma tabela
AUTO CDCdownstream pode perder a história silenciosamente; portanto, projete a retenção de fonte e tópico com isso em mente.
Obter semântica de processamento exatamente uma vez nas bordas
Onde o processamento exactly-once deixa de ser automático é nas bordas do que o pipeline controla diretamente, como gravações em sistemas externos. Ao distribuir para um sistema externo, torne a própria gravação idempotente, por exemplo, fazendo upsert por key no lado receptor, já que um micro-batch repetido poderia, de outra forma, gravar o mesmo lote duas vezes. O coletor a seguir grava cada partição do lote a partir do executor e usa uma key de idempotência para que um lote repetido não seja gravado duas vezes:
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 obter mais informações sobre como gravar em sistemas externos, consulte Sinks em LakeFlow Pipelines.
Desduplicar fontes de pelo menos uma vez
Quando uma fonte pode entregar um registro mais de uma vez, deduplique downstream. Combine uma marca d'água com dropDuplicatesWithinWatermark, que reconhece marcas d'água e não requer estado ilimitado para detectar duplicatas. Deduplique nas colunas que identificam exclusivamente um evento. A identidade pode abranger várias colunas quando nenhuma coluna única é exclusiva por si só. No exemplo a seguir, um número de sequência de clique é exclusivo apenas dentro de sua sessão, portanto, as duas colunas juntas identificam o 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"])
)
Choose those columns from the source's uniqueness contract, not from what looks distinct in sample data. Columns that can legitimately repeat discard real events when you treat them as the identity. A user clicking the same ad twice is a common example: deduplicating on the user and the ad silently drops the second click.
AUTO CDCAs semânticas de upsert baseadas em key também colapsam duplicatas naturalmente, portanto, rotear dados de "pelo menos uma vez" através de um fluxo AUTO CDC baseado em uma key de negócio estável é outra maneira de convergir para um estado de "exatamente uma vez".
Limitações
O processamento exactly-once aplica-se a fluxos Delta-to-Delta gerenciados. Trate as seguintes bordas como at-least-once e adicione lógica explícita de desduplicação ou de gravação idempotente nelas:
foreach_batch_sinke gravações externas personalizadas. O Spark garante que uma tentativa de lote seja feita pelo menos uma vez, mas um lote reexecutado após uma gravação parcial pode deixar algumas linhas visíveis duas vezes no sistema externo. Torne a gravação externa idempotente, por exemplo, fazendo upsert em uma chave natural ou gravando um ID de lote no qual o receptor possa realizar a desduplicação.- Kafka como um sink. Os tópicos do Kafka não oferecem suporte a gravações transacionais de "exatamente uma vez" da mesma forma que o Delta, portanto, um micro-lote reexecutado gravando no Kafka pode produzir mensagens duplicadas. Se os consumidores downstream forem sensíveis a duplicidades, faça a desduplicação no lado do consumidor, por exemplo, por ID de evento.
- Fontes de dados Python personalizadas usadas como fontes. Se as leituras são exatamente uma vez depende se a implementação da sua fonte relata e retoma corretamente a partir dos offsets. Se não rastrear offsets, trate como pelo menos uma vez e deduplique a jusante com
dropDuplicatesem um ID de evento ou confiando na semântica de upsert baseada em key doAUTO CDC.
Como regra geral, se todo o seu pipeline for Delta-para-Delta (tabelas de transmissão e visualizações materializadas lendo e gravando tabelas Delta por meio de fluxos gerenciados), você já tem a garantia de "exatamente uma vez". No momento em que você adicionar um foreach_batch_sink, um coletor não-Delta ou uma fonte personalizada não verificada, trate essa aresta específica como "pelo menos uma vez" e adicione lógica de gravação idempotente ou de deduplicação ali.