Folyamok használata a Lakeflow-adatcsatornákban

A Lakeflow-folyamat adatfolyamai adatokat továbbítanak egy adatfolyam-táblába vagy materializált nézetbe. Az alábbi példák bemutatják, hogyan definiálhat alapértelmezett folyamatokat, hogyan definiálhat egy folyamatot a céljától külön, hogyan írhat streamingtáblába több Kafka-témából, hogyan futtathat egyszeri backfill műveletet, és hogyan helyettesítheti a UNION lekérdezéseket append flow-feldolgozással.

A folyamatok áttekintéséhez tekintse meg az adatok növekményes betöltését és feldolgozását a Lakeflow-folyamatokkal.

Példa: Alapértelmezett folyamat létrehozása

Folyamat létrehozásakor általában egy táblát vagy nézetet határoz meg az azt támogató lekérdezéssel együtt. Például ez a lekérdezés egy customers_silver nevű streamingtáblát hoz létre a(z) customers_bronze beolvasásával. A streamelési tábla és az alapértelmezett folyamat egyetlen lépésben jön létre.

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")

A streamelőtáblák alapértelmezett folyamata egy hozzáfűző folyamat, amely minden frissítéshez új sorokat ad hozzá, és neve megegyezik a cél nevével. Ez a folyamatok használatának leggyakoribb módja – egy folyamat és a cél egyetlen lépésben történő létrehozása –, és az adatok betöltésére vagy átalakítására is használható. A folyamatfogalmakkal kapcsolatos további információkért lásd az adatok növekményes betöltését és feldolgozását a Lakeflow-folyamatokkal.

Példa: Folyamat definiálása a céltól elkülönítve

Létrehozhat folyamatot egy külön definiált táblához is. Az eredmény azonos azzal, mintha egy alapértelmezett folyamatot hozna létre, beleértve azt is, hogy a streaming tábla és a folyamat neve megegyezik:

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);

Egy folyamat céltól külön definiálása lehetővé teszi több olyan folyamat létrehozását, amelyek adatokat fűznek ugyanahhoz a célhoz. A Python-felületen a @dp.append_flow dekorátort, az SQL-felületen pedig a CREATE FLOW...INSERT INTO záradékot használva folyamatokat adhat hozzá például az alábbi feladatokhoz:

  • Olyan streamelési források hozzáadása, amelyek adatokat fűznek egy meglévő streamelési táblához teljes frissítés nélkül. Lehet például, hogy van egy táblázata, amely a regionális adatokat egyesíti minden régióból, amelyben dolgozik. Az új régiók bevezetésekor teljes frissítés nélkül hozzáadhatja az új régióadatokat a táblához. Lásd: Példa: streamingtáblába írás több Kafka-témakörből.
  • Frissítsen egy streamelési táblát hiányzó előzményadatok hozzáfűzésével (backfilling). A(z) INSERT INTO ONCE szintaxis használatával létrehozhat egy egyszer lefutó előzményfeltöltést. Lásd: Példa: Egyszeri adat-visszatöltés végrehajtása és Korábbi adatok visszatöltése folyamatokkal.
  • Több forrásból származó adatokat egyesítsen és írjon egyetlen streamelési táblába a UNION záradék használata helyett egy lekérdezésben. A hozzácsatolási folyamat feldolgozását használva lehetőséged van a céltáblázat fokozatos frissítésére anélkül, hogy teljes körű frissítést kellene futtatnod. Lásd: Példa: a hozzáfűzési folyamat használata UNION helyett.

Python-lekérdezések esetén a create_streaming_table() függvénnyel hozzon létre egy céltáblát.

Important

  • Ha elvárásokkal kell meghatároznia az adatminőségre vonatkozó korlátozásokat, a céltáblán a függvény részeként vagy egy meglévő tábladefinícióban határozza meg az create_streaming_table() elvárásokat. A definícióban @append_flow nem definiálhat elvárásokat.
  • A folyamatokat egy folyamatnév azonosítja, és ez a név a streamelési ellenőrzőpontok azonosítására szolgál. A folyamatnév használata az ellenőrzőpont azonosításához a következőket jelenti:
    • Ha egy meglévő folyamatot egy csővezetékben átneveznek, az ellenőrzési pont nem viszonyul tovább, és az átnevezett folyamat gyakorlatilag teljesen új folyamatnak számít.
    • A folyamat neve nem használható újra, mert a meglévő ellenőrzőpont nem felel meg az új folyamatdefiníciónak.

Példa: Írás streamelési táblába több Kafka-témakörből

Az alábbi példák létrehoznak egy kafka_target nevű streamelési táblát, és ebből a táblából írnak két Kafka-téma adatait.

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');

Az SQL-lekérdezésekben használt táblaértékű függvényről további információt a "read_kafka" fejezetben talál az SQL nyelvi referenciában.

A Pythonban programozott módon hozhat létre több folyamatot, amelyek egyetlen táblát céloznak meg. Az alábbi példa ezt a mintát mutatja be a Kafka-témakörök listájához.

Megjegyzés:

Ez a minta ugyanolyan követelményeket támaszt, mint amikor for hurok segítségével táblázatokat hozunk létre. Explicit módon át kell adnia egy Python-értéket a folyamatot meghatározó függvénynek. Lásd: Táblák létrehozása egy for ciklusban.

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()
    )

Példa: Egyszeri adatvisszatöltés futtatása

Ha egy lekérdezést szeretne futtatni, amely adatokat szeretne hozzáfűzni egy meglévő streamelési táblához, használja a következőt append_flow: .

A meglévő adatok egy készletének hozzáfűzése után több lehetősége is van:

  • Ha azt szeretné, hogy a lekérdezés új adatokat fűzzön hozzá, ha az a backfill könyvtárba érkezik, hagyja a lekérdezést helyben.
  • Ha azt szeretné, hogy ez egyszeri visszatöltés legyen, és soha ne fusson újra, távolítsa el a lekérdezést a folyamat egyszeri futtatása után.
  • Ha azt szeretné, hogy a lekérdezés egyszer fusson, és csak akkor futtassa újra, ha az adatok teljes frissítése folyamatban van, állítsa a once paramétert a hozzáfűzési folyamatra True . Az SQL-ben használja a következőt INSERT INTO ONCE: .

Az alábbi példák egy lekérdezést futtatnak az előzményadatok streamelési táblához való hozzáfűzéséhez:

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"
  );

Részletesebb példa: Az előzményadatok visszatöltése futtatószálakkal.

Példa: Használja a hozzáfűzési folyamat feldolgozását, ahelyett, hogy UNION

Ahelyett, hogy záradékkal rendelkező UNION lekérdezést használ, a hozzáfűző folyamatlekérdezésekkel több forrást egyesíthet, és egyetlen streamelési táblába írhat. Hozzáfűző lekérdezések használatával különböző forrásokból fűzhet hozzá egy stream táblához anélkül, hogy UNION futtatására lenne szükség.

A következő Python-példa egy olyan lekérdezést tartalmaz, amely több adatforrást kombinál egy UNION záradékkal:

@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)

A következő példák a UNION lekérdezést hozzáfűzési folyamathoz tartozó lekérdezésekre cserélik.

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"
  );

Példa: Érzékelő transformWithState szívveréseinek figyelése

Az alábbi példa egy állapotalapú processzort mutat be, amely a Kafkából olvas be, és ellenőrzi, hogy az érzékelők rendszeresen bocsátanak-e ki szívveréseket. Ha 5 percen belül nem érkezik szívverés, a processzor bejegyzést ír a cél Delta-táblába elemzés céljából.

További információért az egyedi állapotú alkalmazások építéséről lásd: Build egy custom stateful application with transformWithState.

Megjegyzés:

A RocksDB az alapértelmezett állapotszolgáltató a Databricks Runtime 17.2-vel kezdődően. Ha a lekérdezés nem támogatott szolgáltatói kivétel miatt meghiúsul, adja hozzá a következő folyamatkonfigurációkat, végezzen teljes frissítést vagy ellenőrzőpont-visszaállítást, majd futtassa újra a folyamatot:

"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'))