update_flow

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.