replace_flow

Important

Deze functie bevindt zich in de bètaversie.

De @dp.replace_flow decorateur maakt een VERVANG GEBRUIK-flow voor een streamingtabel in je pipeline. Bij elke update vervangt de flow alle rijen in de doeltabel die overeenkomen met de replace_using sleutelkolommen en laat alle andere rijen onaangeroerd. De functie moet een Apache Spark-streaming dataframe retourneren. Zie Gedeeltelijke snapshotvervanging met VERVANG MET stromen.

Gebruik @dp.replace_flow wanneer je bron een reeks gedeeltelijke snapshots is die per kolom zijn gecodeerd. Om de doel-tabel en de flow in één enkele instructie te definiëren, geef replace_using en sequence_by dan door naar @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>)

Parameters

Parameter Type Description
function function Required. Een functie die een Streaming DataFrame van Apache Spark retourneert vanuit een door de gebruiker gedefinieerde query.
target str Required. De naam van de streamingtabel die het doel van de flow is.
replace_using list Required. De sleutelkolommen die aangeven welke doelrijen vervangen moeten worden. Specificeer ten minste één kolom. Sleutelkolommen kunnen niet worden herhaald, en het type van elke sleutelkolom moet sorteerbaar zijn.
sequence_by str of Column Required. De kolom die de updates bestelt. Voor elke sleutel wint de hoogste reeks, en een lagere reeks rij overschrijft nooit een hogere die al in het doel staat.
name str De naam van de stroom. Als deze niet is opgegeven, wordt standaard de functienaam gebruikt.
comment str Een beschrijving voor het proces.
spark_conf dict Een lijst met Spark-configuraties voor de uitvoering van deze query.

Voorbeelden

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

Gebruik meer dan één sleutelkolom wanneer een record wordt geïdentificeerd door een combinatie van kolommen:

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

  • Een streamingtabel ondersteunt één enkele REPLACE USING flow en kan niet gecombineerd REPLACE USING worden met een ander flowtype zoals een appendflow, een automatische CDC-flow of een REPLACE WHERE flow.
  • De query moet een streamingquery zijn. @dp.replace_flow Wijst een niet-streaming bron af.
  • REPLACE USING flows vereisen Databricks Runtime 18.2 en hoger.