Använda realtidsläge i Lakeflow-pipelines

Important

Realtidsläget i Lakeflow-pipelines finns i offentlig förhandsversion på Databricks Runtime 18.1.3 på förhandsgranskningskanalen.

Realtidsläge möjliggör databearbetning med extremt låg latens, med svarstid från slutpunkt till slutpunkt så låg som fem millisekunder. Använd realtidsläge för driftarbetsbelastningar som kräver omedelbara svar på strömmande data, till exempel identifiering av bedrägerier och anpassning i realtid.

Realtidsläge är också tillgängligt direkt i Strukturerad direktuppspelning utanför pipelines. Se Realtidsläge i Strukturerad direktuppspelning.

Så här uppnår realtidsläget låg svarstid

Realtidsläget skiljer sig från kontinuerlig standardbearbetning på tre viktiga sätt:

  • Tidskrävande batchar: Systemet bearbetar data när de blir tillgängliga i källan inom långvariga batchar (standardvärdet är fem minuter).
  • Schemaläggning av samtidiga steg: Alla frågesteg schemaläggs samtidigt. Beräkningsresursen måste ha tillräckligt med tillgängliga aktivitetsfack för att täcka alla faser samtidigt. Se Beräkningsstorlek.
  • Direktuppspelningsblandning: Data skickas mellan faser så snart de skapas, i stället för att vänta på att en uppströmsfas ska slutföras innan nedströmsfasen startas.

Kontrollpunktsintervallet (konfigurerat via pipelines.trigger.interval) styr hur ofta tillstånds- och källförskjutningar sparas till beständig lagring. Längre intervall minskar kontrollpunktskostnaderna men ökar återställningstiden efter ett fel och fördröjer rapporteringen av mått. Kortare intervall förbättrar hållbarheten men lägger till omkostnader.

Realtidsläge och kontinuerliga processkedjor

Realtidsläge är en specialiserad typ av kontinuerlig utlösare. Kontinuerligt läge krävs fortfarande. realtidsläge lägger till optimeringar av svarstid på flödesnivå ovanpå. Om du vill använda realtidsläge måste pipelinen först köras i kontinuerligt läge. Realtidsläget tillämpar sedan ytterligare optimeringar på flödesnivå för att uppnå svarstid under sekund utöver vad kontinuerlig standardbearbetning ger.

För att aktivera realtidsläge krävs tre konfigurationssteg:

  1. Ställ in pipelinen i kontinuerligt läge.
  2. Aktivera realtidsläge på pipelinenivå.
  3. Definiera ett uppdateringsflöde i realtid.

Requirements

Requirement Value
Databricks Runtime 18.1.3 på Lakeflow pipelines-förhandsgranskningskanalen
Typ av beräkning Klassisk beräkning eller serverlös

Konfigurera realtidsläge

Steg 1: Ställ in pipelinen i kontinuerligt läge

I pipelineinställningarna anger du Pipeline-läge till Kontinuerlig eller anger det i pipeline-JSON:

{
  "continuous": true
}

Steg 2: Aktivera realtidsläge på pipelinenivå

I pipelineinställningarna lägger du till följande nyckel i Spark-konfigurationen under Advanced > Spark-konfiguration:

spark.databricks.streaming.realTimeMode.enabled = true

Du kan också ange detta i pipeline-JSON:

{
  "continuous": true,
  "spark_conf": {
    "spark.databricks.streaming.realTimeMode.enabled": "true"
  }
}

Steg 3: Definiera ett uppdateringsflöde i realtid

Realtidsläge kräver ett uppdateringsflöde. Använd dp.create_sink() för att definiera utdatamålet och använd sedan dekoratören @dp.update_flow med pipelines.trigger inställd på "RealTime" och target som pekar på mottagaren.

from pyspark import pipelines as dp

# Define the output sink
dp.create_sink(
    "my_kafka_sink",
    "kafka",
    {
        "kafka.bootstrap.servers": "<bootstrap-servers>",
        "topic": "<output-topic>",
    }
)

# Define the real-time update flow targeting the sink
@dp.update_flow(
    name="my_rtm_flow",
    target="my_kafka_sink",
    spark_conf={
        "pipelines.trigger": "RealTime",
        "pipelines.trigger.interval": "5 minutes",  # optional; defaults to 5 minutes
    }
)
def my_real_time_flow():
    return (
        spark.readStream
            .format("kafka")
            .option("kafka.bootstrap.servers", "<bootstrap-servers>")
            .option("subscribe", "<input-topic>")
            .load()
    )

Konfigurationsparametrar på flödesnivå:

Parameter Obligatoriskt Standardinställning Description
pipelines.trigger Yes Ange till "RealTime" för att aktivera realtidsläge för det här flödet.
pipelines.trigger.interval No "5 minutes" Kontrollpunktsintervall. Styr hur ofta tillstånd och offsetvärden sparas. Kortare värden förbättrar återställningsbarheten. längre värden minskar omkostnaderna.

Kodexempel

Kafka till Kafka

Läs från ett Kafka-ämne och skriv till ett Kafka-utdatamål:

from pyspark import pipelines as dp

dp.create_sink("kafka_output_sink", "kafka", {
    "kafka.bootstrap.servers": broker_address,
    "topic": output_topic,
})

@dp.update_flow(
    name="kafka_rtm_flow",
    target="kafka_output_sink",
    spark_conf={
        "pipelines.trigger": "RealTime",
        "pipelines.trigger.interval": "5 minutes",
    }
)
def kafka_rtm_flow():
    return (
        spark.readStream
            .format("kafka")
            .option("kafka.bootstrap.servers", broker_address)
            .option("subscribe", input_topic)
            .option("startingOffsets", "latest")
            .load()
            .selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)", "timestamp")
    )

Berika med en sändningskoppling

Koppla en Kafka-ström mot en statisk uppslagstabell. Endast sändningskopplingar (ström-till-statiska) stöds. Stream-to-stream-kopplingar stöds inte i realtidsläge.

from pyspark import pipelines as dp
from pyspark.sql.functions import broadcast, expr

dp.create_sink("enriched_output_sink", "kafka", {
    "kafka.bootstrap.servers": broker_address,
    "topic": enriched_output_topic,
})

@dp.update_flow(
    name="enriched_events_flow",
    target="enriched_output_sink",
    spark_conf={
        "pipelines.trigger": "RealTime",
        "pipelines.trigger.interval": "5 minutes",
    }
)
def enriched_events():
    lookup = spark.read.table("catalog.schema.lookup_table")
    return (
        spark.readStream
            .format("kafka")
            .option("kafka.bootstrap.servers", broker_address)
            .option("subscribe", input_topic)
            .load()
            .withColumn("event_key", expr("CAST(value AS STRING)"))
            .join(broadcast(lookup), expr("event_key = lookup_key"))
            .select("event_key", "lookup_value", "timestamp")
    )

Aggregering

Räkna händelser efter nyckel med hjälp av en tillståndskänslig groupBy. Ange spark.sql.shuffle.partitions för att matcha antalet indatapartitioner för tillståndskänsliga åtgärder:

from pyspark import pipelines as dp
from pyspark.sql.functions import col

dp.create_sink("event_counts_sink", "kafka", {
    "kafka.bootstrap.servers": broker_address,
    "topic": output_topic,
})

@dp.update_flow(
    name="event_counts_flow",
    target="event_counts_sink",
    spark_conf={
        "pipelines.trigger": "RealTime",
        "pipelines.trigger.interval": "5 minutes",
        "spark.sql.shuffle.partitions": "8",
    }
)
def event_counts():
    return (
        spark.readStream
            .format("kafka")
            .option("kafka.bootstrap.servers", broker_address)
            .option("subscribe", input_topic)
            .load()
            .selectExpr("CAST(key AS STRING) AS event_type", "timestamp")
            .groupBy(col("event_type"))
            .count()
    )

Källor och mottagare som stöds

Connector Som källa Som sänka Notes
Apache Kafka
AWS MSK Använder det Kafka-kompatibla gränssnittet.
Azure Event Hubs (Kafka-anslutningsprogram) Använder det Kafka-kompatibla gränssnittet.
Amazon Kinesis Stöds ej Använd endast för EFO-läge (utökat Fan-Out).
Delta Stöds ej Stöds ej

Beräkna dimensionering

Du kan köra en realtidspipeline per beräkningsresurs om beräkningsresursen har tillräckligt med uppgiftsplatser. Tillgängliga aktivitetsfack måste omfatta alla aktiviteter i alla frågesteg.

Typ av pipeline Configuration Obligatoriska uppgiftsplatser
Tillståndslös i en fas (Kafka-källa + mottagare) maxPartitions = 8 8
Tillståndskänsligt i två steg (Kafka-källa + shuffle) maxPartitions = 8, shuffle partitioner = 20 28 (8 + 20)
Tre steg (Kafka-källa + två omfördelningar) maxPartitions = 8, två blandningssteg på 20 vardera 48 (8 + 20 + 20)

Om du inte anger maxPartitions, använd antalet partitioner i Kafka-ämnet.

Stöd för operatör

Kategori Operator Understödd
Tillståndslös Urval, projektion
UDFs Scala UDF √ (med begränsningar)
UDFs Python-användardefinierad funktion (UDF) √ (med begränsningar)
Aggregering summa, antal, max, min, medelvärde
Windowing Tumlning, glidning
Windowing Session Stöds ej
Deduplication dropDuplicates √ (obundet tillstånd)
Deduplication dropDuplicatesWithinWatermark Stöds ej
Joins Sändningstabellanslutning
Joins Ström-till-ström-sammankoppling Stöds ej
Skräddarsydd transformWithState √ (med beteendemässiga skillnader)
Skräddarsydd union √ (med begränsningar)
Skräddarsydd forEach Stöds ej
Skräddarsydd flatMapGroupsWithState Stöds ej
Skräddarsydd mapPartitions Stöds ej
Skräddarsydd forEachBatch Stöds ej

transformWithState i realtidsläge

transformWithState stöds i realtidsläge med följande skillnader från mikrobatchbearbetning:

  • handleInputRows anropas en gång per rad i stället för en gång per nyckel per batch. Iteratorn inputRows ger ett enda värde per anrop.
  • Timer för händelsetid stöds inte. Timers för bearbetning utlöses när en tidskrävande batch avslutas om inga data har anlänt.
  • transformWithStateInPandas stöds inte.

Pandas-UDF:er i realtidsläge

Om du vill minimera svarstiden med Pandas UDF:er anger du spark.sql.execution.arrow.maxRecordsPerBatch till 1. Detta optimerar för svarstid på bekostnad av dataflödet. Om dataflödet också är viktigt anger du det här värdet till 100 eller högre.

Övervaka prestanda i realtidsläge

Realtidsläge exponerar svarstidsmått i StreamingQueryProgress under fältet latencies . Få åtkomst till dessa mätvärden via en StreamingQueryListener eller genom att inspektera egenskapen lastProgress på strömningsfrågan.

Mått Description
processingLatencyMs Tid mellan när en post läss av flödet och när den bearbetas helt av flödet
sourceQueuingLatencyMs Tiden mellan att en post har skrivits till meddelandebussen med lyckat resultat (till exempel tidpunkten då posten läggs till i loggen i Kafka) och att den först läses av flödet
e2eLatencyMs Total svarstid från slutpunkt till slutpunkt från när posten skapas vid källan till när den bearbetas fullständigt av flödet

Varje mått rapporteras som p50, p90, p95 och p99 percentiler.

Limitations

Ett realtidsflöde per pipeline rekommenderas. Flera flöden tillåts, men konkurrensen om aktivitetsfack mellan flöden ökar svarstiden.

En fullständig lista över operator- och källbegränsningar finns i Begränsningar i realtidsläge.

Ytterligare resurser