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.
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:
- Atur alur ke mode berkelanjutan.
- Aktifkan mode real-time pada tingkat pipeline.
- 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:
-
handleInputRowsdipanggil sekali untuk setiap baris, bukan sekali untuk setiap kunci dalam setiap kelompok. IteratorinputRowsmenghasilkan satu nilai per pemanggilan. - Timer event-time tidak didukung. Timer pemrosesan diaktifkan ketika batch yang berjalan lama berakhir jika tidak ada data yang tiba.
-
transformWithStateInPandastidak 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.