depends_on ile ardışık düzen akışı yürütmesini sıralama

Important

Bu özellik Genel Önizleme aşamasındadır.

depends_on öğesini kullanmak için işlem hattınızı Lakeflow pipelines PREVIEW kanalını kullanacak şekilde yapılandırın. channel bölümünü İşlem Hattı Yapılandırmaları içinde kontrol edin.

Varsayılan olarak, bir boru hattı akışları veri bağımlılıklarına göre planlar: Bir akış başka bir akışın yazdığı tabloyu okursa, okuyucu yazarın ardından çalışır. Akış sıralaması, bir akışın, açıkça depends_on ile bağımlılık bildirerek, veri okumadığı başka bir akışı beklemesine olanak tanır.

depends_on="other_flow" bu akışın ancak başarılı bir şekilde tamamlandıktan other_flow sonra başladığı anlamına gelir. Bu bir zamanlama kenarıdır: bir akışın ne zaman başladığını kontrol eder, nasıl çalıştığını değil. Bildirmek depends_on , bir akışın tetikleyicisini, modunu veya tek seferlik bir akış olup olmadığını değiştirmez. Bunları akışın kendisinde yine de ilan ediyorsunuz, örneğin once=True.

Akış sıralamasını ne zaman kullanmalı?

Ana kullanım durumu, canlı kaynağa geçmeden önce geçmiş verileri boşaltarak akış durumunu korumaktır; örneğin bir tabloyu toplu doldurmadan canlı bir Apache Kafka kaynağına taşımak gibi.

Her iki kaynağın aynı anda okunması sağlıklı çalışmaz: sınırlı geriye dönük doldurma watermark’ı geride tutar; bu da durumun bellekten atılmasını ve pencereli sonuçları geciktirir. Önce arabelleği boşaltmak ve ardından canlı yayını başlatmak bunu önler. Akış sıralandırması bu ikisini sıralar.

Akış sıralaması diğer kalıpları da destekler:

  • Aynı tabloya yapılan birden fazla geriye dönük veri doldurma işleminin katı biçimde sıralanması.
  • İlk yükleme, ardından fark kapatma ve sonra canlı akış gibi aşamalı iş hatları.
  • Farklı tablolara yazan ancak belirli bir sırayla çalışması gereken dizileme akışları.

depends_on ile sipariş akışları

depends_on, Python pipelines API’sindeki @dp.update_flow ve @dp.append_flow dekoratörlerinde mevcuttur. Tek bir akış adını veya bir akış isimleri listesini kabul eder. Bir listeyle, her adlandırılmış akış, bağımlı akışın başlamasından önce tamamlanmalıdır.

Aşağıdaki örnek, tek seferlik bir geri doldurmayı events akış tablosuna aktarır ve ardından, yalnızca geri doldurma tamamlandıktan sonra aynı tabloya canlı bir Kafka akışı başlatır:

from pyspark import pipelines as dp

dp.create_streaming_table(name="events")

# Drain the historical backfill first.
@dp.append_flow(target="events", once=True, name="events_backfill")
def events_backfill():
    return spark.read.table("historical_events")

# Start the live stream only after the backfill completes.
@dp.append_flow(target="events", name="events_live", depends_on="events_backfill")
def events_live():
    return (
        spark.readStream.format("kafka")
        .option("kafka.bootstrap.servers", "<server>:<port>")
        .option("subscribe", "events")
        .load()
    )

Birden fazla öncülü beklemek için bir liste iletin. Aşağıdaki örnek paralel olarak iki geri doldurma işlemi yapar ve canlı yayını ancak her ikisi tamamlandıktan sonra başlatır:

@dp.append_flow(target="events", once=True, name="backfill_2024")
def backfill_2024():
    return spark.read.table("events_2024")

@dp.append_flow(target="events", once=True, name="backfill_2025")
def backfill_2025():
    return spark.read.table("events_2025")

@dp.append_flow(
    target="events",
    name="events_live",
    depends_on=["backfill_2024", "backfill_2025"],
)
def events_live():
    return spark.readStream.format("kafka").option("subscribe", "events").load()

Geri doldurma işlemlerini paralel olarak değil, birbiri ardına çalıştırmak için bunlar arasında depends_on zincirleyin:

@dp.append_flow(target="events", once=True, name="backfill_2024")
def backfill_2024():
    return spark.read.table("events_2024")

@dp.append_flow(
    target="events", once=True, name="backfill_2025", depends_on="backfill_2024"
)
def backfill_2025():
    return spark.read.table("events_2025")

Gereksinimler ve davranış

Akış sıralaması için aşağıdaki kurallar geçerlidir:

  • Akış sırası sadece bir boru hattı içinde çalışır. Boru hattı zamanlayıcısı olmadan sipariş yerine getirilemez.
  • Bir selefin, bir görünüme değil, masaya ya da lavaboya yazmalıdır. Bir görünüme yazan bir akış hiçbir zaman sonlanma durumuna ulaşmaz, bu nedenle ondan sonra gelen bir akış hiçbir zaman başlamaz. Akış masası hedefleri ve foreachBatch sinkler her ikisi de geçerli öncüllerdir.
  • Farklı hedefler arasında sipariş vermeye izin verilir. Bir akış, farklı bir tabloya yazan bir akışa bağlı olabilir.
  • Bilinmeyen akış isimleri ve döngüleri, boru hattı çalışmadan önce doğrulama sırasında yakalanır.

Tetiklenen ve sürekli boru hatlarında öncül ne olabilir?

Öncül olarak işlev görebilecek akış türleri, boru hattının yürütme moduna bağlıdır:

  • Tetiklenen boru hatları: herhangi bir akış öncül olabilir. Tetiklenen bir çalıştırmadaki her akış bir son duruma ulaşır, bu nedenle sıralama her çalıştırma için geçerlidir.
  • Sürekli boru hatları: Öncül, terminal duruma ulaşan tek seferlik (once) bir akış olmalıdır. Sürekli çalışan bir akış asla sonlanmaz, bu yüzden sonrasında sipariş edilen akış asla başlamaz ve boru hattı doğrulamada bunu reddeder.

Bir foreachBatch akış her zaman bir akış yutucusu olduğu ve tek seferlik bir akış olamayacağı için, yalnızca tetiklenen bir boru hattında öncül olarak hareket edebilir. Sürekli bir boru hattında, bir foreachBatch akışı kendisinden önce gelen bir once akışını bekleyebilir, ancak kendisi bir öncül olamaz.

İşletimsel davranış

Bir once öncülün tamamlanma durumu kalıcıdır; bu nedenle yeniden başlatmalardan ve işlem hattı güncellemelerinden sonra da korunur:

  • Yeniden başlatma: Zaten boşaltılmış akışlar boşaltılmış durumda kalır ve atlanır. Boru hattı henüz tamamlanmamış ilk akışta yeniden başlar. Yürütme canlı akışa ulaştıktan sonra, sonraki yeniden başlatmalar sadece canlı akışa devam eder.
  • Tam yenileme: zincirin tamamlanma durumunu temizler ve sırayla baştan tekrar çalıştırır.
  • Tek bir akış için denetim noktası sıfırlama: canlı akışı sıfırlamak, yalnızca canlı akışın yeniden oynatılmasına neden olur. Yukarı akış once dolguları boşaltılı kalıyor ve tekrar çalıştırılmıyor. Canlı bir sorguyu geçmiş verileri yeniden boşaltmadan geri yüklemenin alışılmış yolu budur.

Ek kaynaklar