Kommentar
Åtkomst till den här sidan kräver auktorisering. Du kan prova att logga in eller ändra kataloger.
Åtkomst till den här sidan kräver auktorisering. Du kan prova att ändra kataloger.
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:
- Ställ in pipelinen i kontinuerligt läge.
- Aktivera realtidsläge på pipelinenivå.
- 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:
-
handleInputRowsanropas en gång per rad i stället för en gång per nyckel per batch. IteratorninputRowsger 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.
-
transformWithStateInPandasstö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.