Använd flöden i Lakeflow-pipelines

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:

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_flow definitionen.
  • 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 once till True i tilläggsflödet. I SQL använder du INSERT 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'))