Menggunakan alur di alur Lakeflow

Aliran dalam pipeline Lakeflow memindahkan data ke tabel streaming atau tampilan terwujud. Contoh berikut menunjukkan cara mendefinisikan alur default, mendefinisikan alur secara terpisah dari targetnya, menulis ke tabel streaming dari beberapa topik Kafka, menjalankan backfill satu kali, dan mengganti kueri UNION dengan pemrosesan append flow.

Untuk gambaran umum alur, lihat Memuat dan memproses data secara bertahap dengan alur Lakeflow.

Contoh: Membuat alur default

Saat membuat alur, Anda biasanya menentukan tabel atau tampilan bersama dengan kueri yang mendukungnya. Misalnya, kueri ini membuat tabel streaming bernama customers_silver dengan membaca dari customers_bronze. Tabel streaming dan alur defaultnya dibuat bersama-sama dalam satu langkah.

SQL

CREATE OR REFRESH STREAMING TABLE customers_silver
AS SELECT * FROM STREAM(customers_bronze)

Phyton

from pyspark import pipelines as dp

@dp.table()
def customers_silver():
  return spark.readStream.table("customers_bronze")

Alur default untuk tabel streaming adalah alur penambahan yang menambahkan baris baru dengan setiap pembaruan, dan memiliki nama yang sama dengan target. Ini adalah cara paling umum untuk menggunakan alur—membuat alur dan targetnya dalam satu langkah—dan Anda dapat menggunakannya untuk menyerap atau mengubah data. Untuk informasi selengkapnya tentang konsep alur, lihat Memuat dan memproses data secara bertahap dengan alur Lakeflow.

Contoh: Tentukan alur secara terpisah dari targetnya

Anda juga dapat membuat alur untuk tabel yang Anda tentukan secara terpisah. Hasilnya identik dengan membuat alur default, termasuk menggunakan nama yang sama untuk tabel streaming dan alur:

Phyton

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

Menentukan alur secara terpisah dari targetnya memungkinkan Anda membuat beberapa alur yang menambahkan data ke target yang sama. @dp.append_flow Gunakan dekorator di antarmuka Python atau CREATE FLOW...INSERT INTO klausa di antarmuka SQL untuk menambahkan alur tugas seperti berikut:

Untuk kueri Python, gunakan fungsi create_streaming_table() untuk membuat tabel target.

Important

  • Jika Anda perlu menentukan batasan kualitas data dengan harapan, tentukan harapan pada tabel target sebagai bagian create_streaming_table() dari fungsi atau pada definisi tabel yang ada. Anda tidak dapat menentukan ekspektasi dalam @append_flow definisi.
  • Alur diidentifikasi dengan nama alur, dan nama ini digunakan untuk mengidentifikasi titik pemeriksaan streaming. Penggunaan nama alur untuk mengidentifikasi titik pemeriksaan berarti sebagai berikut:
    • Jika sebuah alur yang ada dalam pipeline diganti namanya, titik pemeriksaan tidak dilanjutkan, dan alur yang diganti namanya tersebut secara efektif menjadi alur baru yang sepenuhnya berbeda.
    • Anda tidak dapat menggunakan kembali nama aliran dalam pipa, karena pos pemeriksaan yang ada tidak akan cocok dengan definisi aliran baru.

Contoh: Menulis ke tabel streaming dari beberapa topik Kafka

Contoh berikut membuat tabel streaming bernama kafka_target dan menulis ke tabel streaming tersebut dari dua topik Kafka:

Phyton

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

Untuk mempelajari selengkapnya tentang fungsi bernilai tabel yang read_kafka() digunakan dalam kueri SQL, lihat read_kafka dalam referensi bahasa SQL.

Di Python, Anda dapat secara terprogram membuat beberapa alur yang menargetkan satu tabel. Contoh berikut menunjukkan pola ini untuk daftar topik Kafka.

Nota

Pola ini memiliki persyaratan yang sama seperti menggunakan perulangan for untuk membuat tabel. Anda harus secara eksplisit meneruskan nilai Python ke fungsi yang menentukan alur. Lihat Buat tabel dalam loopfor.

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

Contoh: Menjalankan isi ulang data satu kali

Jika Anda ingin menjalankan kueri untuk menambahkan data ke tabel streaming yang sudah ada, gunakan append_flow.

Setelah menambahkan sekumpulan data yang sudah ada, Anda memiliki beberapa opsi:

  • Jika Anda ingin kueri menambahkan data baru jika tiba di direktori isi ulang, biarkan kueri di tempatnya.
  • Jika Anda ingin ini menjadi isi ulang satu kali, dan jangan pernah berjalan lagi, hapus kueri setelah menjalankan alur sekali.
  • Jika Anda ingin kueri berjalan satu kali, dan hanya berjalan lagi dalam kasus di mana data mengalami penyegaran penuh, atur once parameter ke True pada alur penyambungan. Di SQL, gunakan INSERT INTO ONCE.

Contoh berikut menjalankan kueri untuk menambahkan data historis ke tabel streaming:

Phyton

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

Untuk contoh yang lebih mendalam, lihat Mengisi ulang data historis dengan alur.

Contoh: Gunakan pemrosesan alur tambahan alih-alih UNION

Daripada menggunakan kueri dengan UNION klausul, Anda dapat menggunakan kueri alur penambahan untuk menggabungkan beberapa sumber dan memasukkannya ke dalam satu tabel streaming. Menggunakan kueri alur tambahan alih-alih UNION memungkinkan Anda menambahkan ke tabel streaming dari beberapa sumber tanpa menjalankan refresh penuh.

Contoh Python berikut menyertakan kueri yang menggabungkan beberapa sumber data dengan UNION klausa:

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

Contoh berikut mengganti kueri pencarian UNION dengan kueri pengaliran tambahan.

Phyton

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

Contoh: Gunakan transformWithState untuk memantau detak jantung sensor

Contoh berikut menunjukkan prosesor stateful yang berbunyi dari Kafka dan memverifikasi bahwa sensor memancarkan heartbeat secara berkala. Jika heartbeat tidak diterima dalam waktu 5 menit, prosesor memancarkan entri ke tabel Delta target untuk analisis.

Untuk informasi lebih lanjut tentang membangun aplikasi stateful kustom, lihat Bangun aplikasi stateful kustom dengan transformWithState.

Nota

RocksDB adalah penyedia status default yang dimulai dengan Databricks Runtime 17.2. Jika kueri gagal karena pengecualian penyedia yang tidak didukung, tambahkan konfigurasi alur berikut, lakukan reset refresh atau titik pemeriksaan penuh, lalu jalankan ulang alur Anda:

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