Not
Bu sayfaya erişim yetkilendirme gerektiriyor. Oturum açmayı veya dizinleri değiştirmeyi deneyebilirsiniz.
Bu sayfaya erişim yetkilendirme gerektiriyor. Dizinleri değiştirmeyi deneyebilirsiniz.
Important
Lakeflow işlem hatlarında gerçek zamanlı mod, önizleme kanalındaki Databricks Runtime 18.1.3'te Genel Önizleme aşamasındadır.
Gerçek zamanlı mod, uçtan uca gecikme süresi beş milisaniyeye kadar düşük olan ultra düşük gecikme süreli veri işlemeye olanak tanır. Sahtekarlık algılama ve gerçek zamanlı kişiselleştirme gibi akış verilerine anında yanıt gerektiren operasyonel iş yükleri için gerçek zamanlı modu kullanın.
Gerçek zamanlı mod, işlem hatlarının dışında doğrudan Yapılandırılmış Akış'ta da kullanılabilir. Bakınız Yapılandırılmış Akış'ta Gerçek Zamanlı Mod.
Gerçek zamanlı mod düşük gecikme süresine nasıl ulaşır?
Gerçek zamanlı mod, standart sürekli işlemeden üç temel yolla farklıdır:
- Uzun süre çalışan toplu işlemler: Sistem, uzun süre çalışan toplu işlemler içinde kaynakta kullanılabilir duruma geldiğinde verileri işler (varsayılan olarak beş dakikadır).
- Eşzamanlı aşama zamanlaması: Tüm sorgu aşamaları aynı anda zamanlanır. İşlem kaynağı, tüm aşamaları eşzamanlı olarak kapsayacak kadar kullanılabilir görev yuvasına sahip olmalıdır. Bkz. İşlem boyutlandırması.
- Akışlı karıştırma: Veriler, aşağı akıştaki aşamanın başlamasından önce yukarı akıştaki bir aşamanın tamamlanmasını beklemek yerine, üretilir üretilmez aşamalar arasında aktarılır.
Denetim noktası aralığı (aracılığıyla pipelines.trigger.intervalyapılandırılır), durum ve kaynak uzaklıklarının dayanıklı depolamada ne sıklıkta kalıcı hale geldiğini denetler. Daha uzun aralıklar denetim noktası ek yükünü azaltır, ancak hatadan sonra kurtarma süresini artırır ve ölçüm raporlamasını geciktirir. Daha kısa aralıklar dayanıklılığı artırır ancak ek yük ekler.
Gerçek zamanlı mod ve sürekli işlem hatları
Gerçek zamanlı mod, özel bir sürekli tetikleyici türüdür. Sürekli mod hala gereklidir; gerçek zamanlı mod, akış düzeyi gecikme süresi iyileştirmelerini en üste ekler. Gerçek zamanlı modu kullanmak için işlem hattının önce sürekli modda çalıştırılması gerekir. Gerçek zamanlı mod daha sonra standart sürekli işlemenin sağladığının ötesinde saniyenin altında gecikme süresi elde etmek için akış düzeyinde ek iyileştirmeler uygular.
Gerçek zamanlı modu etkinleştirmek için üç yapılandırma adımı gerekir:
- İşlem hattını sürekli moda ayarlayın.
- İşlem hattı düzeyinde gerçek zamanlı modu etkinleştirin.
- Gerçek zamanlı bir güncelleştirme akışı tanımlayın.
Requirements
| Requirement | Value |
|---|---|
| Databricks Runtime | Lakeflow işlem hatları önizleme kanalında 18.1.3 |
| İşlem türü | Klasik işlem veya sunucusuz |
Gerçek zamanlı modu yapılandırma
1. Adım: İşlem hattını sürekli moda ayarlama
İşlem hattı ayarlarınızda İşlem hattı modunuSürekli olarak veya işlem hattı JSON'unda ayarlayın:
{
"continuous": true
}
2. Adım: İşlem hattı düzeyinde gerçek zamanlı modu etkinleştirme
İşlem hattı ayarlarınızda , Gelişmiş > Spark yapılandırması altındaki Spark yapılandırmasına aşağıdaki anahtarı ekleyin:
spark.databricks.streaming.realTimeMode.enabled = true
Bunu işlem hattı JSON'unda da ayarlayabilirsiniz:
{
"continuous": true,
"spark_conf": {
"spark.databricks.streaming.realTimeMode.enabled": "true"
}
}
3. Adım: Gerçek zamanlı güncelleştirme akışı tanımlama
Gerçek zamanlı mod bir güncelleştirme akışı gerektirir. Çıkış hedefini tanımlamak için dp.create_sink() kullanın, ardından @dp.update_flow değeri pipelines.trigger olarak ayarlanmış ve "RealTime" sink’i gösterecek şekilde target dekoratörünü kullanın.
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()
)
Akış düzeyi yapılandırma parametreleri:
| Parametre | Zorunlu | Varsayılan | Description |
|---|---|---|---|
pipelines.trigger |
Yes | — |
"RealTime" Bu akış için gerçek zamanlı modu etkinleştirmek için olarak ayarlayın. |
pipelines.trigger.interval |
Hayır | "5 minutes" |
Denetim noktası aralığı. Durum ve ofsetlerin ne sıklıkta kaydedileceğini denetler. Daha kısa değerler kurtarılabilirliği artırır; uzun değerler ek yükü azaltır. |
Kod örnekleri
Kafka'dan Kafka'ya
Bir Kafka topic’inden okuyun ve bir Kafka çıkış hedefine yazın:
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")
)
Yayın katılımı ile zenginleştirme
Kafka akışını statik bir eşleme tablosuyla birleştirme. Yalnızca yayın (akış-statik) birleşimleri desteklenir. Akıştan akışa birleştirmeler gerçek zamanlı modda desteklenmez.
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")
)
Aggregation
Durum bilgisi tutan groupBy kullanarak anahtara göre olay sayımı.
spark.sql.shuffle.partitions öğesini, durum bilgili işlemler için giriş bölüm sayısıyla eşleşecek şekilde ayarlayın:
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()
)
Desteklenen kaynaklar ve havuzlar
| Bağlayıcı | Kaynak olarak | Havuz olarak | Notlar |
|---|---|---|---|
| Apache Kafka | ✓ | ✓ | — |
| AWS MSK | ✓ | ✓ | Kafka uyumlu arabirimini kullanır. |
| Azure Event Hubs (Kafka bağlayıcısı) | ✓ | ✓ | Kafka uyumlu arabirimini kullanır. |
| Amazon Kinesis | ✓ | Desteklenmiyor | Yalnızca EFO (Gelişmiş Fan-Out) modu için kullanın. |
| Delta | Desteklenmiyor | Desteklenmiyor | — |
Hesaplama boyutlandırma
İşlemde yeterli görev yuvası varsa işlem kaynağı başına bir gerçek zamanlı işlem hattı çalıştırabilirsiniz. Kullanılabilir görev yuvaları tüm sorgu aşamalarında tüm görevleri kapsamalıdır.
| İşlem hattı türü | Configuration | Gerekli görev yuvaları |
|---|---|---|
| Tek aşamalı durum bilgisi olmayan (Kafka kaynağı + havuz) |
maxPartitions = 8 |
8 |
| durum bilgisi olan iki aşamalı (Kafka kaynağı + karıştırma) |
maxPartitions = 8, karışık bölümler = 20 |
28 (8 + 20) |
| Üç aşamalı (Kafka kaynağı + iki shuffle) |
maxPartitions = 8, her biri 20 aşamadan oluşan iki karıştırma aşaması |
48 (8 + 20 + 20) |
ayarlamazsanız maxPartitionsKafka konu başlığındaki bölüm sayısını kullanın.
operatör desteği
| Kategori | Operator | Destekleniyor |
|---|---|---|
| Durumsuz | Seçim, Projeksiyon | ✓ |
| UDFs | Scala UDF | ✓ (sınırlamalarla) |
| UDFs | Python Kullanıcı Tanımlı Fonksiyonu (UDF) | ✓ (sınırlamalarla) |
| Aggregation | toplam, sayı, maksimum, minimum, ortalama | ✓ |
| Windowing | Devrilme, Kayma | ✓ |
| Windowing | Oturum | Desteklenmiyor |
| Deduplication | dropDuplicates |
✓ (sınırsız durum) |
| Deduplication | dropDuplicatesWithinWatermark |
Desteklenmiyor |
| Joins | Broadcast tablo birleştirme | ✓ |
| Joins | Akıştan akışa birleştirme | Desteklenmiyor |
| Özelleştirilmiş | transformWithState |
✓ (davranış farklılıkları ile) |
| Özelleştirilmiş | union |
✓ (sınırlamalarla) |
| Özelleştirilmiş | forEach |
Desteklenmiyor |
| Özelleştirilmiş | flatMapGroupsWithState |
Desteklenmiyor |
| Özelleştirilmiş | mapPartitions |
Desteklenmiyor |
| Özelleştirilmiş | forEachBatch |
Desteklenmiyor |
transformWithState gerçek zamanlı modda
transformWithState , mikro toplu işlemden aşağıdaki farklarla gerçek zamanlı modda desteklenir:
-
handleInputRows, her yığındaki her anahtar için bir kez yerine, her satır için bir kez çağrılır. YineleyiciinputRows, çağrı başına tek bir değer verir. - Olay zamanı zamanlayıcıları desteklenmez. Veri gelmemişse, uzun süre çalışan bir toplu iş sonlandığında işleme zamanı zamanlayıcıları tetiklenir.
-
transformWithStateInPandasdesteklenmez.
Pandas UDF'leri gerçek zamanlı modda
pandas UDF’lerle gecikmeyi en aza indirmek için spark.sql.execution.arrow.maxRecordsPerBatch değerini 1 olarak ayarlayın. Bu, aktarım hızı karşılığında gecikme süresini iyileştirir. Aktarım hızı da önemliyse, bu değeri veya daha yüksek bir değere 100 ayarlayın.
Gerçek zamanlı mod performansını izleme
Gerçek zamanlı mod, StreamingQueryProgress içinde latencies alanı altında gecikme metriklerini gösterir. Akış sorgusundaki StreamingQueryListener özelliğini inceleyerek veya lastProgress aracılığıyla bu ölçümlere erişin.
| Metric | Description |
|---|---|
processingLatencyMs |
Kaydın akış tarafından okunması ile akış tarafından tamamen işlenmesi arasında geçen süre |
sourceQueuingLatencyMs |
Bir kaydın ileti veri yoluna başarıyla yazıldığı an (örneğin, Kafka'da günlük ekleme zamanı) ile akış tarafından ilk kez okunduğu an arasındaki süre |
e2eLatencyMs |
Kaydın kaynakta üretilmesinden akış tarafından tamamen işlenmesine kadar olan toplam uçtan uca gecikme süresi |
Her ölçüm p50, p90, p95 ve p99 yüzdebirlik değerleri olarak bildirilir.
Sınırlamalar
İşlem hattı başına bir gerçek zamanlı akış önerilir. Birden çok akışa izin verilir, ancak akışlar arasında görev yuvası çekişmesi gecikme süresini artırır.
İşleç ve kaynak sınırlamalarının tam listesi için bkz. Gerçek zamanlı mod sınırlamaları.