Nota
O acesso a esta página requer autorização. Pode tentar iniciar sessão ou alterar os diretórios.
O acesso a esta página requer autorização. Pode tentar alterar os diretórios.
O @dp.append_flow decorador cria fluxos de acréscimo ou backfills para suas tabelas de pipeline. A função deve retornar um DataFrame de streaming do Apache Spark.
Veja Dados de carga e processamento incrementalmente com fluxos de oleodutos Lakeflow.
Os fluxos adicionais podem direcionar-se a tabelas de streaming, tabelas geridas 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 |
Required. Uma função que retorna um DataFrame de streaming Apache Spark de uma consulta definida pelo usuário. |
target |
str |
Required. 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, como um backfill. O uso once=True altera o fluxo de duas maneiras:
|
depends_on |
str ou list |
Pré-visualização pública. Um ou mais nomes de fluxo que devem ser concluídos com sucesso antes de este fluxo começar. Aceita um único nome de fluxo ou uma lista de nomes. Isto ordena apenas a execução do fluxo; não altera a forma como o fluxo corre. 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 desta consulta |
import_checkpoint |
str |
O caminho para um checkpoint de Streaming Estruturado existente para importar no fluxo, de modo que um fluxo migrado retome 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. |
Examples
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 ponto de verificação 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. Define-o para a checkpointLocation sua consulta de Streaming Estruturado utilizada, que pode ser um armazenamento na cloud, um volume do Unity Catalog ou um caminho DBFS. Na primeira atualização do pipeline, o flow clona esse checkpoint no armazenamento gerido do pipeline. O fluxo retoma então a partir do último deslocamento comprometido com o seu estado (como agregações, chaves de deduplicação e marcas de água) intacto. As atualizações subsequentes do pipeline utilizam o checkpoint clonado do fluxo; o checkpoint original não é modificado.
O fluxo deve direcionar-se para uma tabela gerida criada com create_table ou um sumidouro.
Pare a consulta original de Structured Streaming antes de executar o pipeline. A consulta original de Structured Streaming pode ser reutilizada após a importação, mas é necessário gerir o estado do checkpoint e garantir que o pipeline e a consulta não escrevem na mesma tabela ao mesmo tempo, o que pode produzir dados duplicados.
Recrie a consulta de Structured Streaming como um fluxo pipeline que escreve numa nova tabela e importa o 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; começa a partir 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 anteriormente para a tabela de destino; reutilizar um nome de fluxo existente salta a importação.
Limitações
- Importar um checkpoint para uma tabela que já existe (por exemplo, o alvo original da consulta de Structured Streaming) não é suportado. Foca numa nova tabela que o pipeline cria, ou num sumidouro.
-
import_checkpointé suportado apenas em append_flow.