Stateful processing met gestructureerde streaming

Spark Structured Streaming ondersteunt zowel statusloze bewerkingen als bewerkingen met status. Stateloze bewerkingen, zoals map, filter, en select, verwerken elke invoerrij onafhankelijk. Toestandsafhankelijke bewerkingen, zoals aggregaties, samenvoegingen en deduplicatie, bewaren informatie tussen microbatches zodat Spark resultaten kan bijwerken wanneer nieuwe gebeurtenissen binnenkomen.

In Microsoft Fabric gebruik je dezelfde open-source Apache Spark Structured Streaming-concepten met Fabric Spark-runtimes, notebooks, Spark-jobdefinities, lakehouse-tabellen en andere Fabric data engineering-ervaringen.

Staatswinkel en controleposten

Stateful streaming-queries gebruiken een state store om tussenliggende data tussen microbatches te bewaren. Bijvoorbeeld, een doorlopende telling per apparaat moet de huidige telling voor elke apparaatsleutel bijhouden voordat deze de volgende microbatch kan verwerken.

Spark blijft in de staat bij het streamingcheckpoint. Het checkpoint houdt offsets, zoekvoortgang en state store-gegevens bij. Als een query stopt en opnieuw wordt gestart met dezelfde checkpointlocatie, laadt Spark de toestand opnieuw en hervat deze vanaf de laatst vastgelegde voortgang.

Important

Gebruik een duurzame checkpointlocatie en deel niet hetzelfde checkpoint tussen verschillende streamingqueries. Het toestandsschema en het queryplan maken deel uit van de checkpoint-toestand, dus grote wijzigingen in stateful logica vereisen vaak een nieuw checkpoint.

De status groeit naarmate Spark meer sleutels, vensters of gebufferde join-rijen bijhoudt. Beperk de omvang van de status met watermerken, levensduur (TTL) en zorgvuldig gekozen sleutels, zodat het checkpoint niet onbeperkt groeit.

RocksDB-statusopslag inschakelen

Standaard gebruikt Spark een door HDFS ondersteunde state store waarvan de werkstatus wordt onderhouden in door JVM beheerd geheugen en waarvan de duurzame versies in het checkpoint worden opgeslagen. Schakel de RocksDB state store-provider in zodat Spark de actieve status in het native geheugen en op de lokale schijf van elke executor beheert in plaats van de JVM-heap. RocksDB vermindert de druk op de JVM-heap en zorgt voor voorspelbaarder geheugengebruik. Het bevoordeelt stateful query’s rechtstreeks, en het is een goede standaardkeuze, zelfs voor query’s die weinig status bijhouden.

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

Stel de provider van de statusopslag in vóór de query voor het eerst wordt uitgevoerd en houd deze consistent bij opnieuw opstarten, omdat deze deel uitmaakt van de query waarvan een checkpoint is gemaakt.

Checkpointing van RocksDB-wijzigingslogboeken inschakelen

Changelog checkpointing is een Apache Spark-optimalisatie voor de RocksDB state store. Het uploadt incrementele toestandswijzigingen en maakt periodieke momentopnamen, wat de checkpoint-latentie voor grote toestanden kan verminderen.

spark.conf.set(
    "spark.sql.streaming.stateStore.rocksdb.changelogCheckpointing.enabled",
    True
)

Changelog checkpointing is achterwaarts compatibel met RocksDB snapshot checkpoints. Je kunt het inschakelen voor een bestaande RocksDB-query zonder de status te verwijderen, maar je moet de query opnieuw starten voordat de instelling van kracht wordt.

Vensteraggregaties

Vensteraggregaties groeperen gebeurtenissen op gebeurtenistijd in plaats van verwerkingstijd. Ze helpen je bij het berekenen van meetwaarden zoals aantallen, sommen en gemiddelden over tijdsbereiken.

Veelvoorkomende raamtypen zijn:

  • Vallende ramen. Een tuimelend venster heeft een vaste grootte en overlapt niet met andere vensters. Bijvoorbeeld, een tumblingvenster van vijf minuten telt gebeurtenissen van 10:00 tot en met 10:05 en begint vervolgens een nieuwe telling voor 10:05 tot en met 10:10.
  • Schuivende vensters. Een schuifvenster heeft een vaste grootte en begint op een regelmatig schuifinterval. Bijvoorbeeld, een venster van 10 minuten dat elke minuut schuift, geeft je een rollende telling van 10 minuten.
  • Sessievensters. Een sessievenster groepeert gebeurtenissen die dicht bij elkaar aankomen voor dezelfde sleutel. Bijvoorbeeld, een sessievenster met een onderbreking van 15 minuten sluit een gebruikerssessie nadat die gebruiker 15 minuten geen gebeurtenissen heeft gehad.

Het volgende voorbeeld telt apparaatgebeurtenissen in tumbling-vensters van vijf minuten.

from pyspark.sql import functions as F

windowed_counts = (
    events
    .withWatermark("event_time", "10 minutes")
    .groupBy(
        F.window("event_time", "5 minutes"),
        F.col("device_id")
    )
    .count()
)

Voor schuifvensters voeg je een schuifduur toe. Bereken window("event_time", "10 minutes", "1 minute") bijvoorbeeld elke minuut een venster van 10 minuten. Gebruik voor sessievensters session_window(event_time, '15 minutes') in Spark SQL of session_window("event_time", "15 minutes") in de DataFrame API.

Sessievenster-aggregaties ondersteunen geen update uitvoermodus. Gebruik een outputmodus die wordt ondersteund door de query en sink, en test de resulterende emissietiming met late data.

Watermarks

Een watermark vertelt Spark hoeveel vertraging je verwacht voordat event-timegegevens binnenkomen. Spark gebruikt het watermerk om te bepalen wanneer oude staat verwijderd kan worden en wanneer late rijen te oud zijn om een resultaat bij te werken.

Gebruik withWatermark vóór een statusbehoudende event-time-bewerking:

events_with_watermark = events.withWatermark("event_time", "10 minutes")

Dit voorbeeld maakt het mogelijk dat data tot 10 minuten te laat aankomt, afhankelijk van de event_time kolom. Spark kan uiteindelijk de status verliezen voor ramen die ouder zijn dan het watermerk.

Note

Spark laat geen data vallen die binnen de geconfigureerde watermarkvertraging aankomt. Data ouder dan het watermerk kan nog steeds verwerkt worden, maar Spark garandeert dat niet. Het opschonen van de status gebeurt ook asynchroon in plaats van onmiddellijk wanneer de watermark vooruitgaat.

Kies watermerkvertragingen op basis van het echte brongedrag. Een te korte vertraging kan geldige, vertraagd binnenkomende gegevens verwerpen. Een te lange vertraging behoudt meer status en maakt de checkpoint groter.

Watermerken worden verplaatst wanneer Spark nieuwe input verwerkt. Als er geen gegevens binnenkomen, wordt een watermerk mogelijk niet verder verplaatst, waardoor vensteruitvoer, niet-gematchte rijen van outer joins en statusopschoning kunnen worden uitgesteld totdat een latere microbatch gegevens ontvangt.

Stroom-stroom samenvloeiingen

Een stream-stream join combineert twee streaming-inputs. Spark buffert rijen van beide kanten totdat het kan bepalen of overeenkomende rijen aankomen.

Inner joins kunnen zonder watermerken werken, maar een onbegrensde status kan onbeperkt blijven groeien. Voeg watermerken en een gebeurtenistijdbereik toe zodat Spark oude gebufferde rijen kan verwijderen.

joined_events = (
    impressions.withWatermark("impression_time", "10 minutes")
    .join(
        clicks.withWatermark("click_time", "10 minutes"),
        """
        impression_id = click_impression_id AND
        click_time >= impression_time AND
        click_time <= impression_time + interval 5 minutes
        """,
        "inner"
    )
)

Outer joins vereisen genoeg informatie tijdens het evenement zodat Spark weet wanneer er geen toekomstige match kan arriveren. Gebruik watermerken en een tijdsbereikvoorwaarde. Voor een linker buitenste verbinding, geef je een watermerk aan de rechterkant die mogelijk de nulleerbare match oplevert. Plaats bij een rechter outer join een watermark op de linkerzijde. Voor een volledige buitenste aansluiting geef je een watermerk aan minstens één zijde; watermerk beide kanten wanneer Spark de status voor beide inputs moet opschonen.

Important

Vermijd ruime stream-stream-joins zonder tijdsbereik. Spark moet meer gebuffreerde rijen behouden, en de staat kan groeien zonder een duidelijk opruimpunt.

Niet-gematchte outer-join-rijen worden niet onmiddellijk gegenereerd. Spark wacht totdat het watermerk en de tijdsbereikvoorwaarde aantonen dat een toekomstige match niet meer kan arriveren.

Wanneer een query meerdere invoerstromen combineert, leidt Spark één querywatermerk af van de invoerwatermerken. Standaard regelt de langzaamste invoer de voortgang, zodat Spark geen data voortijdig uit die stroom laat vallen. Een vastgelopen invoer kan daarom het opschonen van de status en de uitvoer van de hele query vertragen.

Deduplicatie in stromen

Streaming deduplicatie verwijdert herhaalde gebeurtenissen door sleutels te onthouden die Spark al heeft verwerkt. Gebruik dit wanneer bronnen dezelfde gebeurtenis opnieuw kunnen verzenden, zoals bij nieuwe pogingen vanuit een berichtensysteem.

dropDuplicates slaat de kolomwaarden op die worden gebruikt voor duplicatendetectie. Zonder watermerk moet Spark die waarden onthouden gedurende de hele zoekopdracht. Om Spark verouderde dropDuplicates status te laten verwijderen, voeg de watermerkkolom toe aan de duplicaatsleutel.

deduplicated_events = (
    events
    .withWatermark("event_time", "10 minutes")
    .dropDuplicates(["event_id", "event_time"])
)

Het opnemen van de gebeurtenistijd betekent dat twee records met dezelfde gebeurtenis-ID maar verschillende tijdstempels geen duplicaten zijn. Gebruik dropDuplicatesWithinWatermark wanneer je runtime het ondersteunt en alleen de event ID een duplicaat definieert. Spark bewaart vervolgens elke gebeurtenis-ID voor de watermerkhorizon:

deduplicated_events = (
    events
    .withWatermark("event_time", "10 minutes")
    .dropDuplicatesWithinWatermark(["event_id"])
)

Kies sleutels die een logische gebeurtenis uniek identificeren. Als de sleutel te breed is, kan Spark afzonderlijke gebeurtenissen verwijderen. Als de sleutel te smal is, kunnen duplicaten erdoorheen gaan.

Aangepaste toestandslogica

Ingebouwde aggregaties, join-bewerkingen en deduplicatie ondersteunen veel statusvolle workloads. Gebruik verwerking met status naar keuze wanneer je aangepaste logica per sleutel nodig hebt, zoals toestandsmachines, waarschuwingsonderdrukking, time-outafhandeling of meerstapsverrijking.

De transformWithState operator vervangt de legacy mapGroupsWithState en flatMapGroupsWithState API's door aangepaste toestandslogica. Je groepeert rijen per sleutel, implementeert een stateful processor en beheert toestandsvariabelen, timers, outputmodus en tijdmodus. Apache Spark werd geïntroduceerd transformWithState in Spark 4.0, en Fabric Runtime 2.0 bevatte het tot en met Spark 4.1.

Note

transformWithState is een nieuwere API. Valideer je code met de Fabric-runtimeversie die je productiewerklast uitvoert.

Gebruik TTL voor aangepaste status wanneer de bedrijfslogica vervaldatum toestaat. TTL helpt Spark om verouderde sleutels te verwijderen zonder te wachten op een specifiek invoerevent.

Debugtools voor status

Fabric Runtime 2.0 bevat Apache Spark 4.1 en de State Data Source for Structured Streaming, die Apache Spark introduceerde in Spark 4.0. Gebruik het om met een aparte batchquery de inhoud van de state store uit een checkpoint te inspecteren.

Note

De State Data Source is experimenteel binnen Apache Spark. De opties en het uitvoerschema kunnen veranderen.

De statestore lezer kan de compatibele status inspecteren in een bestaand checkpoint, inclusief een checkpoint dat is aangemaakt door een eerdere Spark-runtime. Voer de aparte batchlezer uit met Fabric Runtime 2.0 of later; je hoeft het checkpoint niet opnieuw aan te maken of de originele query eerst opnieuw te starten.

state_df = (
    spark.read
    .format("statestore")
    .load("Files/checkpoints/device-counts")
)

state_df.printSchema()
display(state_df)

De state-metadata bron is een afzonderlijk hulpmiddel om operator-ID's, winkelnamen en beschikbare batch-ID's te vinden. Spark maakt deze metadata alleen aan terwijl de streamingquery draait op Spark 4.0 of later. Voor een ouder checkpoint hervat je de oorspronkelijke query vanaf dat checkpoint op Fabric Runtime 2.0 voordat je state-metadata gebruikt.

state_metadata_df = (
    spark.read
    .format("state-metadata")
    .load("Files/checkpoints/device-counts")
)

display(state_metadata_df)

Voor queries met meerdere stateful operators, gebruik metadata om de operator te identificeren die je wilt inspecteren. Als je al de operator ID en de winkelnaam kent, kun je de metadatabron overslaan en die opties direct aan de statestore lezer doorgeven. Voor transformWithState, specificeer de naam van de toestandvariabele wanneer je de staat leest.

Important

Behandel de status van het checkpoint als operationele gegevens. Bewerk checkpointbestanden niet handmatig. Gebruik de State Data Source voor inspectie en gebruik querycodewijzigingen, watermerken of TTL om het gedrag van de toestand te beheersen.

Beste werkwijzen voor afgebakende toestand

Volg deze praktijken om stateful streaming-zoekopdrachten betrouwbaar te houden:

  • Voeg event-time-watermerken toe voor vensteraggregaties, joins tussen streams en deduplicatie wanneer de regels voor te laat binnengekomen gegevens dat toestaan.
  • Gebruik tijdbereikcondities voor stream-stream joins zodat Spark gebufferde rijen kan verwijderen.
  • Kies toetsen met de juiste kardinaliteit. Sleutels met zeer hoge kardinaliteit vergroten de toestandgrootte, terwijl te brede sleutels niet-gerelateerde gebeurtenissen kunnen mengen.
  • Schakel de RocksDB-statusopslagprovider en changelog-checkpointing in om heapdruk, garbagecollectionpauzes en checkpointlatentie te verminderen.
  • Gebruik TTL wanneer transformWithState de aangepaste staat niet voor altijd hoeft te leven.
  • Houd de locaties van checkpoints stabiel gedurende de levensduur van een query en gebruik een nieuw checkpoint voor incompatibele veranderingen in het toestandsschema.
  • Bewaak statusgerelateerde metriek, invoersnelheid, verwerkingssnelheid en de groei van checkpoints tijdens productieruns.
  • Test met realistische late gegevens, dubbele gegevens en scenario's met herstarts voordat je een query met status implementeert.