Remarque
L’accès à cette page nécessite une autorisation. Vous pouvez essayer de vous connecter ou de modifier des répertoires.
L’accès à cette page nécessite une autorisation. Vous pouvez essayer de modifier des répertoires.
Important
Cette fonctionnalité est en version bêta.
Le @dp.replace_flow décorateur crée un flux REMPLACER UTILISANT pour une table de streaming dans votre pipeline. À chaque mise à jour, le flux remplace toutes les lignes de la table cible correspondant aux replace_using colonnes clés et laisse toutes les autres lignes intactes. La fonction doit retourner un DataFrame de streaming Apache Spark. Voir Remplacement partiel de l’instantané par REMPLACER AVEC les flux.
Utilisez @dp.replace_flow lorsque votre source est une série d’instantanés partiels codés par colonne. Pour définir la table cible et le flux dans une seule instruction à la place, passez replace_using et sequence_by à @dp.table.
Syntax
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>)
Paramètres
| Paramètre | Catégorie | Description |
|---|---|---|
| function | function |
Required. Fonction qui retourne un DataFrame de streaming Apache Spark à partir d’une requête définie par l’utilisateur. |
target |
str |
Required. Le nom de la table de flux qui est la cible du flux. |
replace_using |
list |
Required. Les colonnes clés qui identifient les lignes ciblées à 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 |
Required. La colonne qui ordonne les mises à jour. Pour chaque clé, la séquence la plus haute l’emporte, et une ligne de séquence inférieure ne remplace jamais une ligne plus haute déjà dans la cible. |
name |
str |
Nom du flux. S’il n’est pas fourni, la valeur par défaut est le nom de la fonction. |
comment |
str |
Description du flux. |
spark_conf |
dict |
Liste des configurations Spark pour l’exécution de cette requête. |
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")
Utilisez plusieurs colonnes clés 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 flux prend en chargement un seul
REPLACE USINGflux et ne peut pas être combinéeREPLACE USINGavec un autre type de flux tel qu’un flux d’append, un débit CDC automatique ou unREPLACE WHEREdébit. - La requête doit être une requête de diffusion en continu.
@dp.replace_flowrejette une source non diffusée. -
REPLACE USINGles flux nécessitent Databricks Runtime 18.2 et supérieur.