Catatan
Akses ke halaman ini memerlukan otorisasi. Anda dapat mencoba masuk atau mengubah direktori.
Akses ke halaman ini memerlukan otorisasi. Anda dapat mencoba mengubah direktori.
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:
- Tambahkan sumber streaming yang menambahkan data ke tabel streaming yang ada tanpa memerlukan refresh penuh. Misalnya, Anda mungkin memiliki tabel yang menggabungkan data regional dari setiap wilayah tempat Anda beroperasi. Saat wilayah baru diluncurkan, Anda dapat menambahkan data wilayah baru ke tabel tanpa melakukan refresh penuh. Lihat Contoh: Menulis data ke tabel streaming dari beberapa topik Kafka.
- Perbarui tabel streaming dengan menambahkan data historis yang hilang (backfilling). Anda dapat menggunakan
INSERT INTO ONCEsintaksis untuk membuat isi ulang historis yang berjalan satu kali. Lihat Contoh: Menjalankan isi ulang data satu kali dan Mengisi ulang data historis dengan alur. - Gabungkan data dari beberapa sumber dan tulis ke satu tabel streaming alih-alih menggunakan
UNIONklausa dalam kueri. Menggunakan pemrosesan alur tambahan alih-alihUNIONmemungkinkan Anda memperbarui tabel target secara bertahap tanpa menjalankan pembaruan refresh penuh. Lihat Contoh: Gunakan pemrosesan alur tambahan alih-alihUNION.
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_flowdefinisi. - 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
onceparameter keTruepada alur penyambungan. Di SQL, gunakanINSERT 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'))