Aller au contenu principal

replace_flow

info

Bêta

Cette fonctionnalité est en version Bêta.

Le décorateur @dp.replace_flow crée un flux REPLACE USING pour une table de streaming dans votre pipeline. À chaque mise à jour, le flux remplace toutes les lignes de la table cible qui correspondent aux colonnes clés replace_using et laisse toutes les autres lignes intactes. La fonction doit renvoyer un DataFrame de streaming Apache Spark. Voir Remplacement partiel de snapshot avec les flux REPLACE USING.

Utilisez @dp.replace_flow lorsque votre source est une série d'instantanés partiels indexés par colonne. Pour définir la table cible et le flux en une seule instruction, transmettez replace_using et sequence_by à @dp.table.

Syntaxe

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>)

parameter

parameter

Type

Description

fonction

function

Obligatoire. Une fonction qui renvoie un DataFrame Spark streaming Apache à partir d'une query définie par l'utilisateur.

target

str

Obligatoire. Le nom de la table de streaming qui est la cible du flux.

replace_using

list

Obligatoire. Les colonnes clés qui identifient les lignes cibles à remplacer. Spécifiez au moins une colonne. Les colonnes clés ne peuvent pas être répétées et le type de chaque colonne clé doit être triable.

sequence_by

str OU Column

Obligatoire. La colonne qui ordonne les mises à jour. Pour chaque clé, la séquence la plus élevée l’emporte, et une ligne avec une séquence inférieure n’écrase jamais une ligne supérieure déjà présente dans la cible.

name

str

Le nom du flux. S'il n'est pas fourni, default est le nom de la fonction.

comment

str

Une description pour le flux.

spark_conf

dict

Une liste de configurations Spark pour l'exécution de cette query.

parameter

Type

Description

fonction

function

Obligatoire. Une fonction qui renvoie un DataFrame Spark streaming Apache à partir d'une query définie par l'utilisateur.

target

str

Obligatoire. Le nom de la table de streaming qui est la cible du flux.

replace_using

list

Obligatoire. Les colonnes clés qui identifient les lignes cibles à remplacer. Spécifiez au moins une colonne. Les colonnes clés ne peuvent pas être répétées et le type de chaque colonne clé doit être triable.

sequence_by

str OU Column

Obligatoire. La colonne qui ordonne les mises à jour. Pour chaque clé, la séquence la plus élevée l’emporte, et une ligne avec une séquence inférieure n’écrase jamais une ligne supérieure déjà présente dans la cible.

name

str

Le nom du flux. S'il n'est pas fourni, default est le nom de la fonction.

comment

str

Une description pour le flux.

spark_conf

dict

Une liste de configurations Spark pour l'exécution de cette query.

Exemples

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")

Utilisez plus d'une colonne clé lorsqu'un enregistrement est identifié par une combinaison de colonnes :

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")

Limitations

  • Une table de streaming prend en charge un seul flux REPLACE USING et ne peut pas combiner REPLACE USING avec un autre type de flux tel qu'un flux d'ajout, un flux CDC automatique ou un flux REPLACE WHERE.
  • La query doit être un streaming query. @dp.replace_flow rejette une source qui n'est pas en streaming.
  • REPLACE USING les flux nécessitent Databricks Runtime 18.2 et versions ultérieures.