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.
Importante
Esse recurso está em Beta.
O @dp.replace_flow decorador cria um fluxo de SUBSTITUIR USANDO para uma mesa de streaming no seu pipeline. A 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 USING flows.
Use @dp.replace_flow quando sua fonte for uma série de snapshots parciais indexados por coluna. Para definir a tabela alvo e o fluxo em uma única instrução, replace_using passe e sequence_by passe 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>)
Parameters
| Parâmetro | Tipo | Description |
|---|---|---|
| função | function |
Required. Uma função que retorna um DataFrame de streaming do 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 linhas alvo substituir. Especifique pelo menos uma coluna. 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 maior 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 dessa 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 de uma coluna chave quando um registro é 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")
Limitations
- Uma tabela de fluxo suporta um único
REPLACE USINGfluxo e não pode se combinarREPLACE USINGcom outro tipo de fluxo, como um fluxo de anexação, um fluxo automático CDC ou umREPLACE WHEREfluxo. - A consulta deve ser uma consulta de streaming.
@dp.replace_flowrejeita uma fonte que não seja transmitida. -
REPLACE USINGos fluxos requerem Databricks Runtime 18.2 e superiores.