Zbuduj niestandardową aplikację stanową z transformWithState

transformWithState Umożliwia tworzenie stanowych aplikacji przesyłania strumieniowego oraz implementowanie rozwiązań o małych opóźnieniach i niemal w czasie rzeczywistym. Za pomocą niestandardowych operatorów stanowych można utworzyć dowolną logikę stanową, która umożliwia tworzenie nowych przypadków użycia operacyjnych, które nie są możliwe w przypadku tradycyjnego przetwarzania przesyłania strumieniowego ze strukturą.

Notatka

W przypadku operacji stanowych, takich jak agregacje, deduplikacja i sprzężenia strumieniowe, usługa Databricks zaleca używanie wbudowanych operatorów przesyłania strumieniowego ze strukturą zamiast logiki niestandardowej. Co to jest przesyłanie strumieniowe stanowe? Zobacz .

Usługa Databricks zaleca używanie transformWithState zamiast starszych operatorów, takich jak flatMapGroupsWithState i mapGroupsWithState, w przypadku dowolnych przekształceń stanu. Zobacz Starsze dowolne stanowe operatory.

Wymagania

Operatory transformWithState i transformWithStateInPandas mają następujące wymagania:

  • Dostępne w środowisku Databricks Runtime 16.2 lub nowszym.
    • W przypadku trybu czasu rzeczywistego użyj środowiska Databricks Runtime 17.3 LTS lub nowszego. Zobacz koncepcje trybu czasu rzeczywistego.
    • W przypadku standardowego trybu dostępu Python jest dostępna w środowisku Databricks Runtime 16.3 lub nowszym, a język Scala jest dostępny w środowisku Databricks Runtime w wersji 17.3 lub nowszej.
  • RocksDB jest domyślnym dostawcą magazynu stanów w środowisku Databricks Runtime 17.3 lub nowszym.
    • W przypadku środowiska Databricks Runtime 17.2 lub starszego należy skonfigurować dostawcę magazynu stanów bazy danych RocksDB. Usługa Databricks zaleca włączenie bazy danych RocksDB w konfiguracji platformy Spark.

      spark.conf.set("spark.sql.streaming.stateStore.providerClass", "org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider")
      

Co to jest transformWithState?

Operator transformWithState stosuje niestandardowy procesor stanowy do zapytania Structured Streaming. Aby używać transformWithState, należy zaimplementować niestandardowy procesor stanowy. Przesyłanie strumieniowe ze strukturą obejmuje interfejsy API do tworzenia procesora stanowego przy użyciu Python, języka Scala lub Java.

Użyj elementu transformWithState, aby zastosować logikę niestandardową do klucza grupowania. Poniżej opisano projekt wysokiego poziomu:

  • Zdefiniuj co najmniej jedną zmienną stanu.
  • Informacja o stanie jest zachowywana dla każdego klucza grupowania. Dostęp do każdej zmiennej stanu można uzyskać w kodzie zdefiniowanym przez użytkownika.
  • Dla każdej przetworzonej mikropartii wszystkie wiersze dla danego klucza są dostępne w postaci iteratora.
  • Użyj funkcji StatefulProcessorHandle z czasomierzami i warunkami zdefiniowanymi przez użytkownika, aby kontrolować sposób emitowania wierszy.
  • Aby zarządzać wygasaniem stanu i jego rozmiarem, dla wartości stanu można definiować indywidualne czasy życia (TTL).

Ponieważ transformWithState obsługuje ewolucję schematu w magazynie stanów, możesz wprowadzać kolejne zmiany i aktualizować swoje aplikacje produkcyjne bez utraty historycznych informacji o stanie. Po zaktualizowaniu schematu stanu nie jest wymagane ponowne przetwarzanie wierszy, co upraszcza wdrażanie kodu i konserwację. Patrz Ewolucja schematu w przechowywaniu stanów.

Ważny

W dokumentacji usługi Azure Databricks używa się znacznika transformWithState do opisywania implementacji zarówno w języku Python, jak i Scala:

  • Narzędzie PySpark obsługuje zarówno API oparte na wierszach transformWithState, jak i operatory bazujące na bibliotece Pandas transformWithStateInPandas.
    • transformWithStateInPandas nie jest obsługiwany w trybie czasu rzeczywistego. Zamiast tego użyj transformWithState. Aby uzyskać szczegółowe informacje, zobacz transformWithState w trybie czasu rzeczywistego.
    • Interfejs API oparty na wierszach transformWithState obsługuje przetwarzanie asynchroniczne z użyciem asyncio, aby zwiększyć przepustowość. Przetwarzanie asynchroniczne nie jest obsługiwane w środowisku bezserwerowym. Zobacz Przetwarzanie asynchroniczne (Beta).
  • Scala obsługuje tylko interfejs API oparty na wierszach transformWithState.

Implementacje transformWithState języka Scala i Python mają te same możliwości, ale z pewnymi różnicami w składni.

Definiowanie StatefulProcessor

Zdefiniuj procesor stanowy, rozszerzając klasę StatefulProcessor i implementując jej metody.

Spark przekazuje obiekt StatefulProcessorHandle do metody init w Twoim StatefulProcessor. Użyj uchwytu, aby utworzyć zmienne stanu i wchodzić w interakcję z magazynem stanów.

transformWithState obsługuje trzy typy stanów: ValueState, ListStatei MapState. Każdy typ przechowuje stan dla każdego klucza grupowania przy użyciu innej podstawowej struktury danych.

Zaimplementuj następujące metody, aby zdefiniować logikę niestandardową:

  • Zaimplementuj handleInputRows, aby sterować tym, jak aplikacja przetwarza dane, aktualizuje stan i generuje wiersze dla każdej mikropartii. Zobacz Obsługa wierszy wejściowych.
  • Zaimplementuj handleExpiredTimer, aby uruchamiać logikę opartą na czasie niezależnie od tego, czy klucz grupowania otrzymuje nowe wiersze w mikropartii. Zobacz Obsługa wygasłych czasomierzy.
  • Opcjonalnie zaimplementuj handleInitialState, aby wstępnie wypełnić stan, zanim aplikacja przetworzy jakiekolwiek wiersze wejściowe. Zobacz Obsługa stanu początkowego.

W poniższej tabeli porównano zachowania funkcjonalne tych metod:

Zachowanie handleInputRows handleExpiredTimer
Pobieranie, umieszczanie, aktualizowanie lub czyszczenie wartości stanu Tak Tak
Tworzenie lub usuwanie czasomierza Tak Tak
Generuj wiersze Tak Tak
Iteruj po wierszach w bieżącej mikropartii Tak Nie
Logika wyzwalacza oparta na upływie czasu Nie Tak

W razie potrzeby można połączyć zarówno handleInputRows, jak i handleExpiredTimer, aby zaimplementować złożoną logikę.

Można na przykład zaimplementować aplikację, która używa handleInputRows do aktualizowania wartości stanu dla każdej mikropartii i ustawia czasomierz na 10 sekund do przodu. Jeśli nie są przetwarzane żadne dodatkowe wiersze, można użyć handleExpiredTimer, aby wyemitować bieżące wartości zapisane w magazynie stanu. Jeśli nowe wiersze są przetwarzane dla klucza grupowania, możesz wyczyścić istniejący czasomierz i ustawić nowy czasomierz.

StatefulProcessorHandle

W programie StatefulProcessorHandle PySpark klasa umożliwia dostęp do funkcji, które kontrolują sposób korzystania z informacji o stanie w kodzie.

Podczas inicjowania elementu StatefulProcessornależy zawsze importować i przekazywać element StatefulProcessorHandle do zmiennej handle . Zmienna handle łączy zmienną lokalną w klasie Python ze zmienną stanu.

Notatka

Język Scala używa metody getHandle.

Niestandardowe typy stanów

Można zaimplementować wiele obiektów stanu w jednym operatorze stanowym.

Wybierz typ stanu na podstawie pełnej logiki aplikacji. Można na przykład śledzić sesje za pomocą ValueState, pogrupowane według user_id i session_id. Lub, aby ocenić warunki w wielu sesjach, użyj MapState pogrupowanego według user_idsession_id jako klucza mapy.

Jeśli obiekt stanu używa elementu StructType, należy zdefiniować unikatowe nazwy dla każdego pola w strukturze w schemacie. Te nazwy są widoczne podczas odczytu repozytorium stanu. Zobacz Odczytaj informacje o stanie zorganizowanego przesyłu strumieniowego.

W poniższych sekcjach opisano typy stanów obsługiwane przez program transformWithState:

ValueState

ValueState przechowuje wartość dla każdego klucza grupowania.

Stan wartości może obejmować złożone typy, takie jak struktura lub krotka. W przypadku ValueState należy zaimplementować logikę zastępującą całą wartość.

Czas życia (TTL) stanu wartości jest resetowany po zaktualizowaniu wartości. Jeśli przetwarzasz klucz źródłowy dla ValueState bez aktualizowania zapisanego ValueState, czas życia nie zostanie zresetowany.

ListState

ListState przechowuje listę dla każdego klucza grupowania.

Stan listy to kolekcja wartości, z których każda może zawierać typy złożone. Każda wartość na liście ma swój własny czas wygaśnięcia.

Elementy można dodawać do listy, dodając poszczególne elementy, dołączając listę elementów lub zastępując całą listę przy pomocy put. Aby zresetować czas wygaśnięcia, należy wykonać operację put.

MapState

MapState przechowuje mapę dla każdego klucza grupowania. Mapy są odpowiednikiem słownika Pythona w Apache Spark (dict).

Stan mapy to zbiór unikalnych kluczy, z których każdy odpowiada jednej wartości, a każda z tych wartości może obejmować typy złożone. Każda para klucz-wartość na mapie ma swój własny czas wygaśnięcia.

Możesz zaktualizować wartość określonego klucza lub usunąć klucz i jego wartość. Możesz zwrócić pojedynczą wartość przy użyciu klucza, wyświetlić listę wszystkich kluczy, wyświetlić listę wszystkich wartości lub zwrócić iterator, aby pracować z pełnym zestawem par klucz-wartość na mapie.

Ważny

Klucze grupowania opisują pola określone w klauzuli GROUP BY zapytania Strukturalnego Przesyłania Strumieniowego. Stany mapy mogą zawierać dowolną liczbę par klucz-wartość dla klucza grupowania.

Jeśli na przykład zapytanie używa GROUP BY user_id i chcesz zdefiniować mapę dla każdego session_id, kluczem grupowania jest user_id, a klucz MapState to session_id:

Python
class SessionTracker(StatefulProcessor):
  def init(self, handle: StatefulProcessorHandle) -> None:
    self.sessions = handle.getMapState("sessions", "session_id string", "count long")

  def handleInputRows(self, key, rows: Iterator[Row], timerValues) -> Iterator[Row]:
    for row in rows:
      session_key = (row["session_id"],)  # session_id is the MapState key
      count = self.sessions.getValue(session_key)[0] if self.sessions.containsKey(session_key) else 0
      new_count = count + 1
      self.sessions.updateValue(session_key, (new_count,))
    yield from []

  def close(self) -> None:
    pass

df.groupBy("user_id").transformWithState(SessionTracker(), ...) # user_id is the grouping key
Skala
case class Event(userId: String, sessionId: String)

class SessionTracker extends StatefulProcessor[String, Event, (String, Long)] {
  @transient private var sessions: MapState[String, Long] = _

  override def init(outputMode: OutputMode, timeMode: TimeMode): Unit = {
    sessions = getHandle.getMapState[String, Long]("sessions", Encoders.STRING, Encoders.scalaLong, TTLConfig.NONE)
  }

  override def handleInputRows(
      key: String,
      rows: Iterator[Event],
      timerValues: TimerValues): Iterator[(String, Long)] = {
    rows.foreach { event =>
      val count = if (sessions.containsKey(event.sessionId)) sessions.getValue(event.sessionId) else 0L
      sessions.updateValue(event.sessionId, count + 1) // sessionId is the MapState key
    }
    Iterator.empty
  }
}

df.as[Event]
  .groupByKey(_.userId) // userId is the grouping key
  .transformWithState(new SessionTracker(), TimeMode.None(), OutputMode.Update())

Tworzenie niestandardowej zmiennej stanu w obiekcie StatefulProcessor

Podczas inicjowania StatefulProcessortworzysz zmienną lokalną dla każdego obiektu stanu, co pozwala na interakcję z obiektami stanu w twojej logice niestandardowej. Zdefiniuj i zainicjuj zmienne stanu przez przesłonięcie wbudowanej metody init w klasie StatefulProcessor.

Można zdefiniować dowolną liczbę obiektów stanu przy użyciu getValueStatemetod , getListStatei getMapState w obiekcie StatefulProcessor.

Każdy obiekt stanu musi mieć następujące elementy:

  • Unikatowa nazwa
  • Schemat
    • W Python należy określić schemat.
    • W języku Scala można przekazać element Encoder , aby określić schemat stanu.

Opcjonalnie możesz również podać czas życia (TTL) w milisekundach. W przypadku implementowania stanu mapy należy podać oddzielną definicję schematu dla kluczy mapy i wartości.

Notatka

Logika StatefulProcessor obsługuje oddzielnie zapytania, aktualizowanie i emitowanie informacji o stanie. Zobacz Używanie zmiennych stanu w metodach z logiką niestandardową.

Używaj swoich zmiennych stanu w metodach z własną logiką

Obiekty stanu mają metody pobierania stanu, aktualizowania istniejących informacji o stanie i czyszczenia bieżącego stanu.

Każdy klucz grupowania ma własne informacje o jego stanie.

  • StatefulProcessor emituje wiersze na podstawie własnej logiki i określonego schematu wyjściowego. Zobacz Emitowanie wierszy.
  • Użyj odczytnika statestore, aby uzyskać dostęp do wartości w magazynie stanu. Ten czytnik jest przeznaczony do obsługi obciążeń wsadowych i nie jest przeznaczony do obsługi obciążeń wymagających niskich opóźnień. Zobacz Odczytaj informacje o stanie zorganizowanego przesyłu strumieniowego.
  • Logika określona przy użyciu handleInputRows jest wykonywana tylko wtedy, gdy w mikropartii są obecne wiersze dla danego klucza. Zobacz Obsługa wierszy wejściowych.
  • Użyj handleExpiredTimer do implementacji logiki opartej na czasie, która nie zależy od obserwowania wierszy, aby działać. Zobacz Obsługa wygasłych czasomierzy.

Notatka

Obiekty stanu są izolowane przez grupowanie kluczy z następującymi konsekwencjami:

  • Na wartości stanu nie mogą wpływać wiersze powiązane z innym kluczem grupowania.
  • Nie można zaimplementować logiki, która zależy od porównywania wartości lub aktualizowania stanu między kluczami grupowania.

Wartości w kluczu grupowania można porównać. Użyj MapState, aby zaimplementować logikę przy użyciu drugiego klucza, którego może używać logika niestandardowa. Na przykład grupowanie według user_id oraz użycie ip_address jako klucza MapState umożliwia śledzenie jednoczesnych sesji użytkowników.

Zagadnienia zaawansowane dotyczące pracy ze stanem

Aktualizacje stanu są odporne na błędy. Jeśli zadanie ulegnie awarii, zanim zakończy przetwarzanie mikropartii, ponowna próba użyje wartości z ostatniej pomyślnie przetworzonej mikropartii.

W celu zoptymalizowanej wydajności usługa Databricks zaleca przetwarzanie wszystkich wartości w iteratorze dla danego klucza i zatwierdzanie aktualizacji w jednym zapisie. Gdy zapisujesz do zmiennej stanu, powoduje to zapis do RocksDB.

Wartości stanu nie mają wartości domyślnych. Jeśli logika wymaga odczytywania istniejących informacji o stanie, użyj exists metody .

Aby zaimplementować logikę stanu null, MapState zmienne umożliwiają sprawdzanie poszczególnych kluczy lub wyświetlanie listy wszystkich kluczy.

Obsługa wierszy wejściowych

handleInputRows Użyj metody , aby zdefiniować sposób przetwarzania wierszy i aktualizacji wartości stanu przez aplikację. Ta metoda jest uruchamiana za każdym razem, gdy zapytanie przesyłania strumieniowego ze strukturą przetwarza wiersze klucza grupowania.

W przypadku większości aplikacji stanowych zaimplementowanych za pomocą transformWithStatepodstawowa logika jest definiowana przy użyciu handleInputRows.

Dla każdej przetworzonej aktualizacji mikropartii wszystkie wiersze w mikropartii dla danego klucza grupowania są dostępne za pomocą iteratora. Logika zdefiniowana przez użytkownika może wchodzić w interakcję ze wszystkimi wierszami z bieżącej mikropartii oraz z wartościami przechowywanymi w magazynie stanów.

Obsługa wygasłych czasomierzy

Użyj metody handleExpiredTimer, aby zaimplementować własną logikę na podstawie upływu czasu.

W kluczu grupowania czasomierze są jednoznacznie identyfikowane przez znacznik czasu.

Po wygaśnięciu czasomierza wynik jest określany przez logikę zaimplementowaną w aplikacji. Typowe wzorce obejmują:

  • Emitowanie informacji przechowywanych w zmiennej stanu.
  • Usuwanie przechowywanych informacji o stanie.
  • Tworzenie nowego czasomierza.

Wygasłe timery są wyzwalane nawet wtedy, gdy w mikropartii nie są przetwarzane żadne wiersze dla ich powiązanego klucza.

Określanie trybu czasu

Podczas przekazywania StatefulProcessor do elementu transformWithState należy określić tryb czasu za pomocą parametru timeMode.

Obsługiwane są następujące opcje:

Tryb czasowy Opis
ProcessingTime Zarówno timery, jak i TTL są obsługiwane i wyznaczane na podstawie czasu rzeczywistego, gdy Apache Spark przetwarza każdą mikropartię. Użyj ProcessingTime polecenia , gdy czasomierze mają być uruchamiane w stałym interwale względem czasu przetwarzania wierszy, niezależnie od sygnatur czasowych w danych.
EventTime Czasomierze są obsługiwane i oceniane na podstawie limitu czasu zdarzenia. Watermark przesuwa się wraz z tym, jak Apache Spark obserwuje znaczniki czasu w danych wejściowych. TTL nie jest obsługiwane z użyciem EventTime. Użyj funkcji EventTime, gdy dane zawierają znaczniki czasu i chcesz, by timery były wyzwalane na podstawie postępu tych znaczników czasu. W przypadku użycia EventTime należy również określić parametr eventTimeColumnName. Zobacz: eventTimeColumnName.
NoTime lub TimeMode.None() Liczniki czasowe i TTL nie są obsługiwane. Użyj polecenia NoTime, gdy aplikacja z zachowaniem stanu nie wymaga logiki zależnej od czasu.

eventTimeColumnName

W przypadku użycia trybu czasu EventTime parametr eventTimeColumnName określa nazwę kolumny w schemacie wyjściowym, która zawiera znacznik czasu zdarzenia. Apache Spark używa tej kolumny do przekazywania znacznika watermark do strumienia wyjściowego, umożliwiając poprawne operacje zależne od czasu w dalszych etapach przetwarzania.

Python

eventTimeColumnName jest dodatkowym argumentem elementu transformWithState lub transformWithStateInPandas:

q = (
  df.groupBy("key")
    .transformWithState(
      statefulProcessor=MyProcessor(),
      outputStructType=output_schema,
      outputMode="Append",
      timeMode="EventTime",
      eventTimeColumnName="outputTimestamp",
    )
    .writeStream...
)
Skala

transformWithState akceptuje eventTimeColumnName zamiast timeMode. To podejście zawsze używa EventTime trybu:

val q = spark
  .readStream
  .format("delta")
  .load(srcDeltaTableDir)
  .as[(String, String)]
  .groupByKey(x => x._1)
  .transformWithState(
    new MyProcessor(),
    "outputTimestamp",
    OutputMode.Append(),
  )
  .writeStream...

Wbudowane wartości czasomierza

Databricks zdecydowanie odradza wywoływanie zegara systemowego w Twojej niestandardowej aplikacji stanowej, ponieważ może to prowadzić do zawodnych ponownych prób w przypadku niepowodzenia zadania. Używaj metod w klasie TimerValues, gdy musisz uzyskać dostęp do czasu przetwarzania lub znaku wodnego.

TimerValues Opis
getCurrentProcessingTimeInMs Zwraca znacznik czasu przetwarzania dla bieżącej partii w milisekundach od początku epoki.
getCurrentWatermarkInMs Zwraca znacznik czasu dla bieżącej partii w milisekundach od epoki.

Notatka

Czas przetwarzania odnosi się do okresu, w którym mikropartia jest przetwarzana przez Apache Spark. Wiele źródeł przesyłania strumieniowego, takich jak Kafka, obejmuje również czas przetwarzania systemu.

Znaki wodne w zapytaniach przesyłania strumieniowego są często definiowane względem czasu zdarzenia lub czasu przetwarzania źródła przesyłania strumieniowego. Zobacz Zastosuj znaki wodne, aby kontrolować progi przetwarzania danych.

Zarówno znaki wodne, jak i okna mogą być używane w połączeniu z transformWithState. Podobną funkcjonalność można zaimplementować w niestandardowej aplikacji stanowej, korzystając z TTL, czasomierzy oraz funkcjonalności MapState lub ListState.

Czas wygaśnięcia (TTL) dla typów stanów

Aby zapobiec błędom związanym z brakiem pamięci i usuwać nieaktualne wartości typów stanu, transformWithState obsługuje opcjonalną wartość czasu życia (TTL) dla każdej wartości typu stanu. Po wygaśnięciu TTL po cichu usuwa wartości stanu typu. TTL nie uruchamia handleExpiredTimer ani żadnej logiki niestandardowej. Aby uruchomić kod po wygaśnięciu stanu, użyj czasomierza.

Ważny

Jeśli nie wdrożysz mechanizmu TTL, musisz obsługiwać usuwanie stanu, aby uniknąć błędów związanych z brakiem pamięci.

Dla wszystkich typów stanów TTL resetuje się przy aktualizacji informacji o stanie. TTL jest egzekwowany dla każdej wartości danego typu stanu, przy czym dla każdego typu stanu obowiązują inne reguły:

  • Zmienne stanu są ograniczone do grupowania kluczy.
  • W przypadku obiektów ValueState tylko jedna wartość jest przechowywana na klucz grupowania. Czas wygaśnięcia ma zastosowanie do tej wartości.
  • W przypadku obiektów ListState lista może zawierać wiele wartości. TTL ma zastosowanie do każdej wartości na liście niezależnie.
    • Chociaż TTL dotyczy poszczególnych wartości w obiekcie ListState, jedynym sposobem aktualizacji pojedynczej wartości jest metoda put, która nadpisuje całą zawartość zmiennej ListState i resetuje TTL dla wszystkich wartości na liście.
  • Dla obiektów MapState każdy klucz mapy ma skojarzoną wartość stanu. Czas życia (TTL) jest stosowany niezależnie do każdej pary klucz-wartość w mapie.

Notatka

Timery umożliwiają definiowanie niestandardowej logiki wykraczającej poza usuwanie stanu, w tym emitowanie wierszy. Opcjonalnie możesz użyć czasomierzy, aby wyczyścić informacje o stanie dla danej wartości stanu oraz emitować wartości lub wyzwalać logikę warunkową. Zobacz Obsługa wygasłych czasomierzy.

Przykładowa aplikacja stanowa

W poniższym przykładzie zdefiniowano niestandardowy procesor stanowy, SimpleCounterProcessorw tym przykładowe zmienne stanu. SimpleCounterProcessor używa wartości ValueState, ListStatei MapState do zliczania wierszy dla każdego klucza grupowania.

Python (Pandas)

import pandas as pd
from pyspark.sql import Row
from pyspark.sql.streaming import StatefulProcessor, StatefulProcessorHandle
from pyspark.sql.types import StructType, StructField, IntegerType, StringType
from typing import Iterator

spark.conf.set("spark.sql.streaming.stateStore.providerClass","org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider")

output_schema = StructType(
    [
        StructField("id", StringType(), True),
        StructField("countAsString", StringType(), True),
    ]
)

class SimpleCounterProcessor(StatefulProcessor):
  def init(self, handle: StatefulProcessorHandle) -> None:
    value_state_schema = StructType([StructField("count", IntegerType(), True)])
    list_state_schema = StructType([StructField("count", IntegerType(), True)])
    self.value_state = handle.getValueState(stateName="valueState", schema=value_state_schema)
    self.list_state = handle.getListState(stateName="listState", schema=list_state_schema)
    # Schema can also be defined using strings and SQL DDL syntax
    self.map_state = handle.getMapState(stateName="mapState", userKeySchema="name string", valueSchema="count int")

  def handleInputRows(self, key, rows, timerValues) -> Iterator[pd.DataFrame]:
    # Seed the running total from state so the count accumulates across micro-batches
    count = self.value_state.get()[0] if self.value_state.exists() else 0
    for pdf in rows:
      list_state_rows = [(120,), (20,)] # A list of tuples
      self.list_state.put(list_state_rows)
      self.list_state.appendValue((111,))
      self.list_state.appendList(list_state_rows)
      pdf_count = pdf.count()
      count += pdf_count.get("value")
    self.value_state.update((count,)) # Count is passed as a tuple
    iter = self.list_state.get()
    list_state_value = next(iter)[0]
    value = count
    user_key = ("user_key",)
    if self.map_state.exists():
      if self.map_state.containsKey(user_key):
        value += self.map_state.getValue(user_key)[0]
    self.map_state.updateValue(user_key, (value,)) # Value is a tuple
    yield pd.DataFrame({"id": key, "countAsString": str(count)})

q = (df.groupBy("key")
  .transformWithStateInPandas(
    statefulProcessor=SimpleCounterProcessor(),
    outputStructType=output_schema,
    outputMode="Update",
    timeMode="None",
  )
  .writeStream...
)

Python (bazujący na wierszach)

from pyspark.sql import Row
from pyspark.sql.streaming import StatefulProcessor, StatefulProcessorHandle
from pyspark.sql.types import StructType, StructField, IntegerType, StringType
from typing import Iterator

spark.conf.set("spark.sql.streaming.stateStore.providerClass", "org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider")

output_schema = StructType(
  [
    StructField("id", StringType(), True),
    StructField("countAsString", StringType(), True),
  ]
)

class SimpleCounterProcessor(StatefulProcessor):
  def init(self, handle: StatefulProcessorHandle) -> None:
    value_state_schema = StructType([StructField("count", IntegerType(), True)])
    list_state_schema = StructType([StructField("count", IntegerType(), True)])
    self.value_state = handle.getValueState(stateName="valueState", schema=value_state_schema)
    self.list_state = handle.getListState(stateName="listState", schema=list_state_schema)
    self.map_state = handle.getMapState(stateName="mapState", userKeySchema="name string", valueSchema="count int")

  def handleInputRows(self, key, rows: Iterator[Row], timerValues) -> Iterator[Row]:
    # Seed the running total from state so the count accumulates across micro-batches
    count = self.value_state.get()[0] if self.value_state.exists() else 0
    for row in rows:
      list_state_rows = [(120,), (20,)]  # A list of tuples
      self.list_state.put(list_state_rows)
      self.list_state.appendValue((111,))
      self.list_state.appendList(list_state_rows)
      count += 1
    self.value_state.update((count,))  # Count is passed as a tuple
    iter_list = self.list_state.get()
    list_state_value = next(iter_list)[0]
    value = count
    user_key = ("user_key",)
    if self.map_state.exists():
      if self.map_state.containsKey(user_key):
        value += self.map_state.getValue(user_key)[0]
    self.map_state.updateValue(user_key, (value,))  # Value is a tuple
    yield Row(id=key[0], countAsString=str(count))

q = (
  df.groupBy("key")
    .transformWithState(
      statefulProcessor=SimpleCounterProcessor(),
      outputStructType=output_schema,
      outputMode="Update",
      timeMode="None",
    )
    .writeStream...
)

Skala

import org.apache.spark.sql.streaming._
import org.apache.spark.sql.{Dataset, Encoder, Encoders , DataFrame}
import org.apache.spark.sql.types._
import org.apache.spark.sql.functions._

spark.conf.set("spark.sql.streaming.stateStore.providerClass","org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider")

class SimpleCounterProcessor extends StatefulProcessor[String, (String, String), (String, String)] {
  @transient private var countState: ValueState[Int] = _
  @transient private var listState: ListState[Int] = _
  @transient private var mapState: MapState[String, Int] = _

  private val longEncoder = Encoders.scalaLong
  private val intEncoder = Encoders.scalaInt
  private val stringEncoder = Encoders.STRING

  override def init(
      outputMode: OutputMode,
      timeMode: TimeMode): Unit = {
    countState = getHandle.getValueState[Int]("countState",
      intEncoder, TTLConfig.NONE)
    listState = getHandle.getListState[Int]("listState",
      intEncoder, TTLConfig.NONE)
    mapState = getHandle.getMapState[String, Int]("mapState",
      stringEncoder, intEncoder, TTLConfig.NONE)
  }

  override def handleInputRows(
      key: String,
      inputRows: Iterator[(String, String)],
      timerValues: TimerValues): Iterator[(String, String)] = {
    var count = countState.getOption().getOrElse(0)
    for (row <- inputRows) {
      val listData = Array(120, 20)
      listState.put(listData)
      listState.appendValue(count)
      listState.appendList(listData)
      count += 1
    }
    val iter = listState.get()
    var listStateValue = 0
    if (iter.hasNext) {
      listStateValue = iter.next()
    }
    countState.update(count)
    var value = count
    val userKey = "userKey"
    if (mapState.exists()) {
      if (mapState.containsKey(userKey)) {
        value += mapState.getValue(userKey)
      }
    }
    mapState.updateValue(userKey, value)
    Iterator((key, count.toString))
  }
}

val q = spark
        .readStream
        .format("delta")
        .load("$srcDeltaTableDir")
        .as[(String, String)]
        .groupByKey(x => x._1)
        .transformWithState(
            new SimpleCounterProcessor(),
            TimeMode.None(),
            OutputMode.Update(),
        )
        .writeStream...

Uruchom przykład od początku do końca

Notatka

Przykłady, które można uruchomić na tej stronie, tworzą tabele w osobnym schemacie main.stateful_examples, dzięki czemu można je uruchamiać bez wpływu na Twoje istniejące dane. Jeśli nie masz uprawnień do tworzenia schematów w katalogu main , zmień katalog i schemat w przykładach na miejsce, gdzie możesz tworzyć tabele.

Procesor powyżej definiuje logikę stanową, ale nie rozpoczyna zapytania. Aby uruchomić SimpleCounterProcessor po skopiowaniu i wklejeniu, utwórz i wypełnij małą tabelę Delta Lake jako źródło danych strumieniowych, a następnie uruchom zapytanie, które zapisuje dane do ujścia w pamięci. Ten przykład używa Trigger.AvailableNow, więc zapytanie przetwarza wiersze początkowe i zatrzymuje się. Aby zasiać źródło i rozpocząć zapytanie, wykonaj następujące czynności:

import uuid

# Create a dedicated schema for the example tables
spark.sql("CREATE SCHEMA IF NOT EXISTS main.stateful_examples")

# Seed a small Delta table to use as the streaming source
spark.sql("DROP TABLE IF EXISTS main.stateful_examples.tws_counter_source")
spark.createDataFrame(
  [("a", "1"), ("a", "2"), ("a", "3"), ("b", "1"), ("b", "2")],
  "key string, value string",
).write.saveAsTable("main.stateful_examples.tws_counter_source")

df = spark.readStream.table("main.stateful_examples.tws_counter_source")

q = (
  df.groupBy("key")
    .transformWithState(
      statefulProcessor=SimpleCounterProcessor(),
      outputStructType=output_schema,
      outputMode="Update",
      timeMode="None",
    )
    .writeStream.format("memory")
    .queryName("counter_output")
    .option("checkpointLocation", f"/tmp/checkpoint_{uuid.uuid4()}")
    .trigger(availableNow=True)
    .start()
)

q.awaitTermination()

Po zakończeniu zapytania zobacz liczbę dla każdego klucza grupowania:

display(spark.sql("SELECT id, countAsString FROM counter_output ORDER BY id"))

Klucz a ma trzy wiersze, a klucz b dwa, więc zapytanie zwraca:

id  countAsString
a   3
b   2

Aby uzyskać więcej przykładów, zobacz Przykładowe aplikacje stanowe.

Notatka

W Pythonie wartości stanu są krotkami. Przekaż krotki do put i update, a krotek oczekuj z get.

Jeśli na przykład schemat elementu ValueState jest pojedynczą liczbą całkowitą:

current_value_tuple = value_state.get() # Returns the value state as a tuple
current_value = current_value_tuple[0]  # Extracts the first item in the tuple
new_value = current_value + 1           # Calculate a new value
value_state.update((new_value,))        # Pass the new value formatted as a tuple

Użyj tego podejścia także w przypadku elementów w ListState lub wartości w MapState.

Emituj wiersze

Należy użyć handleInputRows lub handleExpiredTimer, aby określić, jak transformWithState emituje wiersze dla każdego klucza grupowania. Zobacz Obsługa wierszy wejściowych i Obsługa wygasłych czasomierzy.

Niestandardowe aplikacje korzystające ze stanu nie narzucają żadnych założeń dotyczących sposobu wykorzystywania informacji o stanie. W przypadku danego warunku aplikacja może nie emitować wierszy, jednego wiersza ani wielu wierszy.

Notatka

Można zaimplementować wiele wartości stanu i zdefiniować wiele warunków emisji wierszy, ale wszystkie wiersze muszą używać tego samego schematu.

Python (Pandas)

Za pomocą transformWithStateInPandas zdefiniuj schemat wyjściowy przy użyciu słowa kluczowego outputStructType.

Emituj wiersze przy użyciu obiektu ramki danych biblioteki pandas i yield.

Opcjonalnie możesz yield pustą ramkę danych. Jeśli używasz trybu wyjściowego update i generujesz pustą ramkę danych, spowoduje to zaktualizowanie wartości klucza grupowania do null.

Python (bazujący na wierszach)

Za pomocą transformWithState zdefiniuj schemat wyjściowy przy użyciu słowa kluczowego outputStructType.

Generuj wiersze przy użyciu obiektu Row i yield.

Opcjonalnie możesz zwrócić pusty iterator. Jeśli używasz trybu wyjściowego update i zwracasz pusty iterator, spowoduje to zaktualizowanie wartości klucza grupowania na null.

Skala

W języku Scala emitujesz wiersze przy użyciu Iterator obiektu. Schemat pochodzi automatycznie ze schematu emitowanych wierszy.

Opcjonalnie możesz zwrócić pusty element Iterator. Jeśli używasz trybu wyjściowego update i emitujesz pusty element Iterator, spowoduje to zaktualizowanie wartości klucza grupowania do null.

Obsługa stanu początkowego

Opcjonalnie, możesz przekazać stan początkowy do pierwszej mikropartii.

Na przykład możesz użyć tego, aby:

  • Migrowanie istniejącego przepływu pracy do nowej aplikacji niestandardowej.
  • Uaktualnij operator stanowy, aby zmienić schemat lub logikę.
  • Napraw awarię, która nie może zostać automatycznie naprawiona i wymaga ręcznej interwencji.

Notatka

Użyj czytnika magazynu stanów, aby wysyłać zapytania o informacje o stanie z istniejącego punktu kontrolnego. Zobacz Odczytaj informacje o stanie zorganizowanego przesyłu strumieniowego.

Jeśli konwertujesz istniejącą tabelę Delta na aplikację stanową, przeczytaj tabelę przy użyciu spark.read.table("table_name") i przekaż wynikowy obiekt DataFrame. Opcjonalnie możesz wybrać lub zmodyfikować pola, aby były zgodne z nową aplikacją stanową.

Stan początkowy należy podać przy użyciu ramki danych z tym samym schematem klucza grupowania co wiersze wejściowe.

Notatka

Python używa handleInitialState do określenia stanu początkowego podczas definiowania StatefulProcessor. Scala używa odrębnej klasy StatefulProcessorWithInitialState.

Poniższy przykład zawiera licznik poszczególnych kluczy z istniejącej tabeli delty:

Python (bazujący na wierszach)

from pyspark.sql import Row
from pyspark.sql.streaming import StatefulProcessor, StatefulProcessorHandle
from pyspark.sql.types import StructType, StructField, IntegerType, StringType
from typing import Iterator

class CounterWithInitialState(StatefulProcessor):
  def init(self, handle: StatefulProcessorHandle) -> None:
    state_schema = StructType([StructField("count", IntegerType(), True)])
    self.count_state = handle.getValueState("countState", state_schema)

  def handleInitialState(self, key, initialState: Row, timerValues) -> None:
    self.count_state.update((initialState["count"],))

  def handleInputRows(self, key, rows: Iterator[Row], timerValues) -> Iterator[Row]:
    count = self.count_state.get()[0] if self.count_state.exists() else 0
    for _ in rows:
      count += 1
    self.count_state.update((count,))
    yield Row(id=key[0], count=count)

  def close(self) -> None:
    pass

output_schema = StructType([
  StructField("id", StringType(), True),
  StructField("count", IntegerType(), True),
])

import uuid

# Create a dedicated schema for the example tables
spark.sql("CREATE SCHEMA IF NOT EXISTS main.stateful_examples")

# Seed existing per-key counts to load as the initial state
spark.sql("DROP TABLE IF EXISTS main.stateful_examples.existing_counts")
spark.createDataFrame(
  [("x", 10)],
  "id string, count int",
).write.saveAsTable("main.stateful_examples.existing_counts")

# Seed a small Delta table to use as the streaming source
spark.sql("DROP TABLE IF EXISTS main.stateful_examples.tws_initial_source")
spark.createDataFrame(
  [("x", "a"), ("x", "b")],
  "id string, value string",
).write.saveAsTable("main.stateful_examples.tws_initial_source")

df = spark.readStream.table("main.stateful_examples.tws_initial_source")

# Load existing counts as initial state — must use the same grouping key as the input
initial_state = spark.read.table("main.stateful_examples.existing_counts").groupBy("id")

q = (
  df.groupBy("id")
    .transformWithState(
      statefulProcessor=CounterWithInitialState(),
      outputStructType=output_schema,
      outputMode="Update",
      timeMode="None",
      initialState=initial_state,
    )
    .writeStream.format("memory")
    .queryName("initial_state_output")
    .option("checkpointLocation", f"/tmp/checkpoint_{uuid.uuid4()}")
    .trigger(availableNow=True)
    .start()
)

q.awaitTermination()

# The initial state seeds "x" with 10, and the source adds two rows, so the count is 12
display(spark.sql("SELECT id, count FROM initial_state_output ORDER BY id"))

Skala

import org.apache.spark.sql.streaming._
import org.apache.spark.sql.Encoders

class CounterWithInitialState
    extends StatefulProcessorWithInitialState[String, (String, String), (String, String), (String, Int)] {

  @transient private var countState: ValueState[Int] = _

  override def init(outputMode: OutputMode, timeMode: TimeMode): Unit = {
    countState = getHandle.getValueState[Int]("countState", Encoders.scalaInt, TTLConfig.NONE)
  }

  override def handleInitialState(
      key: String, initialState: (String, Int), timerValues: TimerValues): Unit = {
    countState.update(initialState._2)
  }

  override def handleInputRows(
      key: String,
      rows: Iterator[(String, String)],
      timerValues: TimerValues): Iterator[(String, String)] = {
    val count = if (countState.exists()) countState.get() else 0
    val newCount = count + rows.size
    countState.update(newCount)
    Iterator((key, newCount.toString))
  }
}

// Load existing counts as initial state — must use the same grouping key as the input
val initialState = spark.read.table("existing_counts")
  .as[(String, Int)]
  .groupByKey(_._1)

val q = spark
  .readStream
  .format("delta")
  .load(srcDeltaTableDir)
  .as[(String, String)]
  .groupByKey(_._1)
  .transformWithState(
    new CounterWithInitialState(),
    TimeMode.None(),
    OutputMode.Update(),
    initialState,
  )
  .writeStream...

Przetwarzanie asynchroniczne (Beta)

Python transformWithState obsługuje przetwarzanie asynchroniczne z wykorzystaniem asyncio do jednoczesnego wykonywania operacji stanów i logiki użytkownika. Przetwarzanie asynchroniczne ma wyższą przepustowość niż synchroniczne i wymaga jedynie drobnych zmian w kodzie, bez żadnych zewnętrznych bibliotek asynchronicznych. Aby użyć przetwarzania asynchronicznego, zaimplementuj an AsyncStatefulProcessor zamiast synchronicznego StatefulProcessor. Zobacz Przetwarzanie asynchroniczne z transformWithState (Beta).

Użyj transformWithState w potokach Lakeflow

Użyj operatora transformWithState w potokach Lakeflow, aby implementować dowolną logikę z utrzymywaniem stanu w potokach strumieniowych w języku Python.

W tym celu wykonaj następujące czynności:

  1. Zdefiniuj schemat wyjściowy i logikę procesora stanowego dla dowolnych przekształceń stanowych. Przykłady można znaleźć w temacie Przykładowe aplikacje stanowe.
  2. Utwórz przepływ w potoku Lakeflow, który wywołuje operator transformWithState na ramce danych. Zobacz Samouczek: Utwórz swój pierwszy potok przy użyciu edytora Lakeflow Pipelines.
  3. Uruchom potok danych i zweryfikuj oraz sprawdź wyniki w tabeli docelowej lub w odbiorniku.

Aby zapoznać się z przykładem wykorzystującym transformWithState do monitorowania sygnałów czujnika, zobacz Przykład: monitorowanie sygnałów czujnika przy użyciu transformWithState.