Een aangepaste 'stateful' toepassing bouwen

U kunt transformWithState gebruiken om stateful streamingtoepassingen op te bouwen en om oplossingen met lage latentie en bijna real-time te implementeren. Met aangepaste stateful operators kunt u willekeurige stateful logica maken waarmee u nieuwe operationele use cases kunt bouwen die niet mogelijk zijn met traditionele Structured Streaming-verwerking.

Notitie

Voor statusafhankelijke bewerkingen, zoals aggregaties, deduplicatie en streaming-joins, raadt Databricks aan ingebouwde Structured Streaming-operators te gebruiken in plaats van aangepaste logica. Zie Wat is 'stateful streaming'?.

Databricks raadt aan om in plaats van verouderde operators te gebruiken transformWithState , zoals flatMapGroupsWithState en mapGroupsWithStatevoor willekeurige statustransformaties. Zie Verouderde willekeurige stateful operators.

Eisen

De transformWithState en transformWithStateInPandas operators hebben de volgende vereisten:

  • Beschikbaar in Databricks Runtime 16.2 en hoger.
    • Voor realtime-modus gebruikt u Databricks Runtime 17.3 LTS of hoger. Bekijk de real-time modus in Structured Streaming.
    • Voor de standaardtoegangsmodus is Python beschikbaar in Databricks Runtime 16.3 en hoger en Scala is beschikbaar in Databricks Runtime 17.3 en hoger.
  • RocksDB is de standaardprovider voor statusopslag in Databricks Runtime 17.3 en hoger.
    • Voor Databricks Runtime 17.2 en lager moet u de provider van het rocksDB-statusarchief configureren. Databricks raadt het inschakelen van RocksDB aan in de Spark-configuratie.

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

Wat is transformWithState?

De operator transformWithState past een aangepaste stateful processor toe op een Structured Streaming-query. U moet een aangepaste stateful processor implementeren om transformWithStatete gebruiken. Structured Streaming bevat API's voor het bouwen van uw stateful processor met behulp van Python, Scala of Java.

Gebruik transformWithState dit om aangepaste logica toe te passen op een groeperingssleutel. Hieronder wordt het ontwerp op hoog niveau beschreven:

  • Definieer een of meer statusvariabelen.
  • Statusinformatie blijft behouden voor elke groeperingssleutel. U kunt elke statusvariabele openen in door de gebruiker gedefinieerde code.
  • Voor elke verwerkte microbatch zijn alle rijen voor de sleutel beschikbaar in de vorm van een iterator.
  • Gebruik de StatefulProcessorHandle functie met timers en door de gebruiker gedefinieerde voorwaarden om te bepalen hoe u rijen verzendt.
  • Als u statusverlooptijd en statusgrootte wilt beheren, ondersteunen statuswaarden afzonderlijke TTL-definities (Time-to-Live).

Omdat transformWithState de ontwikkeling van schema's in het statusarchief wordt ondersteund, kunt u uw productietoepassingen herhalen en bijwerken zonder historische statusgegevens te verliezen. Na het bijwerken van het statusschema hoeft u geen rijen opnieuw te verwerken, wat code-implementaties en onderhoud vereenvoudigt. Bekijk de ontwikkeling van schema's in het statusarchief.

Belangrijk

Azure Databricks documentatie gebruikt transformWithState om zowel Python als Scala-implementaties te beschrijven:

  • PySpark ondersteunt zowel de op rijen gebaseerde transformWithState API als de pandas-operator transformWithStateInPandas .
  • Scala ondersteunt alleen de API op basis van transformWithState rijen.

De Scala- en Python-implementaties hebben transformWithState dezelfde mogelijkheden, maar met enkele verschillen in syntaxis.

Een StatefulProcessor definiëren

U definieert een stateful processor door de StatefulProcessor klasse uit te breiden en de methoden ervan te implementeren.

Spark geeft een StatefulProcessorHandle door aan de methode init van uw StatefulProcessor. Gebruik de handle om statusvariabelen te maken en met de statusopslag te interageren.

transformWithState ondersteunt drie statustypen: ValueState, ListStateen MapState. Elk type slaat de status voor elke groeperingssleutel op met behulp van een andere onderliggende gegevensstructuur.

Implementeer de volgende methoden om uw aangepaste logica te definiëren:

  • Implementeer handleInputRows om te bepalen hoe uw toepassing gegevens verwerkt, status bijwerken en rijen verzendt voor elke microbatch. Zie Invoerrijenverwerken.
  • Implementeer handleExpiredTimer om logica op basis van tijd uit te voeren, ongeacht of de groeperingssleutel nieuwe rijen in een microbatch ontvangt. Zie Verlopen timers verwerken.
  • U kunt eventueel implementeren handleInitialState om de status vooraf in te vullen voordat uw toepassing invoerrijen verwerkt. Zie Beginstatus Afhandelen.

In de volgende tabel worden de functionele gedragingen van deze methoden vergeleken:

Gedrag handleInputRows handleExpiredTimer
Toestandswaarden ophalen, plaatsen, bijwerken of wissen Ja Ja
Een timer maken of verwijderen Ja Ja
Rijen verzenden Ja Ja
Itereren over de rijen in de huidige microbatch Ja Nee
Triggerlogica op basis van verstreken tijd Nee Ja

U kunt zowel handleInputRows als handleExpiredTimer combineren om indien nodig complexe logica te implementeren.

U kunt bijvoorbeeld een toepassing implementeren die gebruikmaakt van handleInputRows om statuswaarden voor elke microbatch bij te werken en in de toekomst een timer van 10 seconden in te stellen. Als er geen extra rijen worden verwerkt, kunt u handleExpiredTimer gebruiken om de huidige waarden in de statusopslag uit te sturen. Als nieuwe rijen worden verwerkt voor de groeperingssleutel, kunt u de bestaande timer wissen en een nieuwe timer instellen.

StatefulProcessorHandle

In PySpark kunt u met de StatefulProcessorHandle klasse toegang krijgen tot functies die bepalen hoe uw code statusgegevens gebruikt.

Bij het initialiseren van een StatefulProcessor, moet u altijd de StatefulProcessorHandle variabele importeren en doorgeven aan de handle variabele. De variabele handle koppelt de lokale variabele in uw Python klasse aan de statusvariabele.

Notitie

Scala maakt gebruik van de methode getHandle.

aangepaste statustypen

U kunt meerdere statusobjecten implementeren in één stateful operator.

Kies een statustype op basis van de volledige toepassingslogica. U kunt bijvoorbeeld sessies bijhouden met een ValueState, gegroepeerd op basis van user_id en session_id. Of, om voorwaarden over meerdere sessies te evalueren, gebruikt u een MapState gegroepeerd op basis van user_id met session_id als mapsleutel.

Als uw statusobject een StructTypeobject gebruikt, moet u unieke namen definiëren voor elk veld in de struct voor het schema. Deze namen zijn zichtbaar bij het lezen van de statusopslag. Zie Lees informatie over de status van Gestructureerd Streamen.

In de volgende secties worden de statustypen beschreven die worden ondersteund door transformWithState:

ValueState

ValueState slaat een waarde op voor elke groeperingssleutel.

Een waardestatus kan complexe typen bevatten, zoals een struct of tuple. Voor ValueState moet u logica implementeren om de volledige waarde te vervangen.

De time-to-live voor een waardestatus wordt opnieuw ingesteld wanneer de waarde wordt bijgewerkt. Als u een bronsleutel voor ValueState verwerkt zonder de opgeslagen ValueState bij te werken, wordt de time-to-live niet opnieuw ingesteld.

ListState

ListState slaat een lijst op voor elke groeperingssleutel.

Een lijststatus is een verzameling waarden, die elk complexe typen kunnen bevatten. Elke waarde in een lijst heeft een eigen time-to-live.

U kunt items toevoegen aan een lijst door afzonderlijke items toe te voegen, een lijst met items toe te voegen of de hele lijst te overschrijven met een put. Als u time-to-live opnieuw wilt instellen, moet u een put bewerking gebruiken.

MapState

MapState slaat een kaart op voor elke groeperingssleutel. Kaarten zijn het Apache Spark-equivalent aan een Python woordenlijst (dict).

Een mapstatus is een verzameling unieke sleutels die elk aan een waarde zijn gekoppeld, waarvan elk complexe typen kan bevatten. Elk sleutel-waardepaar in een map heeft zijn eigen time-to-live.

U kunt de waarde van een specifieke sleutel bijwerken of u kunt een sleutel en de bijbehorende waarde verwijderen. U kunt een afzonderlijke waarde retourneren met behulp van de sleutel, alle sleutels weergeven, alle waarden weergeven of een iterator retourneren om te werken met de volledige set sleutel-waardeparen in de kaart.

Belangrijk

Met groeperingssleutels worden de velden beschreven die zijn opgegeven in de GROUP BY-component van de Structured Streaming-query. Map-statussen kunnen een willekeurig aantal sleutel-waardeparen bevatten bij een groeperingssleutel.

Als uw query bijvoorbeeld gebruikmaakt van GROUP BY user_id en u voor elke session_id een toewijzing wilt definiëren, is uw groeperingssleutel user_id en is de MapState-sleutel session_id:

Python
class SessionTracker(StatefulProcessor):
  def init(self, handle: StatefulProcessorHandle) -> None:
    self.sessions = handle.getMapState("sessions", StringType(), LongType())

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

  def close(self) -> None:
    pass

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

class SessionTracker extends StatefulProcessor[String, Event, Row] {
  @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[Row] = {
    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())

Een aangepaste statusvariabele maken in de StatefulProcessor

Wanneer u uw StatefulProcessorinitialiseert, maakt u een lokale variabele voor elk statusobject waarmee u kunt communiceren met statusobjecten in uw aangepaste logica. Definieer en initialiseer statusvariabelen door de ingebouwde methode in init de StatefulProcessor klasse te overschrijven.

U kunt een willekeurig aantal statusobjecten definiëren met behulp van de getValueState, getListStateen getMapState methoden in uw StatefulProcessor.

Elk statusobject moet het volgende hebben:

  • Een unieke naam
  • Een schema
    • In Python moet u het schema opgeven.
    • In Scala kunt u een Encoder doorgeven om het statusschema op te geven.

U kunt desgewenst ook een TTL-duur (Time to Live) opgeven in milliseconden. Als u een kaartstatus implementeert, moet u een afzonderlijke schemadefinitie opgeven voor de kaartsleutels en de waarden.

Notitie

De StatefulProcessor logica wordt afzonderlijk verwerkt voor het opvragen, bijwerken en verzenden van statusinformatie. Zie Uw statusvariabelen gebruiken in methoden met aangepaste logica.

Uw statusvariabelen gebruiken in methoden met aangepaste logica

Statusobjecten hebben methoden voor het ophalen van status, het bijwerken van bestaande statusgegevens en het wissen van de huidige status.

Elke groeperingssleutel heeft toegewezen statusinformatie.

  • De StatefulProcessor genereert rijen op basis van je aangepaste logica en het opgegeven uitvoerschema. Zie Rijen verzenden.
  • Gebruik de statestore lezer om toegang te krijgen tot waarden in het statusarchief. Deze lezer is bedoeld voor batchworkloads en is niet bedoeld voor workloads met lage latentie. Zie Lees informatie over de status van Gestructureerd Streamen.
  • Logica die is opgegeven met handleInputRows wordt alleen uitgevoerd als rijen voor de sleutel aanwezig zijn in een micro-batch. Zie Invoerrijenverwerken.
  • Gebruik handleExpiredTimer dit om op tijd gebaseerde logica te implementeren die niet afhankelijk is van het observeren van rijen die moeten worden geactiveerd. Zie Verlopen timers verwerken.

Notitie

Statusobjecten worden geïsoleerd door sleutels te groeperen met de volgende implicaties:

  • Statuswaarden kunnen niet worden beïnvloed door rijen die zijn gekoppeld aan een andere groeperingssleutel.
  • Je kan geen logica implementeren die afhankelijk is van het vergelijken van waarden of het bijwerken van de status tussen groeperingssleutels.

U kunt waarden in een groeperingssleutel vergelijken. Gebruik een MapState om logica te implementeren met een tweede sleutel die uw aangepaste logica kan gebruiken. Als u bijvoorbeeld op basis van user_id groepeert en ip_address gebruikt voor uw MapState-sleutel, kunt u gelijktijdige gebruikerssessies bijhouden.

Geavanceerde overwegingen voor het omgaan met status

Statusupdates zijn fouttolerant. Als een taak vastloopt voordat een microbatch is verwerkt, gebruikt de nieuwe poging de waarde van de laatste geslaagde microbatch.

Voor geoptimaliseerde prestaties raadt Databricks u aan om alle waarden in de iterator te verwerken voor een bepaalde sleutel en updates door te voeren in één schrijfbewerking. Wanneer u naar een statusvariabele schrijft, wordt hiermee een schrijfbewerking naar RocksDB geactiveerd.

Statuswaarden hebben geen standaardwaarden. Als voor uw logica bestaande statusgegevens moeten worden gelezen, gebruikt u de exists methode.

Als u logica voor null-status wilt implementeren, MapState kunt u met variabelen controleren op afzonderlijke sleutels of alle sleutels weergeven.

Invoerrijen verwerken

Gebruik de handleInputRows methode om te definiëren hoe uw toepassing rijen verwerkt en statuswaarden bijwerken. Deze methode wordt telkens uitgevoerd wanneer uw Structured Streaming-query rijen verwerkt voor een groeperingssleutel.

Voor de meeste stateful toepassingen die zijn geïmplementeerd met transformWithState, wordt de kernlogica gedefinieerd met behulp van handleInputRows.

Voor elke microbatch-update die wordt verwerkt, zijn alle rijen in de microbatch voor een bepaalde groeperingssleutel beschikbaar met behulp van een iterator. Door de gebruiker gedefinieerde logica kan communiceren met alle rijen uit de huidige microbatch en waarden in de statestore.

Verlopen timers verwerken

Gebruik de handleExpiredTimer methode om aangepaste logica te implementeren op basis van verstreken tijd.

Binnen een groeperingssleutel worden timers uniek geïdentificeerd door hun tijdstempel.

Wanneer een timer verloopt, wordt het resultaat bepaald door de logica die in uw toepassing is geïmplementeerd. Veelvoorkomende patronen zijn:

  • Informatie verzenden die is opgeslagen in een statusvariabele.
  • Opgeslagen statusgegevens verwijderen.
  • Een nieuwe timer maken.

Verlopen timers gaan af, zelfs als er in een micro-batch geen rijen voor de bijbehorende sleutel worden verwerkt.

De tijdmodus opgeven

Wanneer u uw StatefulProcessor aan transformWithState doorgeeft, moet u de tijdmodus opgeven met de parameter timeMode.

De volgende opties worden ondersteund:

Tijdmodus Beschrijving
ProcessingTime Timers en TTL worden beide ondersteund en geëvalueerd op basis van de wandkloktijd wanneer Apache Spark elke microbatch verwerkt. Gebruik ProcessingTime deze optie als u wilt dat timers met een vast interval worden geactiveerd ten opzichte van wanneer rijen worden verwerkt, ongeacht tijdstempels in de gegevens.
EventTime Timers worden ondersteund en geëvalueerd op basis van het event-time-watermerk. Het watermerk gaat verder naarmate Apache Spark tijdstempels in de invoergegevens bekijkt. TTL wordt niet ondersteund met EventTime. Gebruik EventTime wanneer uw gegevens tijdstempels bevatten en u wilt dat timers worden geactiveerd op basis van de voortgang van deze tijdstempels. Wanneer u deze EventTimegebruikt, moet u ook de eventTimeColumnName parameter opgeven. Zie eventTimeColumnName.
NoTime of TimeMode.None() Timers en TTL worden niet ondersteund. Gebruik NoTime wanneer uw statusvolle toepassing geen tijdsgebaseerde logica vereist.

eventTimeColumnName

Wanneer u de EventTime tijdmodus gebruikt, geeft de eventTimeColumnName parameter de naam van de kolom in het uitvoerschema op die de tijdstempel van de gebeurtenis bevat. Apache Spark gebruikt deze kolom om het watermerk door te geven aan de uitvoerstroom, waardoor juiste downstreambewerkingen op basis van tijd worden ingeschakeld.

Python

eventTimeColumnName is een aanvullend argument voor transformWithState of transformWithStateInPandas:

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

transformWithState accepteert eventTimeColumnName in plaats van timeMode. Deze aanpak gebruikt altijd de modus EventTime:

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

Ingebouwde timerwaarden

Databricks raadt af om de systeemklok in uw aangepaste stateful toepassing aan te roepen, omdat dit kan leiden tot onbetrouwbare herhalingen bij het mislukken van een taak. Gebruik de methoden in de klasse TimerValues wanneer u toegang moet hebben tot de verwerkingstijd of het watermerk:

TimerValues Beschrijving
getCurrentProcessingTimeInMs Retourneert de tijdstempel van de verwerkingstijd voor de huidige batch in milliseconden sinds epoch.
getCurrentWatermarkInMs Retourneert de tijdstempel van het watermerk voor de huidige batch in milliseconden sinds epoch.

Notitie

De verwerkingstijd beschrijft de tijd die de microbatch door Apache Spark wordt verwerkt. Veel streamingbronnen, zoals Kafka, bevatten ook systeemverwerkingstijd.

Watermerken voor streamingquery's worden vaak gedefinieerd op basis van gebeurtenistijd of de verwerkingstijd van de streamingbron. Zie Watermerken toepassen om drempelwaarden voor gegevensverwerking te beheren.

Zowel watermerken als vensters kunnen worden gebruikt in combinatie met transformWithState. U kunt vergelijkbare functionaliteit implementeren in uw aangepaste stateful toepassing door gebruik te maken van TTL, timers en MapState of ListState functionaliteit.

Time-to-live (TTL) voor statustypen

Als u fouten met onvoldoende geheugen wilt voorkomen en verouderde statustypewaarden wilt verwijderen, transformWithState ondersteunt u een optionele TTL-waarde (Time to Live) voor elke statustypewaarde. Na het verlopen verwijdert TTL stilzwijgend waarden van het statetype. TTL voert handleExpiredTimer of aangepaste logica niet uit. Als u code wilt uitvoeren wanneer de status verloopt, gebruikt u in plaats daarvan een timer.

Belangrijk

Als u TTL niet implementeert, moet u statusverwijdering afhandelen om geheugenfouten te voorkomen.

Voor alle statustypen wordt TTL opnieuw ingesteld bij het bijwerken van statusgegevens. TTL wordt afgedwongen voor elke statustypewaarde, met verschillende regels voor elk statustype:

  • Statusvariabelen zijn gericht op het groeperen van sleutels.
  • Voor ValueState objecten wordt slechts één waarde per groeperingssleutel opgeslagen. TTL is van toepassing op deze waarde.
  • Voor ListState objecten kan de lijst veel waarden bevatten. TTL is onafhankelijk van toepassing op elke waarde in een lijst.
    • Hoewel TTL is gericht op afzonderlijke waarden in een ListState, is de enige manier om een afzonderlijke waarde bij te werken met de put methode, die de volledige inhoud van de ListState variabele overschrijft en TTL opnieuw instelt voor alle waarden in de lijst.
  • Voor MapState objecten heeft elke kaartsleutel een bijbehorende statuswaarde. TTL is onafhankelijk van toepassing op elk sleutel-waardepaar in een kaart.

Notitie

Met timers kunt u aangepaste logica definiëren buiten verwijdering van statussen, inclusief het verzenden van rijen. U kunt timers ook gebruiken om statusinformatie voor een bepaalde statuswaarde te wissen en waarden te verzenden of voorwaardelijke logica te activeren. Zie Verlopen timers verwerken.

Voorbeeld van een "stateful" toepassing

In het volgende voorbeeld wordt een aangepaste stateful processor gedefinieerd, SimpleCounterProcessorinclusief voorbeeldstatusvariabelen. SimpleCounterProcessor gebruikt ValueState, ListStateen MapState om rijen te tellen voor elke groeperingssleutel.

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]:
    count = 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 (op regels gebaseerd)

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]:
    count = 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, countAsString=str(count))

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

Scala

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

Zie Voorbeeld van stateful toepassingenvoor meer voorbeelden.

Notitie

In Python zijn statuswaarden tuples. Geef tuples door aan put en update, en verwacht tuples uit get.

Als bijvoorbeeld het schema voor uw ValueState uit één geheel getal bestaat:

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

Gebruik deze methode ook voor items in een ListState of waarden in een MapState .

Rijen verzenden

U moet handleInputRows of handleExpiredTimer gebruiken om te definiëren hoe transformWithState rijen produceert voor elke groeperingssleutel. Zie Invoerrijen verwerken en Verlopen timers verwerken.

Aangepaste toepassingen met status gaan niet uit van hoe statusinformatie wordt gebruikt. Voor een bepaalde voorwaarde kan de toepassing geen rijen, één rij of veel rijen verzenden.

Notitie

U kunt meerdere statuswaarden implementeren en meerdere voorwaarden definiëren voor het verzenden van rijen, maar alle rijen moeten hetzelfde schema gebruiken.

Python (Pandas)

Met transformWithStateInPandas, definieer je je uitvoerschema met het outputStructType sleutelwoord.

Rijen verzenden met behulp van een Pandas DataFrame-object en yield.

U kunt yield eventueel een leeg DataFrame gebruiken. Als u de uitvoermodus gebruikt update en een leeg DataFrame verzendt, worden de waarden voor de groeperingssleutel bijgewerkt null.

Python (op regels gebaseerd)

Met transformWithState, definieer je je uitvoerschema met het outputStructType sleutelwoord.

Genereer rijen met een Row-object en yield.

U kunt eventueel een lege iterator retourneren. Als u de uitvoermodus gebruikt update en een lege iterator verzendt, worden de waarden voor de groeperingssleutel bijgewerkt null.

Scala

In Scala verzendt u rijen met behulp van een Iterator object. Het schema wordt automatisch afgeleid van het schema van de verzonden rijen.

U kunt eventueel een lege Iteratorwaarde retourneren. Als u de update-uitvoermodus gebruikt en een lege Iterator genereert, worden de waarden voor de groeperingssleutel bijgewerkt naar null.

Begintoestand verwerken

U kunt desgewenst een initiële status doorgeven aan de eerste microbatch.

U kunt dit bijvoorbeeld gebruiken om het volgende te doen:

  • Een bestaande werkstroom migreren naar een nieuwe aangepaste toepassing.
  • Werk een stateful operator bij om uw schema of logica te wijzigen.
  • Herstel een fout die niet automatisch kan worden hersteld en waarvoor handmatige tussenkomst is vereist.

Notitie

Gebruik de lezer van het statusarchief om statusgegevens van een bestaand controlepunt op te vragen. Zie Lees informatie over de status van Gestructureerd Streamen.

Als u een bestaande Delta-tabel converteert naar een stateful applicatie, gebruik dan spark.read.table("table_name") om de tabel te lezen en geef het resulterende DataFrame door. U kunt desgewenst velden selecteren of wijzigen om aan te passen aan uw nieuwe statusgevoelige toepassing.

U geeft een initiële status op met behulp van een DataFrame met hetzelfde groeperingssleutelschema als de invoerrijen.

Notitie

Python gebruikt handleInitialState om de initiële status op te geven tijdens het definiëren van een StatefulProcessor. Scala maakt gebruik van de afzonderlijke klasse StatefulProcessorWithInitialState.

In het volgende voorbeeld wordt een teller per sleutel geïnitialiseerd vanuit een bestaande Delta-tabel:

Python (op regels gebaseerd)

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),
])

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

q = (
  df.groupBy("id")
    .transformWithState(
      statefulProcessor=CounterWithInitialState(),
      outputStructType=output_schema,
      outputMode="Update",
      timeMode="None",
      initialState=initial_state,
    )
    .writeStream...
)

Scala

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

Gebruik transformWithState in Lakeflow-pijplijnen

Gebruik de transformWithState operator in Lakeflow-pijplijnen om willekeurige stateful logica in uw streamingpijplijnen te implementeren met behulp van Python.

Voer hiervoor de volgende stappen uit:

  1. Definieer het uitvoerschema en de stateful-processorlogica voor uw willekeurige stateful-transformaties. Zie Voorbeeld van stateful toepassingen.
  2. Maak een Lakeflow-pijplijnstroom die de transformWithState operator aanroept op een DataFrame. Zie Zelfstudie: Uw eerste pijplijn maken met behulp van de Lakeflow Pipelines Editor.
  3. Voer uw pijplijn uit en valideer de resultaten op de doeltabel of de doellocatie.

Zie voor een voorbeeld dat transformWithState gebruikt om sensor heartbeats te bewaken, Voorbeeld: Gebruik transformWithState om sensor heartbeats te bewaken.