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.
Gunakan dekorator @dp.update_flow untuk membuat alur pembaruan. Alur pembaruan menulis ke sink menggunakan mode output pembaruan, hanya memancarkan baris yang berubah di setiap batch. Tidak seperti alur penambahan, mereka mendukung agregasi stateful tanpa memerlukan marka air.
Alur pembaruan hanya dapat menargetkan sink. Tabel delta tidak didukung.
Syntax
from pyspark import pipelines as dp
dp.create_sink("<sink-name>", "<format>", {"<key>": "<value>"})
@dp.update_flow(
target = "<sink-name>",
name = "<flow-name>", # optional, defaults to function name
spark_conf = {"<key>" : "<value>", "<key>" : "<value>"}, # optional
comment = "<comment>", # optional
import_checkpoint = "<checkpoint-path>") # optional
def <function-name>():
return (<streaming-query>)
Parameters
| Parameter | Type | Description |
|---|---|---|
| fungsi | function |
Required. Fungsi yang mengembalikan DataFrame streaming Apache Spark dari kueri yang ditentukan pengguna. |
target |
str |
Required. Nama sink tempat alur ini menulis. |
name |
str |
Nama alur. Jika tidak disediakan, akan otomatis menggunakan nama fungsi. |
comment |
str |
Deskripsi untuk alur. |
spark_conf |
dict |
Dict konfigurasi Spark untuk eksekusi kueri ini. Konfigurasi ini mengambil alih confs set untuk tujuan, alur, atau kluster. |
import_checkpoint |
str |
Jalur titik pemeriksaan eksternal untuk diimpor sebelum memulai alur. Diimpor hanya sekali, ketika direktori titik pemeriksaan alur belum ada. |
Examples
Agregasi ke sink Kafka
Tulis hasil agregasi stateful ke sink Kafka:
from pyspark import pipelines as dp
from pyspark.sql.functions import col
dp.create_sink("event_counts_sink", "kafka", {
"kafka.bootstrap.servers": broker_address,
"topic": output_topic,
})
@dp.update_flow(
name="event_counts_flow",
target="event_counts_sink",
)
def event_counts():
return (
spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", broker_address)
.option("subscribe", input_topic)
.load()
.selectExpr("CAST(key AS STRING) AS event_type")
.groupBy(col("event_type"))
.count()
)
Mode Real-time
Important
Mode real time ada di Pratinjau Umum.
Gunakan spark_conf untuk mengonfigurasi alur pembaruan untuk mode real-time:
from pyspark import pipelines as dp
dp.create_sink("my_kafka_sink", "kafka", {
"kafka.bootstrap.servers": broker_address,
"topic": output_topic,
})
@dp.update_flow(
name="my_rtm_flow",
target="my_kafka_sink",
spark_conf={
"pipelines.trigger": "RealTime",
"pipelines.trigger.interval": "5 minutes",
}
)
def my_real_time_flow():
return (
spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", broker_address)
.option("subscribe", input_topic)
.load()
)
Keterbatasan
- Sink tabel Delta tidak didukung sebagai target untuk alur pembaruan.