append_flow

Dekorator @dp.append_flow tworzy dołączanie przepływów lub uzupełniania dla tabel potoków. Funkcja musi zwrócić strumieniową ramkę danych Apache Spark. Zobacz Ładowanie i przetwarzanie danych przyrostowo za pomocą przepływów potoku lakeflow.

Przepływy dodawania mogą celować w tabele strumieniowe, tabele zarządzane lub sinki.

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.append_flow(
  target = "<target-table-name>",
  name = "<flow-name>", # optional, defaults to function name
  once = False, # optional
  depends_on = "<flow-name>", # optional, Public Preview
  spark_conf = {"<key>" : "<value", "<key" : "<value>"}, # optional
  comment = "<comment>", # optional
  import_checkpoint = "<checkpoint-path>") # optional
def <function-name>():
  return (<streaming-query>) #

Parametry

Parameter Typ Description
funkcja function To jest wymagane. Funkcja, która zwraca strumieniowy DataFrame w Apache Spark na podstawie zapytania zdefiniowanego przez użytkownika.
target str To jest wymagane. Nazwa tabeli lub ujścia będącego elementem docelowym przepływu dołączania.
name str Nazwa przepływu. Jeśli nie zostanie podana, wartość domyślna to nazwa funkcji.
once bool Opcjonalnie zdefiniuj przepływ jako przepływ jednorazowy, taki jak wypełnienie wsteczne. Używanie once=True zmienia przepływ na dwa sposoby:
  • Wartość zwracana. streaming-query. w tym przypadku musi być wsadową ramką danych, a nie strumieniową ramką danych.
  • Domyślnie przepływ jest uruchamiany jeden raz. W przypadku zaktualizowania pipeline'u przez pełne odświeżenie, przepływ ONCE zostanie uruchomiony ponownie w celu odtworzenia danych.
depends_on str lub list Publiczna wersja zapoznawcza. Jedna lub więcej nazw przepływów, które muszą zakończyć się pomyślnie, zanim ten przepływ się rozpocznie. Akceptuje pojedynczą nazwę przepływu lub listę nazw. To nakazuje tylko wykonywanie przepływu; nie zmienia sposobu działania przepływu. Zobacz Wykonywanie przepływu potoku zamówień z depends_on.
comment str Opis przepływu.
spark_conf dict Lista konfiguracji platformy Spark na potrzeby wykonywania tego zapytania
import_checkpoint str Ścieżka do istniejącego punktu kontrolnego Structured Streaming jest importowana do przepływu, tak aby migrowany strumień wznawiał się od ostatniego zadeklarowanego offsetu zamiast ponownie przetwarzać źródło. Import punktu kontrolnego jest w fazie beta. Zobacz checkpoint Migracja Structured Streaming.

Przykłady

from pyspark import pipelines as dp

# Create a sink for an external Delta table
dp.create_sink("my_sink", "delta", {"path": "/tmp/delta_sink"})

# Add an append flow to an external Delta table
@dp.append_flow(name = "flow", target = "my_sink")
def flowFunc():
  return <streaming-query>

# Add a backfill
@dp.append_flow(name = "backfill", target = "my_sink", once = True)
def backfillFlowFunc():
    return (
      spark.read
      .format("json")
      .load("/path/to/backfill/")
    )

# Create a Kafka sink
dp.create_sink(
  "my_kafka_sink",
  "kafka",
  {
    "kafka.bootstrap.servers": "host:port",
    "topic": "my_topic"
  }
)

# Add an append flow to a Kafka sink
@dp.append_flow(name = "flow", target = "my_kafka_sink")
def myFlow():
  return read_stream("xxx").select(F.to_json(F.struct("*")).alias("value"))

Migracja punktu kontrolnego Structured Streaming

Ważna

Import punktu kontrolnego jest w fazie beta.

Używam do import_checkpoint migracji istniejącego obciążenia Structured Streaming do potoku bez ponownego przetwarzania źródła. Ustaw to na checkpointLocation zapytanie Structured Streaming, które może być przechowywaniem w chmurze, woluminem katalogu Unity lub ścieżką DBFS. Podczas pierwszej aktualizacji potoku przepływ klonuje ten punkt kontrolny do zarządzanej pamięci magazynowej. Następnie przepływ wznawia się od ostatniego zadeklarowanego przesunięcia z nienaruszonym stanem (takim jak agregacje, klucze deduplikacyjne i znaki wodne). Kolejne aktualizacje potoku wykorzystują sklonowany punkt kontrolny przepływu; oryginalny punkt kontrolny nie jest modyfikowany.

Przepływ musi być skierowany do zarządzanej tabeli utworzonej za pomocą create_table lub zlewa.

Zatrzymaj oryginalne zapytanie Structured Streaming przed uruchomieniem pipeline. Oryginalne zapytanie Structured Streaming można ponownie użyć po imporcie, ale musisz zarządzać jego stanem checkpointów i upewnić się, że pipeline i zapytanie nie zapisują się jednocześnie do tej samej tabeli, co może generować duplikaty danych.

Odtworz zapytanie Structured Streaming jako przepływ potoku, który zapisuje się do nowej tabeli i importuje swój punkt kontrolny:

from pyspark import pipelines as dp

# Create a new managed table for the pipeline
dp.create_table("target_table")

# Continue from the imported checkpoint instead of reprocessing the source.
@dp.append_flow(
  target = "target_table",
  import_checkpoint = "/Volumes/my_catalog/my_schema/checkpoints/my_stream",
)
def migrate():
  # The same source your original query read from.
  return spark.readStream.table("source_table")

Punkt kontrolny jest importowany tylko raz, przy pierwszej aktualizacji potoku; późniejsze aktualizacje ignorują import_checkpoint. Pełne odświeżenie nie importuje ponownie punktu kontrolnego; zaczyna się od nowego, pustego punktu kontrolnego i przetwarza źródło ponownie. Aby zaimportować inny punkt kontrolny, użyj nazwy przepływu, która wcześniej nie była używana dla tabeli docelowej; ponowne użycie istniejącej nazwy przepływu pomija import.

Ograniczenia

  • Nie jest obsługiwane importowanie punktu kontrolnego do tabeli, która już istnieje (na przykład oryginalnego celu zapytania Structured Streaming). Celuj w nową tabelę, którą tworzy potok, lub w zlew.
  • import_checkpoint jest obsługiwany tylko na append_flow.