replace_flow

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 USING flux et ne peut pas être combinée REPLACE USING avec un autre type de flux tel qu’un flux d’append, un débit CDC automatique ou un REPLACE WHERE débit.
  • La requête doit être une requête de diffusion en continu. @dp.replace_flow rejette une source non diffusée.
  • REPLACE USING les flux nécessitent Databricks Runtime 18.2 et supérieur.