Yazılara geri dön
Yazı

Spark Yapılandırılmış Akışta Watermark ve allowedLateness Kullanarak Veri Kaybını Nasıl Önleriz?

Spark Yapılandırılmış Akış'ta watermark ve allowedLateness kullanarak geç gelen olayları nasıl yöneteceğinizi ve olay-zaman işlemede veri kaybını nasıl önleyeceğinizi öğrenin.

Big DataSpark Structured StreamingwatermarkallowedLatenessevent-time processingdata loss prevention

Spark Structured Streaming ile gerçek zamanlı analizler için çalışmaya başladığımda, işim çalıştığı sürece tüm olayların doğru şekilde işleneceğini varsaydım. Ancak, olay-zamanı işleminin sessiz bir risk getirdiğini hızla öğrendim: gecikmiş veri. Akışınız, ağ sorunları, toplu veri alma gecikmesi veya çevrimdışı mobil cihazların senkronizasyonu gibi nedenlerle gecikmeleri hesaba katmazsa, herhangi bir hata göstermeden değerli olayları kaybedebilirsiniz. Bu yazıda, hata ayıkladığım gerçek olaylara dayanarak, üretim hattımda veri kaybını önlemek için nasıl watermark ve allowedLateness kullandığımı paylaşacağım.

Olay-Zamanı İşlemeyi Anlamak ve Geciken Veri Sorunu

Olay-zamanı işlemede, Spark, veriyi ne zaman geldiği (işleme zamanı) temel alarak işlemez; yerine, olayın kendisine gömülü zaman damgasına göre işler. Bu, kullanıcı oturumlandırma, dolandırıcılık tespiti veya IoT telemetrisi gibi kullanım senaryolarında çok önemlidir — çünkü olayın gerçekleşme zamanı, onu gördüğümüz zamandan daha önemli olur.

Ancak gerçek dünya sistemleri mükemmel değildir. Olaylar dakika, saat veya hatta günler gecikebilir. Akış işiniz zaten olay-zamanı filigranını (watermark) bir gecikmiş olayın zaman damgası ötesine taşıdıysa, Spark onu sessizce atar — uyarı vermez, metrikte ani artış olmaz, sadece aşağı akımdaki toplamalarda veri eksik olur.

Bir perakende analitiği hattında bu durumu yaşadım: bağlantısı zor olan bölgelerdeki mobil uygulama kullanıcılarından gelen olaylar, en fazla 4 saat gecikiyordu. İlk filigran ayarımız 10 dakika olduğu için, akşam sattığımız olayların %15’i kayboluyordu. Gelir raporlamasının doğruluğuna bağlı olduğu bir sistemde bu kabul edilemez bir durumdu.

Spark Yapılandırılmış Akışta Watermark Nasıl Çalışır

Watermark, Spark'ın gecikmiş veriler için ne kadar bekleyeceğimizi olay zamanına göre takip ettiği mekanizmadır. Aşağıdaki gibi tanımlanır:

watermark = görülen_maximum_olay_zamanı - gecikme_sınırı

Burada delay_threshold (gecikme_sınırı), beklediğiniz maksimum gecikmeyi ifade eder. Güncel watermark'tan daha eski bir zaman damgasına sahip olan her olay, çok geç kabul edilir ve atılır.

Bu değeri sorgunuzda şu şekilde ayarlarsınız:

val stream = spark
  .readStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "broker1:9092,broker2:9092")
  .option("subscribe", "user-events")
  .load()
  .selectExpr("CAST(value AS STRING)")
  .select(from_json(col("value"), schema).as("data"))
  .select("data.userId", "data.eventTime", "data.action")
  .withWatermark("eventTime", "2 hours")

val aggregated = stream
  .groupBy(window(col("eventTime", "1 hour"), "15 minutes"), col("userId"))
  .agg(count("*").as("eventCount"))

aggregated.writeStream
  .outputMode("update")
  .format("console")
  .start()

Bu örnekte, Spark'a şöyle diyoruz: "Gördüğünüz en son olay zamanına göre, 2 saat gerideki olaylar için durum bilgisi tutulsun." Eğer bir olayın eventTime değeri 10:00 ise ve gördüğünüz en son olay zamanı 14:00 ise, watermark 12:00 olur — bu yüzden 10:00'luk olay hâlâ işlenir. Ancak bu olay 14:00 sonrasında gelirse, atılır.

allowedLateness Rolü: Geç Verilere İkinci Şans

Watermark tek başına katıdır: bir kez geçildiğinde, veri kaybolur. Ancak bazen, gecikmenin nadir ama tahmin edilebilir gecikmeler nedeniyle watermark eşiğini aşabileceğini bilirsiniz — örneğin, ortak sistemlerden gelen 3'te çalışan günlük toplu yüklemeler.

Burada allowedLateness devreye girer. window() veya dropDuplicates() gibi durumlu işlemler için watermark'ın ötesindeki pencereyi genişletir. Watermark'tan sonra gelen ancak izin verilen gecikme süresi içinde olan olaylar, hâlâ durumu güncelleyebilir.

Ödeme対账流中的使用方式如下:

val payments = spark.readStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "payments-broker:9092")
  .option("subscribe", "payment-events")
  .load()
  .selectExpr("CAST(value AS STRING)")
  .select(from_json(col("value"), paymentSchema).as("data"))
  .select("data.transactionId", "data.eventTime", "data.amount", "data.status")
  .withWatermark("eventTime", "30 minutes")

val dailySummary = payments
  .groupBy(window(col("eventTime", "1 day"), "1 hour"), col("status"))
  .agg(sum("amount").as("totalAmount"))
  .allowLateness("6 hours")  // Güncellemeleri 6 saat kadar gecikmeye izin ver

val query = dailySummary
  .writeStream
  .outputMode("update")
  .format("delta")
  .option("checkpointLocation", "/mnt/delta/checkpoints/payments_daily")
  .trigger(Trigger.ProcessingTime("10 minutes"))
  .start()

Bu durumda, dün bir ödeme eventi bugün saat 9'da (15 saat geç) gelirse, pencerenin bitiş zamanına göre 6 saatlik izin verilen gecikme penceresinin içinde olduğu sürece, günlük toplamı hâlâ güncelleyebilir. allowLateness("6 hours") olmadan, bu olaylar göz ardı edilir.

Üretimde İzlediğim En İyi Uygulamalar

Zaman içinde, veri kaybını önlerken durum boyutunu yönetilebilir tutmak için bazı kurallar geliştirdim:

  • Gözlemlenen gecikmeyle başla: Staging ortamında gerçek olay gecikmesini ölç. eventTime ile ingestionTime (Kafka tarafından alındığı zaman) arasındaki farkı kullanarak gerçekçi bir watermark ayarla.
  • Tahminde aşırı olma: Maksimum gecikmenin 20 dakika olduğu bir sistemde watermark'i "6 saat" olarak ayarlamak, durum boyutunu ve arızalardan sonraki kurtarma süresini gereksiz yere artırır.
  • Watermark'i allowedLateness ile birleştir: Watermark'i durum büyümesini sınırlamak için kullan, allowedLateness'i ise bilinen outlier'ları işlemek için kullan.
  • Dropped event'leri izle: Watermark'in altına düşen olayları loglamak için foreachBatch havuzu ekle:
stream.writeStream
  .foreachBatch { (batchDF, batchId) =>
    val lateEvents = batchDF.filter(col("eventTime").lt(currentWatermark()))
    if (lateEvents.count() > 0) {
      logWarning(s"Dropped ${lateEvents.count()} late events in batch $batchId")
      lateEvents.write.mode("append").json("/mnt/late-events/")
    }
  }
  .start()
  • Amaca yönelik gecikmeyle test et: Test ortamında, geçmiş zaman damgalı olaylar enjekte etmek için bir betik kullan ve bu olayların doğru şekilde işlendiğini doğrula.

Pitfalls I’ve Encountered

Başlangıçta yaptığım bir hata, allowedLateness parametresinin watermark'i etkilediğini varsaymaktı. Oysa öyle değil — sadece durum güncellemeleri için pencereyi uzatır. Watermark hâlâ max_event_time_seen - delay_threshold formülüne göre ilerler. Eğer allowedLateness değerini watermark gecikmesinden daha yüksek ayarlarsanız, bekleme süresinizi artırmıyorsunuz — sadece mevcut duruma gecikmiş güncellemeler yapmanıza izin veriyorsunuz.

Başka bir dikkat edilmesi gereken nokta: outputMode("complete") kullanıyorsanız, allowedLateness hiçbir etkisi yoktur çünkü tüm durum her seferinde yeniden hesaplanır ve yeniden yazılır. Gecikmeyi etkili bir şekilde kullanmak için update veya append modunu tercih edin.

Son olarak, durum temizliğini dikkatli yapın. Spark, watermark'li anahtarlar için durumu otomatik olarak temizler; ancak özel durum haritaları veya flatMapGroupsWithState kullanıyorsanız, durumun silinmesini watermark temelinde kendiniz yönetmelisiniz.

Sonuç

WaterMark ve allowedLateness sadece yapılandırma düğmeleri değildir — güvenilir olay-zamanı ardışık düzenleri oluşturmak için temel araçlardır. Gerçek gözlemlenen gecikmelere göre ayarlamalar yapıp sınır durumlarını test ederek, yapılandırılmış akış işlerimde iki haneli yüzde veri kaybından %0,1'in altında veri kaybına kadar ilerledim.

Spark Structured Streaming'i üretim ortamında çalıştırıyorsanız ve bu parametreleri henüz ayarlamadıysanız, gerçek gecikmenizi ölçerek başlayın. 15 dakikalık bir waterMark belki de ihtiyacınız olan her şeydir — ya da benim gibi keşfedebilirsiniz: mobil kullanıcılarınızın bağlı kalabilmesi için 2 saatlik bir pencereye ihtiyaç duyabilir. Anahtar, varsayımlarla değil, gerçeklikle uyumlu ayarlar yapmaktır.

Önceki yazımda [akış gecikmesini izleme](https://furkanikkan.com/urun/perfetto-ile-mikro-saniye-cozunurlukte-cekirdek-fonksiyon-cagrilarinin-profili-olusturulmasi-83) konusunda bahsettiğim gibi, gözlemlenebilirlik savaşın yarısıdır. Diğer yarısı ise, gerçek dünya verisinin karışık doğasına karşı dayanıklı bir akış tasarlamaktır.


Kapak görseli: USDAgov · PDM (Openverse / kamu malı) · https://www.flickr.com/photos/41284017@N08/54674581419