Gunakan mode real-time di pipeline Lakeflow

Important

Mode waktu nyata dalam pipeline Lakeflow tersedia dalam Pratinjau Publik di Databricks Runtime 18.1.3 pada saluran pratinjau.

Mode real time memungkinkan pemrosesan data latensi ultra rendah, dengan latensi end-to-end serendah lima milidetik. Gunakan mode real time untuk beban kerja operasional yang memerlukan respons langsung terhadap data streaming, seperti deteksi penipuan dan personalisasi real-time.

Mode real time juga tersedia langsung di Streaming Terstruktur di luar alur. Lihat Mode real time di Streaming Terstruktur.

Bagaimana mode real time mencapai latensi rendah

Mode real time berbeda dari pemrosesan berkelanjutan standar dengan tiga cara utama:

  • Batch yang berjalan lama: Sistem memproses data saat tersedia di sumber dalam batch yang berjalan lama (defaultnya adalah lima menit).
  • Penjadwalan tahap simultan: Semua tahap kueri dijadwalkan secara bersamaan. Sumber daya komputasi harus memiliki slot tugas yang cukup tersedia untuk mencakup semua tahapan secara bersamaan. Lihat Ukuran komputasi.
  • Pengacakan streaming: data diteruskan antartahap segera setelah dihasilkan, alih-alih menunggu tahap hulu selesai sebelum memulai tahap hilir.

Interval titik pemeriksaan (dikonfigurasi melalui pipelines.trigger.interval) mengontrol seberapa sering status dan offset sumber dipertahankan ke penyimpanan yang tahan lama. Interval yang lebih panjang mengurangi overhead checkpointing, tetapi meningkatkan waktu pemulihan setelah kegagalan dan menunda pelaporan metrik. Interval yang lebih pendek meningkatkan durabilitas tetapi menambahkan overhead.

Mode real time dan alur berkelanjutan

Mode real time adalah jenis pemicu berkelanjutan khusus. Mode berkelanjutan masih diperlukan; Mode real time menambahkan pengoptimalan latensi tingkat aliran di atasnya. Untuk menggunakan mode real time, alur harus terlebih dahulu berjalan dalam mode berkelanjutan. Mode real time kemudian menerapkan pengoptimalan tambahan pada tingkat alur untuk mencapai latensi sub-detik di luar apa yang disediakan pemrosesan berkelanjutan standar.

Mengaktifkan mode real time memerlukan tiga langkah konfigurasi:

  1. Atur alur ke mode berkelanjutan.
  2. Aktifkan mode real-time pada tingkat pipeline.
  3. Tentukan alur pembaruan real-time.

Requirements

Requirement Value
Databricks Runtime 18.1.3 pada saluran pratinjau alur Lakeflow
Jenis komputasi Komputasi klasik atau tanpa server

Konfigurasikan mode real-time

Langkah 1: Atur alur ke mode berkelanjutan

Di pengaturan alur Anda, atur mode Alur ke Berkelanjutan, atau atur di alur JSON:

{
  "continuous": true
}

Langkah 2: Aktifkan mode real-time di tingkat alur

Di pengaturan alur Anda, tambahkan kunci berikut ke konfigurasi Spark di bawah Konfigurasi Spark Tingkat Lanjut>:

spark.databricks.streaming.realTimeMode.enabled = true

Anda juga dapat mengatur ini di alur JSON:

{
  "continuous": true,
  "spark_conf": {
    "spark.databricks.streaming.realTimeMode.enabled": "true"
  }
}

Langkah 3: Tentukan alur pembaruan real-time

Mode real time memerlukan alur pembaruan. Gunakan dp.create_sink() untuk menentukan target output, lalu gunakan dekorator @dp.update_flow dengan pipelines.trigger diatur ke "RealTime" dan target mengarah ke sink.

from pyspark import pipelines as dp

# Define the output sink
dp.create_sink(
    "my_kafka_sink",
    "kafka",
    {
        "kafka.bootstrap.servers": "<bootstrap-servers>",
        "topic": "<output-topic>",
    }
)

# Define the real-time update flow targeting the sink
@dp.update_flow(
    name="my_rtm_flow",
    target="my_kafka_sink",
    spark_conf={
        "pipelines.trigger": "RealTime",
        "pipelines.trigger.interval": "5 minutes",  # optional; defaults to 5 minutes
    }
)
def my_real_time_flow():
    return (
        spark.readStream
            .format("kafka")
            .option("kafka.bootstrap.servers", "<bootstrap-servers>")
            .option("subscribe", "<input-topic>")
            .load()
    )

Parameter konfigurasi tingkat alur:

Parameter Required Default Description
pipelines.trigger Yes Atur ke "RealTime" untuk mengaktifkan mode real time untuk alur ini.
pipelines.trigger.interval No "5 minutes" Interval pos pemeriksaan. Mengontrol seberapa sering status dan offset disimpan. Nilai yang lebih pendek meningkatkan pemulihan; nilai yang lebih panjang mengurangi overhead.

Contoh kode

Kafka ke Kafka

Baca dari topik Kafka dan tulis ke target output Kafka:

from pyspark import pipelines as dp

dp.create_sink("kafka_output_sink", "kafka", {
    "kafka.bootstrap.servers": broker_address,
    "topic": output_topic,
})

@dp.update_flow(
    name="kafka_rtm_flow",
    target="kafka_output_sink",
    spark_conf={
        "pipelines.trigger": "RealTime",
        "pipelines.trigger.interval": "5 minutes",
    }
)
def kafka_rtm_flow():
    return (
        spark.readStream
            .format("kafka")
            .option("kafka.bootstrap.servers", broker_address)
            .option("subscribe", input_topic)
            .option("startingOffsets", "latest")
            .load()
            .selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)", "timestamp")
    )

Memperkaya dengan gabungan siaran

Gabungkan aliran Kafka terhadap tabel pencarian statis. Hanya gabungan siaran (stream-to-static) yang didukung. Gabungan stream-to-stream tidak didukung dalam mode real time.

from pyspark import pipelines as dp
from pyspark.sql.functions import broadcast, expr

dp.create_sink("enriched_output_sink", "kafka", {
    "kafka.bootstrap.servers": broker_address,
    "topic": enriched_output_topic,
})

@dp.update_flow(
    name="enriched_events_flow",
    target="enriched_output_sink",
    spark_conf={
        "pipelines.trigger": "RealTime",
        "pipelines.trigger.interval": "5 minutes",
    }
)
def enriched_events():
    lookup = spark.read.table("catalog.schema.lookup_table")
    return (
        spark.readStream
            .format("kafka")
            .option("kafka.bootstrap.servers", broker_address)
            .option("subscribe", input_topic)
            .load()
            .withColumn("event_key", expr("CAST(value AS STRING)"))
            .join(broadcast(lookup), expr("event_key = lookup_key"))
            .select("event_key", "lookup_value", "timestamp")
    )

Agregasi

Hitung peristiwa menurut kunci menggunakan stateful groupBy. Atur spark.sql.shuffle.partitions agar sesuai dengan jumlah partisi input untuk operasi stateful:

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",
    spark_conf={
        "pipelines.trigger": "RealTime",
        "pipelines.trigger.interval": "5 minutes",
        "spark.sql.shuffle.partitions": "8",
    }
)
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", "timestamp")
            .groupBy(col("event_type"))
            .count()
    )

Sumber dan penampung yang didukung

Connector Sebagai sumber Sebagai sink Notes
Apache Kafka
AWS MSK Menggunakan antarmuka yang kompatibel dengan Kafka.
Azure Event Hubs (konektor Kafka) Menggunakan antarmuka yang kompatibel dengan Kafka.
Amazon Kinesis Tidak didukung Gunakan untuk mode EFO (Enhanced Fan-Out) saja.
Delta Tidak didukung Tidak didukung

Penentuan Ukuran Komputasi

Anda dapat menjalankan satu alur real-time per sumber daya komputasi jika komputasi memiliki slot tugas yang cukup. Slot tugas yang tersedia harus mencakup semua tugas di semua tahap kueri.

Jenis alur Configuration Slot tugas yang diperlukan
Stateless tahap tunggal (sumber Kafka + penampungan) maxPartitions = 8 8
Stateful dua tahap (sumber Kafka + acak) maxPartitions = 8, partisi shuffle = 20 28 (8 + 20)
Tiga tahap (sumber Kafka + dua shuffle) maxPartitions = 8, dua tahap pengacakan masing-masing 20 48 (8 + 20 + 20)

Jika Anda tidak mengatur maxPartitions, gunakan jumlah partisi dalam topik Kafka.

Dukungan operator

Kategori Operator Dukungan
Tanpa Negara Seleksi, Proyeksi
UDFs Scala UDF ✓ (dengan batasan)
UDFs Python UDF (User Defined Function) ✓ (dengan batasan)
Agregasi jumlah, hitungan, maks, min, rata-rata
Windowing Tumbling, Geser
Windowing Session Tidak didukung
Deduplication dropDuplicates ✓ (status tidak terbatas)
Deduplication dropDuplicatesWithinWatermark Tidak didukung
Joins Gabungan tabel siaran
Joins Penggabungan stream ke stream Tidak didukung
Kustom transformWithState ✓ (dengan perbedaan perilaku)
Kustom union ✓ (dengan batasan)
Kustom forEach Tidak didukung
Kustom flatMapGroupsWithState Tidak didukung
Kustom mapPartitions Tidak didukung
Kustom forEachBatch Tidak didukung

transformWithState dalam mode waktu nyata

transformWithState didukung pada mode real-time dengan perbedaan berikut dibandingkan dengan pemrosesan mikro-batch:

  • handleInputRows dipanggil sekali untuk setiap baris, bukan sekali untuk setiap kunci dalam setiap kelompok. Iterator inputRows menghasilkan satu nilai per pemanggilan.
  • Timer event-time tidak didukung. Timer pemrosesan diaktifkan ketika batch yang berjalan lama berakhir jika tidak ada data yang tiba.
  • transformWithStateInPandas tidak didukung.

Pandas UDF dalam mode waktu nyata

Untuk meminimalkan latensi saat menggunakan pandas UDFs, atur spark.sql.execution.arrow.maxRecordsPerBatch ke 1. Ini mengoptimalkan latensi dengan mengorbankan throughput. Jika throughput juga penting, atur nilai ini ke 100 atau lebih tinggi.

Memantau performa mode real-time

Mode real-time menampilkan metrik latensi di bawah kolom StreamingQueryProgress pada latencies. Akses metrik ini melalui StreamingQueryListener atau dengan memeriksa lastProgress properti pada kueri streaming.

Metrik Description
processingLatencyMs Waktu antara ketika rekaman dibaca oleh alur dan ketika sepenuhnya diproses oleh alur
sourceQueuingLatencyMs Waktu antara saat catatan berhasil ditulis ke bus pesan (misalnya, waktu penambahan log di Kafka) dan saat pertama kali dibaca oleh alur
e2eLatencyMs Total latensi ujung ke ujung dari saat rekaman data diproduksi di sumber hingga diproses sepenuhnya oleh alur

Setiap metrik dilaporkan sebagai persentil p50, p90, p95, dan p99.

Keterbatasan

Disarankan satu aliran waktu nyata untuk setiap pipeline. Beberapa alur diizinkan, tetapi pertikaian slot tugas di seluruh alur meningkatkan latensi.

Untuk daftar lengkap batasan operator dan sumber, lihat Batasan mode real time.

Sumber daya tambahan