Veelvoorkomende Structured Streaming-patronen

Spark Structured Streaming ondersteunt verschillende manieren om één logische stroom naar meerdere bestemmingen te routeren, stromen te correleren en de voortgang van querys te observeren. Dit artikel toont veelvoorkomende patronen voor Microsoft Fabric, waaronder fan-out met foreachBatch, onafhankelijke streamingqueries, streamstream-joins, duurzame Delta-tabellen en StreamingQueryListener.

De voorbeelden gaan ervan uit dat events een stromend DataFrame is met order_id, amount, en event_time kolommen.

Important

start() en toTable() start een streamingquery asynchroon en geeft een StreamingQuery handle terug. Een actieve query voorkomt niet dat een getriggerde notebook- of Spark-taakdefinitie een terminaltoestand bereikt. Roep query.awaitTermination() elke query aan waarop de getriggerde taak moet wachten. Bij meerdere altijd-actieve query's blokkeert spark.streams.awaitAnyTermination() totdat één query stopt, zodat de taak de beëindiging kan afhandelen.

Kies een uitwaaierpatroon

Kies het patroon op basis van de isolatie-, latentie- en duurzaamheidsvereisten van elke bestemming.

Patroon Wanneer gebruiken Afweging
Eén zoekopdracht met foreachBatch Bestemmingen gebruiken batchwriters en zouden dezelfde microbatch samen moeten verwerken. De bestemmingen delen een levenscyclus van query’s en hun schrijfbewerkingen zijn niet atomair als geheel.
Onafhankelijke streaming-queries Bestemmingen hebben aparte transformaties, controlepunten of herstartgedrag nodig. Elke query leest en verwerkt de bron onafhankelijk.
Eenmaal opnemen in een bronze Delta-tabel Downstream streams hebben duurzame herhaling, andere schema's of aparte faaldomeinen nodig. De extra Delta write voegt opslag en latentie toe.

Gebruik niet hetzelfde controlepunt voor vertakkende query's. Elke streamingquery heeft een eigen checkpointlocatie nodig.

Schrijf naar meerdere sinks met foreachBatch

foreachBatch geeft elke microbatch DataFrame en zijn batch-ID door aan een functie. Gebruik de functie om batchgewijs DataFrame-schrijvers, MERGE-bewerkingen of schrijfbewerkingen naar meerdere bestemmingen toe te passen.

foreachBatch biedt standaard ten minste één keer verwerking. Het volgende voorbeeld blijft de microbatch behouden, zodat Spark deze niet opnieuw berekent voor elke sink. Delta-transactie-identificaties zorgen ervoor dat elke tabel idempotent schrijft wanneer Spark dezelfde batch opnieuw probeert.

def append_idempotently(batch_df, table_name, transaction_id, batch_id):
    (
        batch_df.write
        .format("delta")
        .mode("append")
        .option("txnAppId", transaction_id)
        .option("txnVersion", batch_id)
        .saveAsTable(table_name)
    )


def write_to_sinks(batch_df, batch_id):
    cached_batch = batch_df.persist()
    try:
        append_idempotently(
            cached_batch,
            "bronze_orders",
            "orders-fanout-bronze",
            batch_id,
        )
        append_idempotently(
            cached_batch.filter("amount >= 1000"),
            "high_value_orders",
            "orders-fanout-high-value",
            batch_id,
        )
    finally:
        cached_batch.unpersist()


query = (
    events.writeStream
    .queryName("orders-fanout")
    .outputMode("append")
    .option("checkpointLocation", "Files/checkpoints/orders-fanout")
    .foreachBatch(write_to_sinks)
    .start()
)

query.awaitTermination()

Houd elke transactie-ID stabiel terwijl je het checkpoint hergebruikt. Als je de query start met een nieuw checkpoint en de batch-ID's opnieuw starten, gebruik dan nieuwe transactie-ID's zodat Delta de nieuwe batches niet als transacties behandelt die het al heeft gecommitted.

De schrijfbewerkingen binnen foreachBatch vormen niet één enkele transactie voor alle bestemmingen. Als een later schrijven mislukt, probeert Spark de batch opnieuw. Maak elke bestemming idempotent zodat eerdere schrijfopdrachten veilig opnieuw kunnen draaien. Voor niet-Delta bestemmingen gebruik een stabiele gebeurtenissleutel, een upsert-operatie of een bestemmingsspecifieke transactie-identificatie.

foreachBatch gebruikt het microbatch-uitvoeringsmodel. Gebruik het niet met Real-time Mode.

Voer onafhankelijke streamingqueries uit

Start een aparte query wanneer elke bestemming een eigen checkpoint, trigger, outputmodus of foutafhandeling nodig heeft. De queries kunnen dezelfde streaming DataFrame-definitie gebruiken, maar Spark voert elke query onafhankelijk uit.

bronze_query = (
    events.writeStream
    .queryName("orders-bronze")
    .format("delta")
    .outputMode("append")
    .option("checkpointLocation", "Files/checkpoints/orders-bronze")
    .toTable("bronze_orders")
)

high_value_query = (
    events.filter("amount >= 1000")
    .writeStream
    .queryName("orders-high-value")
    .format("delta")
    .outputMode("append")
    .option("checkpointLocation", "Files/checkpoints/orders-high-value")
    .toTable("high_value_orders")
)

spark.streams.awaitAnyTermination()

Gebruik onafhankelijke queries voor een klein aantal bestemmingen wanneer bronherlezingen en herhaalde transformaties acceptabel zijn. Controleer en herstart elke query afzonderlijk.

awaitAnyTermination() keert terug wanneer een van beide zoekopdrachten stopt. Inspecteer de querystatus en uitzondering, stop of herstart vervolgens de resterende query in plaats van deze zonder toezicht te laten draaien. Voor meerdere begrensde query's van het type available-now start je eerst alle query's en roep je vervolgens awaitTermination() aan op elke query-handle, zodat de taak op alle query's wacht.

Voor grotere fan-out topologieën neem je de brongegevens één keer op in een bronzen Delta-tabel. Start downstream-streams vanuit die tabel met afzonderlijke checkpoints. Dit patroon geeft elke consumer duurzame herafspeelmogelijkheden en voorkomt dat een langzame bestemming de gegevensinname blokkeert. Standaard Delta-streamingleesbewerkingen verwachten commits met alleen toevoegingen; gebruik de Change Data Feed wanneer downstream streams updates of verwijderingen moeten ontvangen.

Voeg twee stromen met begrensde toestand samen

Een stream-stream join correleert gebeurtenissen die onafhankelijk binnenkomen, zoals een bestelling en de betaling ervan. Beide inputs kunnen late data blijven ontvangen, dus voeg watermerken en een tijdbereikvoorwaarde toe om Spark een punt te geven waarop het de niet-geëvenaarde toestand kan verwijderen.

Het volgende voorbeeld gaat ervan uit dat orders bevat order_id en order_time, en payments bevat order_id, payment_id, en payment_time.

from pyspark.sql import functions as F

watermarked_orders = (
    orders
    .withWatermark("order_time", "10 minutes")
    .alias("orders")
)

watermarked_payments = (
    payments
    .withWatermark("payment_time", "10 minutes")
    .alias("payments")
)

matched_orders = (
    watermarked_orders.join(
        watermarked_payments,
        F.expr("""
            orders.order_id = payments.order_id
            AND payments.payment_time >= orders.order_time
            AND payments.payment_time
                <= orders.order_time + INTERVAL 30 MINUTES
        """),
        "inner",
    )
    .select(
        F.col("orders.order_id").alias("order_id"),
        F.col("orders.order_time"),
        F.col("payments.payment_id"),
        F.col("payments.payment_time"),
    )
)

query = (
    matched_orders.writeStream
    .format("delta")
    .outputMode("append")
    .option("checkpointLocation", "Files/checkpoints/matched-orders")
    .toTable("matched_orders")
)

query.awaitTermination()

Stream-stream joins ondersteunen de append uitvoermodus. Inner joins kunnen zonder watermerken werken, maar de state ervan kan onbegrensd groeien. Outer joins hebben extra watermerkvereisten en vertragen ongeëvenaarde output totdat Spark kan bewijzen dat er geen match kan komen. Voor details, zie Stream-stream joins.

Real-time Mode ondersteunt alleen specifieke stream-stream joinshapes en vereist extra configuratie. Bevestig de Real-time Mode ondersteuningsmatrix voordat je dit patroon verandert naar een Real-time Mode trigger.

Bekijk zoekopdrachten met StreamingQueryListener

A StreamingQueryListener ontvangt levenscyclus- en voortgangsgebeurtenissen voor elke streamingquery in de Spark-sessie. Gebruik het om querystatistieken en details over beëindiging naar je logsysteem of monitoringsysteem te sturen.

import json

from pyspark.sql.streaming import StreamingQueryListener


class QueryMetricsListener(StreamingQueryListener):
    def onQueryStarted(self, event):
        print(json.dumps({
            "event": "started",
            "id": str(event.id),
            "runId": str(event.runId),
            "name": event.name,
        }))

    def onQueryProgress(self, event):
        progress = event.progress
        print(json.dumps({
            "event": "progress",
            "name": progress.name,
            "batchId": progress.batchId,
            "numInputRows": progress.numInputRows,
            "inputRowsPerSecond": progress.inputRowsPerSecond,
            "processedRowsPerSecond": progress.processedRowsPerSecond,
        }, default=str))

    def onQueryIdle(self, event):
        pass

    def onQueryTerminated(self, event):
        print(json.dumps({
            "event": "terminated",
            "id": str(event.id),
            "runId": str(event.runId),
            "exception": event.exception,
        }))


query_listener = QueryMetricsListener()
spark.streams.addListener(query_listener)

Listener-callbacks worden uitgevoerd op de driver en kunnen door verschillende threads worden aangeroepen. Houd onQueryStarted niet-blokkerend, omdat het starten van de query daarop wacht. Houd de andere callbacks kort, bescherm de gedeelde veranderbare toestand en sluit telemetrie in voor asynchrone levering in plaats van langzame netwerkoproepen te maken tijdens een callback.

Voortgangsgebeurtenissen worden gegenereerd wanneer Spark queryvoortgang publiceert. Voor de realtime-modus gebeurt dat bij de overgang tussen langlopende batches, niet voor elk record. Gebruik bronachterstand en monitoring van downstreamlatentie tussen voortgangsgebeurtenissen. Voor details, zie Real-time queries monitoren.

Verwijder in een interactief notitieblok een listener voordat je een vervangende listener registreert, zodat het opnieuw uitvoeren van een cel geen dubbele callbackfuncties veroorzaakt:

spark.streams.removeListener(query_listener)