Kommentar
Åtkomst till den här sidan kräver auktorisering. Du kan prova att logga in eller ändra kataloger.
Åtkomst till den här sidan kräver auktorisering. Du kan prova att ändra kataloger.
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 USINGflöde och kan inte kombinerasREPLACE USINGmed en annan flödestyp som ett appendflöde, ett automatiskt CDC-flöde eller ettREPLACE WHEREflöde. - Frågan måste vara en direktuppspelningsfråga.
@dp.replace_flowFörkastar en icke-strömmande källa. -
REPLACE USINGflöden kräver Databricks Runtime 18.2 och senare.