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
  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.