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.
Flöden i en Lakeflow-pipeline flyttar data till en strömmande tabell eller materialiserad vy. Följande exempel visar hur man definierar standardflöden, definierar ett flöde separat från sitt mål, skriver till en strömmande tabell från flera Kafka-topics, kör en engångsbackfill och ersätter UNION-frågor med append-flödesbearbetning.
En översikt över flöden finns i Läsa in och bearbeta data stegvis med Lakeflow-pipelineflöden.
Exempel: Skapa ett standardflöde
När du skapar en pipeline definierar du vanligtvis en tabell eller vy tillsammans med frågan som stöder den. Den här frågan skapar till exempel en strömmande tabell med namnet customers_silver genom att läsa från customers_bronze. Strömningstabellen och dess standardflöde skapas tillsammans i ett enda steg.
SQL
CREATE OR REFRESH STREAMING TABLE customers_silver
AS SELECT * FROM STREAM(customers_bronze)
python
from pyspark import pipelines as dp
@dp.table()
def customers_silver():
return spark.readStream.table("customers_bronze")
Standardflödet för en strömmande tabell är ett tilläggsflöde som lägger till nya rader med varje uppdatering och har samma namn som målet. Det här är det vanligaste sättet att använda pipelines – att skapa ett flöde och dess mål i ett enda steg – och du kan använda det för att mata in eller transformera data. Mer information om flödesbegrepp finns i Läsa in och bearbeta data stegvis med Lakeflow-pipelineflöden.
Exempel: Definiera ett flöde separat från målet
Du kan också skapa ett flöde för en tabell som du har definierat separat. Resultatet är identiskt med att skapa ett standardflöde, inklusive att använda samma namn för strömningstabellen och flödet:
python
from pyspark import pipelines as dp
# create streaming table
dp.create_streaming_table("customers_silver")
# add a flow
@dp.append_flow(
target = "customers_silver")
def customer_silver():
return spark.readStream.table("customers_bronze")
SQL
-- create a streaming table
CREATE OR REFRESH STREAMING TABLE customers_silver;
-- add a flow
CREATE FLOW customers_silver
AS INSERT INTO customers_silver BY NAME
SELECT * FROM STREAM(customers_bronze);
Genom att definiera ett flöde separat från målet kan du skapa flera flöden som lägger till data till samma mål. Använd dekoratören @dp.append_flow i Python-gränssnittet eller CREATE FLOW...INSERT INTO -satsen i SQL-gränssnittet för att lägga till flöden för uppgifter, till exempel följande:
- Lägg till strömmande källor som lägger till data i en befintlig strömmande tabell utan att kräva en fullständig uppdatering. Du kan till exempel ha en tabell som kombinerar regionala data från varje region som du arbetar i. När nya regioner distribueras kan du lägga till nya regiondata i tabellen utan att utföra en fullständig uppdatering. Se Exempel: Skriv till en streamingtabell från flera Kafka-ämnen.
- Uppdatera en strömmande tabell genom att lägga till saknade historiska data (återfyllnad). Du kan använda syntaxen
INSERT INTO ONCEför att skapa en historisk återfyllnad som körs en gång. Se Exempel: Kör en engångsåterfyllning av data och Återfyllning av historiska data med pipelines. - Kombinera data från flera källor och skriv till en enda strömmande tabell i stället för att använda
UNION-satsen i en fråga. Om du använder tilläggsflödesbearbetning i stället förUNIONkan du uppdatera måltabellen inkrementellt utan att köra en fullständig uppdatering. Se Exempel: Använd append-flödesbearbetning i stället förUNION.
För Python-frågor använder du funktionen create_streaming_table() för att skapa en måltabell.
Important
- Om du behöver definiera datakvalitetsbegränsningar med förväntningar definierar du förväntningarna på måltabellen som en del av
create_streaming_table()funktionen eller i en befintlig tabelldefinition. Du kan inte definiera förväntningar i@append_flowdefinitionen. - Flöden identifieras med ett flödesnamn och det här namnet används för att identifiera kontrollpunkter för strömning. Användning av flödesnamnet för att identifiera kontrollpunkten innebär följande:
- Om ett befintligt flöde i en pipeline byter namn överförs inte kontrollpunkten och det omdöpta flödet är i praktiken ett helt nytt flöde.
- Du kan inte återanvända ett flödesnamn i en pipeline eftersom den befintliga kontrollpunkten inte matchar den nya flödesdefinitionen.
Exempel: Skriva till en strömmande tabell från flera Kafka-ämnen
I följande exempel skapas en strömmande tabell med namnet kafka_target och skriver till den strömmande tabellen från två Kafka-ämnen.
python
from pyspark import pipelines as dp
dp.create_streaming_table("kafka_target")
# Kafka stream from multiple topics
@dp.append_flow(target = "kafka_target")
def topic1():
return (
spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "host1:port1,...")
.option("subscribe", "topic1")
.load()
)
@dp.append_flow(target = "kafka_target")
def topic2():
return (
spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "host1:port1,...")
.option("subscribe", "topic2")
.load()
)
SQL
CREATE OR REFRESH STREAMING TABLE kafka_target;
CREATE FLOW
topic1
AS INSERT INTO
kafka_target BY NAME
SELECT * FROM
read_kafka(bootstrapServers => 'host1:port1,...', subscribe => 'topic1');
CREATE FLOW
topic2
AS INSERT INTO
kafka_target BY NAME
SELECT * FROM
read_kafka(bootstrapServers => 'host1:port1,...', subscribe => 'topic2');
Mer information om den read_kafka() tabellvärdesfunktion som används i SQL-frågorna finns i read_kafka i SQL-språkreferensen.
I Python kan du programmatiskt skapa flera flöden som riktar sig mot en enda tabell. I följande exempel visas det här mönstret för en lista över Kafka-ämnen.
Anmärkning
Det här mönstret har samma krav som att använda en for loop för att skapa tabeller. Du måste uttryckligen skicka ett Python-värde till funktionen som definierar flödet. Se Skapa tabeller i en for loop.
from pyspark import pipelines as dp
dp.create_streaming_table("kafka_target")
topic_list = ["topic1", "topic2", "topic3"]
for topic_name in topic_list:
@dp.append_flow(target = "kafka_target", name=f"{topic_name}_flow")
def topic_flow(topic=topic_name):
return (
spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "host1:port1,...")
.option("subscribe", topic)
.load()
)
Exempel: Kör en engångsdatapåfyllning
Om du vill köra en fråga för att lägga till data i en befintlig strömmande tabell använder du append_flow.
När du har lagt till en uppsättning befintliga data har du flera alternativ:
- Om du vill att frågan ska lägga till ny data om den kommer till katalogen för återfyllnad, förblir frågan aktiv.
- Om du vill att detta ska vara en engångsefterfyllnad och aldrig köras igen tar du bort frågan när du har kört pipelinen en gång.
- Om du vill att frågan ska köras en gång och bara köras igen i de fall där data uppdateras fullständigt anger du parametern
oncetillTruei tilläggsflödet. I SQL använder duINSERT INTO ONCE.
I följande exempel körs en fråga för att lägga till historiska data i en strömmande tabell:
python
from pyspark import pipelines as dp
@dp.table()
def csv_target():
return spark.readStream
.format("cloudFiles")
.option("cloudFiles.format","csv")
.load("path/to/sourceDir")
@dp.append_flow(
target = "csv_target",
once = True)
def backfill():
return spark.read
.format("cloudFiles")
.option("cloudFiles.format","csv")
.load("path/to/backfill/data/dir")
SQL
CREATE OR REFRESH STREAMING TABLE csv_target
AS SELECT * FROM
read_files(
"path/to/sourceDir",
"csv"
);
CREATE FLOW
backfill
AS INSERT INTO ONCE
csv_target BY NAME
SELECT * FROM
read_files(
"path/to/backfill/data/dir",
"csv"
);
För ett mer detaljerat exempel, se Att fylla i historiska data med hjälp av pipelines.
Exempel: Använd tilläggsflödesbearbetning i stället för UNION
I stället för att använda en fråga med en UNION sats kan du använda tilläggsflödesfrågor för att kombinera flera källor och skriva till en enda strömmande tabell. Använda tilläggsflödesfrågor i stället för UNION tillåter att lägga till poster i en strömmande tabell från flera källor utan behov av att köra en fullständig uppdatering.
Följande Python-exempel innehåller en fråga som kombinerar flera datakällor med en UNION sats:
@dp.create_table(name="raw_orders")
def unioned_raw_orders():
raw_orders_us = (
spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "csv")
.load("/path/to/orders/us")
)
raw_orders_eu = (
spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "csv")
.load("/path/to/orders/eu")
)
return raw_orders_us.union(raw_orders_eu)
Följande exempel ersätter UNION frågan med tilläggsflödesfrågor:
python
dp.create_streaming_table("raw_orders")
@dp.append_flow(target="raw_orders")
def raw_orders_us():
return spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "csv")
.load("/path/to/orders/us")
@dp.append_flow(target="raw_orders")
def raw_orders_eu():
return spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "csv")
.load("/path/to/orders/eu")
# Additional flows can be added without the full refresh that a UNION query would require:
@dp.append_flow(target="raw_orders")
def raw_orders_apac():
return spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "csv")
.load("/path/to/orders/apac")
SQL
CREATE OR REFRESH STREAMING TABLE raw_orders;
CREATE FLOW
raw_orders_us
AS INSERT INTO
raw_orders BY NAME
SELECT * FROM
STREAM read_files(
"/path/to/orders/us",
format => "csv"
);
CREATE FLOW
raw_orders_eu
AS INSERT INTO
raw_orders BY NAME
SELECT * FROM
STREAM read_files(
"/path/to/orders/eu",
format => "csv"
);
-- Additional flows can be added without the full refresh that a UNION query would require:
CREATE FLOW
raw_orders_apac
AS INSERT INTO
raw_orders BY NAME
SELECT * FROM
STREAM read_files(
"/path/to/orders/apac",
format => "csv"
);
Exempel: Använd transformWithState för att övervaka sensor pulsslag
I följande exempel visas en tillståndskänslig processor som läser från Kafka och verifierar att sensorer sänder pulsslag med jämna mellanrum. Om ett hjärtslag inte tas emot inom 5 minuter ger processorn ifrån sig en post till måldeltatabellen för analys.
För mer information om att bygga anpassade tillståndsfulla applikationer, se Bygg en anpassad tillståndsbaserad applikation med transformWithState.
Anmärkning
RocksDB är den nya standardleverantören av tillstånd från och med Databricks Runtime 17.2. Om frågan misslyckas på grund av ett undantag från providern som inte stöds lägger du till följande pipelinekonfigurationer, utför en fullständig uppdatering eller kontrollpunktsåterställning och kör sedan pipelinen igen:
"configuration": {
"spark.sql.streaming.stateStore.providerClass": "com.databricks.sql.streaming.state.RocksDBStateStoreProvider",
"spark.sql.streaming.stateStore.rocksdb.changelogCheckpointing.enabled": "true"
}
from typing import Iterator
import pandas as pd
from pyspark import pipelines as dp
from pyspark.sql.functions import col, from_json
from pyspark.sql.streaming import StatefulProcessor, StatefulProcessorHandle
from pyspark.sql.types import StructType, StructField, LongType, StringType, TimestampType
KAFKA_TOPIC = "<your-kafka-topic>"
output_schema = StructType([
StructField("sensor_id", LongType(), False),
StructField("sensor_type", StringType(), False),
StructField("last_heartbeat_time", TimestampType(), False)])
class SensorHeartbeatProcessor(StatefulProcessor):
def init(self, handle: StatefulProcessorHandle) -> None:
# Define state schema to store sensor information (sensor_id is the grouping key)
state_schema = StructType([
StructField("sensor_type", StringType(), False),
StructField("last_heartbeat_time", TimestampType(), False)])
self.sensor_state = handle.getValueState("sensorState", state_schema)
# State variable to track the previously registered timer
timer_schema = StructType([StructField("timer_ts", LongType(), False)])
self.timer_state = handle.getValueState("timerState", timer_schema)
self.handle = handle
def handleInputRows(self, key, rows, timerValues) -> Iterator[pd.DataFrame]:
# Process one row from input and update state
pdf = next(rows)
row = pdf.iloc[0]
# Store or update the sensor information in state using current timestamp
current_time = pd.Timestamp(timerValues.getCurrentProcessingTimeInMs(), unit='ms')
self.sensor_state.update((
row["sensor_type"],
current_time
))
# Delete old timer if already registered
if self.timer_state.exists():
old_timer = self.timer_state.get()[0]
self.handle.deleteTimer(old_timer)
# Register a timer for 5 minutes from current processing time
expiry_time = timerValues.getCurrentProcessingTimeInMs() + (5 * 60 * 1000)
self.handle.registerTimer(expiry_time)
# Store the new timer timestamp in state
self.timer_state.update((expiry_time,))
# No output on input processing, output only on timer expiry
return iter([])
def handleExpiredTimer(self, key, timerValues, expiredTimerInfo) -> Iterator[pd.DataFrame]:
# Emit output row based on state store
if self.sensor_state.exists():
state = self.sensor_state.get()
output = pd.DataFrame({
"sensor_id": [key[0]], # Use grouping key as sensor_id
"sensor_type": [state[0]],
"last_heartbeat_time": [state[1]]
})
# Remove the entry for the sensor from the state store
self.sensor_state.clear()
# Remove the timer state entry
self.timer_state.clear()
yield output
def close(self) -> None:
pass
dp.create_streaming_table("sensorAlerts")
# Define the schema for the Kafka message value
sensor_schema = StructType([
StructField("sensor_id", LongType(), False),
StructField("sensor_type", StringType(), False),
StructField("sensor_value", LongType(), False)])
@dp.append_flow(target = "sensorAlerts")
def kafka_delta_flow():
return (
spark.readStream
.format("kafka")
.option("subscribe", KAFKA_TOPIC)
.option("startingOffsets", "earliest")
.load()
.select(from_json(col("value").cast("string"), sensor_schema).alias("data"), col("timestamp"))
.select("data.*", "timestamp")
.withWatermark('timestamp', '1 hour')
.groupBy(col("sensor_id"))
.transformWithStateInPandas(
statefulProcessor = SensorHeartbeatProcessor(),
outputStructType = output_schema,
outputMode = 'update',
timeMode = 'ProcessingTime'))