Observação
O acesso a essa página exige autorização. Você pode tentar entrar ou alterar diretórios.
O acesso a essa página exige autorização. Você pode tentar alterar os diretórios.
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:
|
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.