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
depends_on = "<flow-name>", # optional, Public Preview
spark_conf = {"<key>" : "<value>", "<key>" : "<value>"}, # optional
comment = "<comment>") # 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. |
depends_on |
str atau list |
Pratinjau Publik Satu atau lebih nama alur yang harus berhasil diselesaikan sebelum alur ini dimulai. Menerima satu nama alur atau daftar nama. Ini hanya mengatur eksekusi alur; ini tidak mengubah cara alur berjalan. Lihat Eksekusi aliran pipeline pesanan dengan depends_on. |
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. |
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.