replace_flow

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 USING fluxo e não pode se combinar REPLACE USING com outro tipo de fluxo, como um fluxo de anexação, um fluxo automático CDC ou um REPLACE WHERE fluxo.
  • A consulta deve ser uma consulta de streaming. @dp.replace_flow rejeita uma fonte que não seja transmitida.
  • REPLACE USING os fluxos requerem Databricks Runtime 18.2 e superiores.