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.
Importante
Este recurso está em versão Beta.
O @dp.replace_flow decorador cria um fluxo de SUBSTITUIR USANDO para uma mesa de streaming no seu pipeline. Em cada atualização, o fluxo substitui todas as linhas da tabela de destino que correspondem às replace_using colunas-chave e deixa todas as outras linhas intocadas. A função deve retornar um DataFrame de streaming do Apache Spark. Veja Substituição parcial de snapshot com SUBSTITUIR fluxos UTILIZANTES.
Use @dp.replace_flow quando a sua fonte for uma série de snapshots parciais indexados por coluna. Para definir a tabela alvo e o fluxo numa única instruçã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 | 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 de streaming que é o alvo do fluxo. |
replace_using |
list |
Required. As colunas-chave que identificam quais as filas-alvo a substituir. Especifique pelo menos uma coluna. As colunas chave não podem ser repetidas, e o tipo de cada coluna chave deve ser ordenável. |
sequence_by |
str ou Column |
Required. A coluna que ordena as atualizações. Para cada chave, a sequência mais alta vence, e uma linha de sequência inferior nunca sobrescreve uma linha superior já no alvo. |
name |
str |
O nome do fluxo. Se não for fornecido, o padrão 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 consulta. |
Examples
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 do que uma coluna chave quando um registo é 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 fluxo suporta um único
REPLACE USINGfluxo e não pode combinar-seREPLACE USINGcom outro tipo de fluxo, como um fluxo de anexação, um fluxo automático CDC ou umREPLACE WHEREfluxo. - A consulta tem de ser uma consulta de streaming.
@dp.replace_flowrejeita uma fonte que não seja streaming. -
REPLACE USINGos fluxos requerem Databricks Runtime 18.2 e superiores.