replace_flow
Beta
Esse recurso está em Beta.
O decorador @dp.replace_flow cria um fluxo REPLACE USING para uma tabela de transmissão em seu pipeline. A cada atualização, o fluxo substitui todas as linhas na tabela de destino que correspondem às colunas de chave replace_using e deixa todas as outras linhas inalteradas. A função deve retornar um DataFrame de Spark streaming do Apache. Consulte Substituição parcial de Snapshot com fluxos REPLACE USING.
Utilize @dp.replace_flow quando sua origem for uma série de Snapshots parciais indexados por coluna. Para definir a tabela de destino e o fluxo em uma única declaração, passe replace_using e sequence_by para @dp.table.
Sintaxe
from pyspark import pipelines as dp
dp.create_streaming_table("<target-table-name>") # Required only if the target table doesn't exist.
@dp.replace_flow(
target = "<target-table-name>",
replace_using = ["<key-column>", "<key-column>"],
sequence_by = "<sequence-column>",
name = "<flow-name>", # optional, defaults to function name
comment = "<comment>", # optional
spark_conf = {"<key>" : "<value>", "<key>" : "<value>"}) # optional
def <function-name>():
return (<streaming-query>)
Parâmetros
Parâmetro | Tipo | Descrição |
|---|---|---|
função |
| Obrigatório. Uma função que retorna um DataFrame de transmissão do Apache Spark a partir de uma query definida pelo usuário. |
|
| Obrigatório. O nome da tabela de transmissão que é o destino do fluxo. |
|
| Obrigatório. As colunas chave que identificam quais linhas de destino substituir. Especifique pelo menos uma coluna. As colunas chave não podem ser repetidas, e o tipo de cada coluna chave deve ser classificável. |
|
| Obrigatório. A coluna que ordena as atualizações. Para cada key, a sequência mais alta vence, e uma linha de sequência menor nunca sobrescreve uma mais alta que já esteja no destino. |
|
| O nome do fluxo. Se não for informado, o default será o nome da função. |
|
| Uma descrição para o fluxo. |
|
| Uma lista de configurações do Spark para a execução desta query. |
Exemplos
from pyspark import pipelines as dp
# Keep the latest row for each order from a stream of partial snapshots
dp.create_streaming_table("orders_current")
@dp.replace_flow(
target = "orders_current",
replace_using = ["order_id"],
sequence_by = "updated_at"
)
def orders_flow():
return spark.readStream.table("order_updates")
Use mais de uma coluna de key quando um registro for identificado por uma combinação de colunas:
from pyspark import pipelines as dp
dp.create_streaming_table("accounts_current")
@dp.replace_flow(
target = "accounts_current",
replace_using = ["region", "account_id"],
sequence_by = "updated_at"
)
def accounts_flow():
return spark.readStream.table("account_updates")
Limitações
- Uma tabela de transmissão suporta um único fluxo
REPLACE USINGe não pode combinarREPLACE USINGcom outro tipo de fluxo, como um fluxo de acréscimo, um fluxo de CDC automático ou um fluxoREPLACE WHERE. - A query deve ser uma query de transmissão.
@dp.replace_flowrejeita uma fonte que não seja de transmissão. REPLACE USINGflows requerem Databricks Runtime 18.2 e acima.