replace_flow

Important

Den här funktionen finns i Beta.

Dekoratören @dp.replace_flow skapar ett flöde ERSÄTT MED ATT ANVÄNDA MED FÖR EN STRÖMNINGSTABELL I DIN PIPELINE. Vid varje uppdatering ersätter flödet alla rader i måltabellen som matchar replace_using nyckelkolumnerna och lämnar alla andra rader orörda. Funktionen måste returnera en Apache Spark streaming-Dataram. Se Delvis snapshot-ersättning med ERSÄTT MED flöden.

Använd @dp.replace_flow när din källa är en serie partiella snapshots som är kodade efter kolumn. För att definiera måltabellen och flödet i en enda sats skickas replace_using istället och sequence_by till @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. En funktion som returnerar en Apache Spark-strömmande DataFrame från en användardefinierad fråga.
target str Required. Namnet på strömningstabellen som är målet för flödet.
replace_using list Required. De nyckelkolumner som identifierar vilka målrader som ska ersättas. Ange minst en kolumn. Nyckelkolumner kan inte upprepas, och varje nyckelkolumns typ måste vara sorterbar.
sequence_by str eller Column Required. Kolumnen som beställer uppdateringarna. För varje nyckel vinner den högsta sekvensen, och en rad i en lägre sekvens skriver aldrig över en högre som redan finns i målet.
name str Flödets namn Om det inte anges används funktionsnamnet som standard.
comment str En beskrivning av flödet.
spark_conf dict En lista över Spark-konfigurationer för körning av den här frågan.

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

Använd mer än en nyckelkolumn när en post identifieras av en kombination av kolumner:

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

  • En strömningstabell stöder ett enda REPLACE USING flöde och kan inte kombineras REPLACE USING med en annan flödestyp som ett appendflöde, ett automatiskt CDC-flöde eller ett REPLACE WHERE flöde.
  • Frågan måste vara en direktuppspelningsfråga. @dp.replace_flow Förkastar en icke-strömmande källa.
  • REPLACE USING flöden kräver Databricks Runtime 18.2 och senare.