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.

Dołączanie przepływów może dotyczyć tabel przesyłania strumieniowego lub ujścia.

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
  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.
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.