replace_flow
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
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 |
| Obligatoire. Une fonction qui renvoie un DataFrame Spark streaming Apache à partir d'une query définie par l'utilisateur. |
|
| Obligatoire. Le nom de la table de streaming qui est la cible du flux. |
|
| 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. |
|
| 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. |
|
| Le nom du flux. S'il n'est pas fourni, default est le nom de la fonction. |
|
| Une description pour le flux. |
|
| Une liste de configurations Spark pour l'exécution de cette query. |
Exemples
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 :
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 USINGet ne peut pas combinerREPLACE USINGavec un autre type de flux tel qu'un flux d'ajout, un flux CDC automatique ou un fluxREPLACE WHERE. - La query doit être un streaming query.
@dp.replace_flowrejette une source qui n'est pas en streaming. REPLACE USINGles flux nécessitent Databricks Runtime 18.2 et versions ultérieures.