Pular para o conteúdo principal

replace_flow

info

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

Python
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

function

Obrigatório. Uma função que retorna um DataFrame de transmissão do Apache Spark a partir de uma query definida pelo usuário.

target

str

Obrigatório. O nome da tabela de transmissão que é o destino do fluxo.

replace_using

list

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.

sequence_by

str ou Column

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.

name

str

O nome do fluxo. Se não for informado, o default será o nome da função.

comment

str

Uma descrição para o fluxo.

spark_conf

dict

Uma lista de configurações do Spark para a execução desta query.

Parâmetro

Tipo

Descrição

função

function

Obrigatório. Uma função que retorna um DataFrame de transmissão do Apache Spark a partir de uma query definida pelo usuário.

target

str

Obrigatório. O nome da tabela de transmissão que é o destino do fluxo.

replace_using

list

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.

sequence_by

str ou Column

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.

name

str

O nome do fluxo. Se não for informado, o default será o nome da função.

comment

str

Uma descrição para o fluxo.

spark_conf

dict

Uma lista de configurações do Spark para a execução desta query.

Exemplos

Python
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:

Python
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 USING e não pode combinar REPLACE USING com outro tipo de fluxo, como um fluxo de acréscimo, um fluxo de CDC automático ou um fluxo REPLACE WHERE.
  • A query deve ser uma query de transmissão. @dp.replace_flow rejeita uma fonte que não seja de transmissão.
  • REPLACE USING flows requerem Databricks Runtime 18.2 e acima.