replace_flow

Important

Ta funkcja jest dostępna w wersji beta.

Dekorator @dp.replace_flow tworzy flow REPLACE USING dla tabeli streamingowej w twoim pipeline. Przy każdej aktualizacji flow zastępuje wszystkie wiersze w docelowej tabeli, które odpowiadają kluczowym kolumnom replace_using , a pozostałe wiersze pozostają nietknięte. Funkcja musi zwrócić strumieniową ramkę danych Apache Spark. Zobacz Częściowa wymiana migawki za pomocą przepływów ZASTÓW UŻYWAJĄC.

Używaj, @dp.replace_flow gdy źródło to seria częściowych migawek przypisanych kolumnom. Aby zdefiniować docelową tabelę i flow w jednym zaleceniu, przekaż replace_using i sequence_by do @dp.table.

Składnia

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

Parametr Typ Description
funkcja function Required. Funkcja, która zwraca strumieniowy DataFrame w Apache Spark na podstawie zapytania zdefiniowanego przez użytkownika.
target str Required. Nazwa tabeli strumieniowej, która jest celem przepływu.
replace_using list Required. Kluczowe kolumny określające, które docelowe wiersze należy zastąpić. Określ co najmniej jedną kolumnę. Kolumny klucza nie mogą być powtarzane, a typ każdej kolumny klucza musi być sortowalny.
sequence_by str lub Column Required. Kolumna, która nakazuje aktualizacje. Dla każdego klucza wygrywa najwyższa sekwencja, a wiersz niższej sekwencji nigdy nie nadpisuje wyższego już znajdującego się w celu.
name str Nazwa przepływu. Jeśli nie zostanie podana, wartość domyślna to nazwa funkcji.
comment str Opis przepływu.
spark_conf dict Lista konfiguracji platformy Spark na potrzeby wykonywania tego zapytania.

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

Używaj więcej niż jednej kolumny klucza, gdy rekord jest identyfikowany kombinacją kolumn:

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

Ograniczenia

  • Tabela strumieniowa obsługuje pojedynczy REPLACE USING przepływ i nie może się łączyć REPLACE USING z innym typem przepływu, takim jak przepływ dodatkowy, automatyczny przepływ CDC czy przepływ.REPLACE WHERE
  • Zapytanie musi być zapytaniem przesyłanym strumieniowo. @dp.replace_flow odrzuca źródło niebędące streamingiem.
  • REPLACE USING przepływy wymagają Databricks Runtime 18.2 i wyższych.