append_flow

O decorador @dp.append_flow cria fluxos de acréscimo ou de backfill para suas tabelas de pipeline. A função deve retornar um DataFrame de streaming do Apache Spark. Veja Carregar e processar dados incrementalmente com fluxos de pipeline do Lakeflow.

Fluxos de anexação podem direcionar tabelas de streaming, tabelas gerenciadas ou sumidouros.

Sintaxe

from pyspark import pipelines as dp

dp.create_streaming_table("<target-table-name>") # Required only if the target table doesn't exist.

@dp.append_flow(
  target = "<target-table-name>",
  name = "<flow-name>", # optional, defaults to function name
  once = False, # optional
  depends_on = "<flow-name>", # optional, Public Preview
  spark_conf = {"<key>" : "<value", "<key" : "<value>"}, # optional
  comment = "<comment>", # optional
  import_checkpoint = "<checkpoint-path>") # optional
def <function-name>():
  return (<streaming-query>) #

Parâmetros

Parâmetro Tipo Description
função function Obrigatório Uma função que retorna um DataFrame de streaming do Apache Spark de uma consulta definida pelo usuário.
target str Obrigatório O nome da tabela ou coletor que é o destino do fluxo de acréscimo.
name str O nome do fluxo. Se não for fornecido, o padrão será o nome da função.
once bool Opcionalmente, defina o fluxo como um fluxo único, por exemplo, um preenchimento retroativo. Usar once=True altera o fluxo de duas maneiras:
  • O valor retornado. streaming-query. deve ser um DataFrame em lote nesse caso, não um DataFrame de streaming.
  • O fluxo é executado uma vez por padrão. Se a linha de processamento for atualizada com um recarregamento completo, então o fluxo ONCE será executado novamente para recriar os dados.
depends_on str ou list Prévia Pública. Um ou mais nomes de fluxo que devem ser concluídos com sucesso antes que esse fluxo comece. Aceita um único nome de fluxo ou uma lista de nomes. Essa ordem é apenas a execução do fluxo; não altera a forma como o fluxo funciona. Veja execução de fluxo de pipeline de ordens com depends_on.
comment str Uma descrição para o fluxo.
spark_conf dict Uma lista de configurações do Spark para a execução dessa consulta
import_checkpoint str O caminho até um checkpoint de Streaming Estruturado existente para importar no fluxo, então um fluxo migrado retoma a partir do seu último offset comprometido em vez de reprocessar a fonte. Importar um checkpoint está em Beta. Veja Migrar um ponto de verificação de Streaming Estruturado.

Exemplos

from pyspark import pipelines as dp

# Create a sink for an external Delta table
dp.create_sink("my_sink", "delta", {"path": "/tmp/delta_sink"})

# Add an append flow to an external Delta table
@dp.append_flow(name = "flow", target = "my_sink")
def flowFunc():
  return <streaming-query>

# Add a backfill
@dp.append_flow(name = "backfill", target = "my_sink", once = True)
def backfillFlowFunc():
    return (
      spark.read
      .format("json")
      .load("/path/to/backfill/")
    )

# Create a Kafka sink
dp.create_sink(
  "my_kafka_sink",
  "kafka",
  {
    "kafka.bootstrap.servers": "host:port",
    "topic": "my_topic"
  }
)

# Add an append flow to a Kafka sink
@dp.append_flow(name = "flow", target = "my_kafka_sink")
def myFlow():
  return read_stream("xxx").select(F.to_json(F.struct("*")).alias("value"))

Migre um checkpoint de Streaming Estruturado

Importante

Importar um checkpoint está em Beta.

Use import_checkpoint para migrar uma carga de trabalho de Streaming Estruturado existente para um pipeline sem reprocessar a fonte. Defina para a checkpointLocation consulta de Streaming Estruturado usada, que pode ser um armazenamento em nuvem, um volume do Unity Catalog ou um caminho DBFS. Na primeira atualização do pipeline, o fluxo clona esse checkpoint no armazenamento gerenciado do pipeline. O fluxo então retoma a partir do último deslocamento comprometido com seu estado (como agregações, chaves de deduplicação e marcas d'água) intacto. Atualizações subsequentes do pipeline usam o checkpoint clonado do fluxo; o checkpoint original não é modificado.

O fluxo deve ter como alvo uma tabela gerenciada criada com create_table ou um sumidouro.

Pare a consulta original de Structured Streaming antes de executar o pipeline. A consulta original de Streaming Estruturado pode ser reutilizada após a importação, mas você precisa gerenciar o estado do checkpoint e garantir que o pipeline e a consulta não escrevam na mesma tabela ao mesmo tempo, o que pode gerar dados duplicados.

Recrie a consulta de Streaming Estruturado como um fluxo pipeline que escreve em uma nova tabela e importa seu checkpoint:

from pyspark import pipelines as dp

# Create a new managed table for the pipeline
dp.create_table("target_table")

# Continue from the imported checkpoint instead of reprocessing the source.
@dp.append_flow(
  target = "target_table",
  import_checkpoint = "/Volumes/my_catalog/my_schema/checkpoints/my_stream",
)
def migrate():
  # The same source your original query read from.
  return spark.readStream.table("source_table")

O checkpoint é importado apenas uma vez, na primeira atualização do pipeline; atualizações posteriores ignoram import_checkpoint. Uma atualização completa não reimporta o checkpoint; ele começa de um novo checkpoint vazio e reprocessa a fonte. Para importar um checkpoint diferente, use um nome de fluxo que não tenha sido usado para a tabela de destino antes; reutilizar um nome de fluxo existente pula a importação.

Limitações

  • Importar um checkpoint para uma tabela que já existe (por exemplo, o alvo original da consulta de Streaming Estruturado) não é suportado. Foque em uma nova tabela que o pipeline cria, ou em um sumidouro.
  • import_checkpoint é suportado apenas em append_flow.