Ordenar a execução do fluxo do pipeline com depends_on
Visualização
Esse recurso está em Prévia Pública.
Para usar depends_on, configure seu pipeline para usar o canal PREVIEW do LakeFlow Pipelines. Consulte channel em Configurações de pipeline.
Por default, um pipeline programará fluxos por suas dependências de dados: se um fluxo lê a tabela na qual outro fluxo grava, o leitor é executado após o gravador. A ordenação de fluxos permite que um fluxo aguarde outro do qual ele não lê, declarando explicitamente a dependência com depends_on.
depends_on="other_flow" significa que esse fluxo começa apenas após a conclusão bem-sucedida de other_flow. É uma aresta de agendamento: ela controla quando um fluxo começa, não como ele é executado. A declaração de depends_on não altera o trigger, o modo ou o fato de ser um fluxo único de um pipeline. Você ainda os declara no próprio fluxo, por exemplo, com once=True.
Quando usar a ordenação de fluxo
O caso de uso principal é drenar data histórica antes de alternar para uma fonte ativa, preservando o estado de transmissão, como a migração de uma tabela de um preenchimento em lote para uma fonte ativa do Apache Kafka.
A leitura simultânea de ambas as fontes não funciona bem: o preenchimento retroativo limitado retarda a marca d'água, o que atrasa a remoção de estado e os resultados com janelas. O esvaziamento do preenchimento retroativo primeiro e, em seguida, a inicialização da transmissão ao vivo evitam isso. A ordenação de fluxos sequencia os dois.
A ordenação do fluxo também é compatível com outros padrões:
- Ordenação rigorosa de vários preenchimentos na mesma tabela.
- Pipelines em fases, como uma carga inicial, seguida de recuperação (catch-up) e, em seguida, um feed dinâmico.
- Sequenciamento de fluxos que gravam em tabelas diferentes , mas devem ser executados em uma ordem definida.
Ordenar fluxos com depends_on
depends_on está disponível nos decoradores @dp.append_flow e @dp.update_flow na API de pipelines do Python. Ele aceita um único nome de fluxo ou uma lista de nomes de fluxo. Com uma lista, cada fluxo nomeado deve ser concluído antes que o fluxo dependente comece.
O exemplo a seguir esvazia um preenchimento único na tabela de transmissão events e, em seguida, inicia uma transmissão ativa do Kafka na mesma tabela somente após a conclusão do preenchimento:
from pyspark import pipelines as dp
dp.create_streaming_table(name="events")
# Drain the historical backfill first.
@dp.append_flow(target="events", once=True, name="events_backfill")
def events_backfill():
return spark.read.table("historical_events")
# Start the live stream only after the backfill completes.
@dp.append_flow(target="events", name="events_live", depends_on="events_backfill")
def events_live():
return (
spark.readStream.format("kafka")
.option("kafka.bootstrap.servers", "<server>:<port>")
.option("subscribe", "events")
.load()
)
Para aguardar mais de um predecessor, passe uma lista. O exemplo a seguir executa dois preenchimentos em paralelo e inicia a transmissão ao vivo somente após a conclusão de ambos:
@dp.append_flow(target="events", once=True, name="backfill_2024")
def backfill_2024():
return spark.read.table("events_2024")
@dp.append_flow(target="events", once=True, name="backfill_2025")
def backfill_2025():
return spark.read.table("events_2025")
@dp.append_flow(
target="events",
name="events_live",
depends_on=["backfill_2024", "backfill_2025"],
)
def events_live():
return spark.readStream.format("kafka").option("subscribe", "events").load()
Para ordenar preenchimentos retroativos um após o outro em vez de em paralelo, encadeie depends_on entre eles:
@dp.append_flow(target="events", once=True, name="backfill_2024")
def backfill_2024():
return spark.read.table("events_2024")
@dp.append_flow(
target="events", once=True, name="backfill_2025", depends_on="backfill_2024"
)
def backfill_2025():
return spark.read.table("events_2025")
Requisitos e comportamento
As regras a seguir se aplicam à ordenação de fluxos:
- A ordenação de fluxos funciona apenas dentro de um pipeline. Sem o programador do pipeline, a ordenação não pode ser respeitada.
- Um predecessor deve gravar em uma tabela ou coletor, e não em uma view. Um fluxo que grava em uma view nunca atinge um estado terminal, de modo que um fluxo ordenado após ele nunca começaria. Destinos de tabela de transmissão e coletores
foreachBatchsão ambos predecessores válidos. - O pedido entre destinos é permitido. Um fluxo pode depender de um fluxo que grava em uma tabela diferente.
- Nomes de fluxo desconhecidos e ciclos são detectados na validação , antes da execução do pipeline.
O que pode ser um predecessor em pipelines Trigger e contínuos
Os tipos de fluxo que podem atuar como predecessores dependem do modo de execução do pipeline:
- Pipelines acionados: qualquer fluxo pode ser um predecessor. Cada fluxo em uma execução acionada atinge um estado terminal, de modo que a ordenação se aplica a cada execução.
- Pipelines contínuos: um predecessor deve ser um fluxo de execução única (
once) que atinge um estado terminal. Um fluxo que é executado continuamente nunca é encerrado, portanto, um fluxo ordenado após ele nunca começaria, e o pipeline o rejeita na validação.
Como um fluxo foreachBatch é sempre um destino de transmissão e não pode ser um fluxo de execução única, ele pode atuar como predecessor apenas em um pipeline com trigger. Em um pipeline contínuo, um fluxo foreachBatch pode aguardar um predecessor once, mas ele próprio não pode ser um predecessor.
Comportamento operacional
O estado de conclusão de um predecessor once é durável, portanto, ele sobrevive a reinicializações e atualizações de pipeline:
- Reiniciar: os fluxos já drenados permanecem drenados e são ignorados. O pipeline é retomado no primeiro fluxo que ainda não foi concluído. Depois que a execução atingir o fluxo dinâmico (live flow), as reinicializações posteriores retomam apenas o fluxo dinâmico.
- Full refresh: limpa o estado de conclusão da cadeia e a executa novamente desde o início, em ordem.
- Checkpoint Reset para um único fluxo: Reset o fluxo ativo reproduz apenas o fluxo ativo. Os preenchimentos retroativos
onceupstream permanecem esvaziados e não são executados novamente. Esta é a maneira usual de recuperar uma query ativa sem esvaziar a história novamente.