Lakeflow işlem hatlarında gerçek zamanlı modu kullanma

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:

  1. İşlem hattını sürekli moda ayarlayın.
  2. İşlem hattı düzeyinde gerçek zamanlı modu etkinleştirin.
  3. 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. Yineleyici inputRows , ç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.
  • transformWithStateInPandas desteklenmez.

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

Ek kaynaklar