Yapılandırılmış Akış için üretimle ilgili dikkat edilmesi gerekenler

Azure Databricks üzerinde üretim ortamındaki Structured Streaming iş yüklerini zamanlanmış Lakeflow Jobs olarak çalıştırın. Bakınız Lakeflow İşleri.

Databricks her zaman aşağıdakileri yapılandırmanızı önerir:

  • display ve countgibi sonuçları döndürebilecek gereksiz kodu not defterlerinden kaldırın.
  • Structured Streaming iş yüklerini genel amaçlı işlem kaynakları üzerinde çalıştırmayın. Akışları her zaman iş işlemeyi kullanarak Lakeflow Jobs olarak planlayın.
  • Lakeflow İşlerini Continuousmodunu kullanarak zamanlayın. Azure Databricks İşleri zamanlama özelliğini, Yapılandırılmış Akış trigger interval'ından ayıran budur.
  • Yapılandırılmış Akış işleri için işlem için otomatik ölçeklendirmeyi etkinleştirmeyin.

Bazı iş yükleri aşağıdakilerden yararlanıyor:

Databricks, Yapılandırılmış Akış iş yükleri için üretim altyapısını yönetme karmaşıklıklarını azaltmak için Lakeflow işlem hatlarını kullanıma sunar. Databricks, yeni Yapılandırılmış Akış işlem hatları için Lakeflow işlem hatlarının kullanılmasını önerir. Bkz. Spark Bildirimsel İşlem Hatları.

Not

İşlem otomatik ölçeklendirmesi, Yapılandırılmış Akış iş yükleri için küme boyutunu azaltmayla ilgili sınırlamalara sahiptir. Databricks, akış iş yükleri için geliştirilmiş otomatik ölçeklendirme ile Lakeflow üzerinde Spark Bildirimli İşlem Hatlarının kullanılmasını önerir. Bkz. Otomatik ölçeklendirme ile Lakeflow işlem hattı kümesi kullanımını iyileştirme.

:::note Sunucusuz işlem

Sunucusuz işlemde yalnızca Trigger.AvailableNow() ve Trigger.Once() desteklenir. Databricks Trigger.AvailableNow() öneriyor.

Sunucusuz işlemde sürekli akış için, sürekli modda Tetiklenen ile sürekli işlem hattı modu karşılaştırmasını kullanın.

Bkz . Akış sınırlamaları.

:::

Operasyonel akış için gecikmeyi azaltın

Operasyonel akış iş yükleri, verileri neredeyse gerçek zamanlı olarak alır, dönüştürür ve üzerinde işlem yapar. Yaygın örnekler arasında dolandırıcılık tespiti, anomali tespiti, kişiselleştirme ve gecikmeli işlemlerin doğrudan iş sonuçlarını etkilediği gerçek zamanlı izleme ve uyarı bulunur. Bu iş yükleri için düşük gecikme genellikle onlarca ila yüzlerce milisaniye anlamına gelir, ancak birçok ekip yüksek yüzde ile arasındaki değişkenliği hesaba katmak için hizmet seviyesi anlaşmaları (SLA) saniyeler aralığında belirler.

En düşük uçtan uca gecikme için, uçtan uca gecikmeyi bir saniyenin altında, yaygın durumlarda ise yaklaşık 300 milisaniye elde eden gerçek zamanlı mod kullanın. Gerçek zamanlı mod kavramlarına bakınız.

Gerçek zamanlı mod iş yükünüze uymadığında, aşağıdaki en iyi uygulamalar mikro-toplu Yapılandırılmış Akış için gecikmeyi azaltır:

  • Çıkış modu: Sorgu operatörlerinizin ve sink'inizin desteklediği güncelleme modunu kullanın. Güncelleme modu, her tetiklemeden sonra güncellenmiş satırları çıktı olarak verir ve watermark süresi sona erene kadar bunları güncellemeyi sürdürür; bu nedenle, aşağı akış hedefinizin güncellenmiş sonuçları işleyebilmesi için idempotent olmasını sağlayın. Güncelleme modunun desteklemediği iş yükleri için ekleme modunu kullanın, örneğin akış-akış birleşimleri veya geç gelen verileri bırakabildiğiniz zamanlar. Düşük gecikme için tam mod kullanmayın. Bkz. Yapılandırılmış Akışiçin çıkış modu seçme.
  • Tetikleyici: Önceki mikro toplu iş biter bitmez ve yeni veriler kullanılabilir olduğunda bir sonraki mikro toplu işi başlatan, 0 aralığına sahip bir processingTime tetikleyici kullanın. Bu, en düşük mikro parti gecikmesini sağlar, ancak bulut depolama API maliyetlerini artırır. Once, Continuous veya AvailableNow öğelerini operasyonel iş yükleri için kullanmayın. Bkz . Yapılandırılmış Akış tetikleyici aralıklarını yapılandırma.
  • Watermark: Watermark değerini, iş yükünüzün atlamaması gereken gecikmeli gelen verileri kapsayacak kadar yüksek ayarlayın. Watermark, sorgunun sırasız olay zamanı verilerini, bunları düşürüp durumu bellekten atmadan önce ne kadar süreyle kabul edeceğini belirler; bu nedenle çok kısa bir watermark, geçerli ancak gecikmiş kayıtları sessizce eler. Bu kısıtlama içinde, daha kısa bir filigran gecikmeyi azaltır ve daha az durum korur; daha uzun bir filigran ise gecikme ve durum pahasına daha fazla geç veriyi tolere eder. Gecikme SLA’nizin 2 katı gibi küçük bir kat değeri, ince ayar yapmak için makul bir başlangıç noktasıdır. Bkz Veri işleme eşiklerini denetlemek için filigranları uygulama.
  • Kaynaklar ve yutucular: Mesaj veri yolları (Apache Kafka, Amazon Kinesis, Apache Pulsar veya Google Cloud Pub/Sub) gibi düşük gecikmeli kaynaklardan okuyun veya Delta Lake ve Apache Iceberg tablolarından veri akışlarını değiştirin. Düşük gecikmeli, yüksek işlem hacimli hedeflere; örneğin ileti veri yollarına, işlemsel veritabanlarına veya foreach hedeflerine yazın. Alt akıştaki tüketicilerin yinelenen ve geç gelen verileri işleyebilmesi için alıcı işlemlerini idempotent olacak şekilde tasarlayın.
  • Durum ve kontrol noktası: Durumlu sorgular için, hem changelog kontrol noktası hem de asenkron durum kontrol noktası için gerekli olan RocksDB state store'u kullanın. Yalnızca artımlı durum değişikliklerini kalıcı hale getirmek için değişiklik günlüğü denetim noktası oluşturmayı etkinleştirin. Durum denetim noktası oluşturma, toplu iş sürenizde darboğaza neden oluyorsa, hata kurtarma ve küme yeniden boyutlandırmayla ilgili sınırlamalarını gözden geçirdikten sonra, denetim noktası yazımlarını bir sonraki mikro toplu işlemle örtüştürmek için eşzamansız durum denetim noktası oluşturmayı etkinleştirin. Her sorguya dayanıklı bulut depolamada kendi kontrol noktası dizinini verin. Bkz. Azure Databricks'te RocksDB durum deposunu yapılandırma, durum bilgisi olan sorgular için eşzamansız durum denetim noktası oluşturma ve Structured Streaming denetim noktaları.
  • Ofset yönetimi: Sürekli akışlarda ofset kontrol noktalarından kaynaklanan gecikmeyi azaltmak için, veri işleme engellenmeden ofset ve commit loglarını güncelleyen asenkron ilerleme takibi etkinleştirin. Once veya AvailableNow tetikleyicileriyle uyumlu değildir. Bkz. Zaman uyumsuz ilerleme izleme.
  • Depolama atlamaları: Mümkün olduğunca hesaplamayı tek bir akış boru hattı içinde tutun. Mantığı birden fazla iş veya iş hattına bölmek, gecikmeyi artıran ek depolama geçişlerine neden olur.

Akış iş yüklerini hata bekleyebileceğiniz şekilde tasarlama

Databricks, akış işlerini her zaman hata durumunda otomatik olarak yeniden başlatacak şekilde yapılandırmanızı önerir. Şema evrimi de dahil olmak üzere bazı özellikler, Yapılandırılmış Akış iş yüklerinin otomatik olarak yeniden denemesini gerektirir. Hata durumunda akış sorgularını yeniden başlatmak için bkz. Yapılandırılmış Akış işlerini yapılandırma.

Bazı işlemler foreachBatch gibi, tam olarak bir kez yerine en az bir kez garanti verir. Bu işlemler için işlem hattınızın idempotent olduğundan emin olun. Bkz. Rastgele veri havuzlarına yazmak için foreachBatch kullanma.

Not

Bir sorgu yeniden başlatıldığında, önceki çalıştırma sırasında planlanan mikro toplu işlem gerçekleştirilir. İşiniz bellek yetersiz hatası nedeniyle başarısız olduysa veya büyük boyutlu bir mikro toplu işlem nedeniyle işi el ile iptal ettiyseniz, mikro toplu işlemi başarıyla işlemek için işlemin ölçeğini artırmanız gerekebilir.

Çalıştırmalar arasındaki yapılandırmaları değiştirirseniz, bu yapılandırmalar planlanan ilk yeni toplu işleme uygulanır. Bkz. Yapılandırılmış Akış sorgusunda değişikliklerden sonra kurtarma işlemini.

Bir görev yeniden denendiğinde

bir Azure Databricks işinin parçası olarak birden çok görev zamanlayabilirsiniz. Sürekli tetikleyiciyi kullanarak bir işi yapılandırdığınızda, görevler arasında bağımlılık ayarlayamazsınız.

Aşağıdaki yaklaşımlardan birini kullanarak tek bir işte çok sayıda akışı zamanlamayı tercih edebilirsiniz:

  • Birden çok görev: Sürekli tetikleyiciyi kullanarak akış iş yüklerini çalıştıran birden çok görev içeren bir iş tanımlayın.
  • Birden çok sorgu: Tek bir görev için kaynak kodunda birden çok akış sorgusu tanımlayın.

Ayrıca bu stratejileri birleştirebilirsiniz. Aşağıdaki tablo bu yaklaşımları karşılaştırır.

Strateji Birden çok görev Birden çok sorgu
Hesaplama kaynakları nasıl paylaşılır? Databricks, her bir akış görevine uygun şekilde boyutlandırılmış işleri dağıtmanızı önerir. İsteğe bağlı olarak görevler arasında işlem gücü paylaşımı yapabilirsiniz. Tüm sorgular aynı hesaplamayı paylaşır. İsteğe bağlı olarak zamanlayıcı havuzlarına sorgu atayabilirsiniz.
Yeniden denemeler nasıl yönetilir? İş yeniden denenmeden önce tüm görevlerin başarısız olması gerekir. Herhangi bir sorgu başarısız olursa görev yeniden denenir.

Birden çok görev veya sorguyla çalışma hakkında daha fazla ayrıntı için bkz. Aynı kümede birden çok Yapılandırılmış Akış sorgusu çalıştırma.

Hata durumunda akış sorgularını yeniden başlatmak için Yapılandırılmış Akış işlerini yapılandırma

Databricks, sürekli tetikleyiciyi kullanarak tüm akış iş yüklerini yapılandırmanızı önerir. Bkz İşleri sürekli çalıştırma.

Sürekli tetikleyici varsayılan olarak aşağıdaki davranışa sahiptir:

  • İşin birden fazla eşzamanlı çalışmasını önler.
  • Önceki çalıştırma başarısız olduğunda yeni bir çalıştırma başlatır.
  • Yeniden denemeler için üstel geri çekilme stratejisi kullanır.

Databricks, iş akışlarını zamanlarken her zaman çok amaçlı işlem yerine iş hesaplama kaynaklarının kullanılmasını önerir. İş hatası ve yeniden deneme sırasında yeni işlem kaynakları dağıtılır.

Not

Databricks, streamingQuery.awaitTermination() veya spark.streams.awaitAnyTermination() kullanmamanızı önerir. Bkz . Ne zaman kullanılır awaitTermination()?

Ne zaman kullanılır? awaitTermination()

streamingQuery.awaitTermination() ve spark.streams.awaitAnyTermination() mevcut iş parçacığını, akış sorgusu sonlandığında kadar bloklar. Bu işlevlerin kullanılıp kullanılmaymayacağı, yürütme ortamınıza bağlıdır.

Lakeflow Jobs için streamingQuery.awaitTermination() veya spark.streams.awaitAnyTermination() kullanmayın. İşler hizmeti bir akış sorgusu etkin olduğunda bir çalıştırmanın tamamlanmasını otomatik olarak engellediğinden bu işlevler gerekli değildir. Her iki işlev de not defteri hücrelerinin tamamlanmasını engeller ve Görevler servisinin canlı veri sorgusunu izlemesine mani olur, bu da kapsam ölçümlerini ve iş bildirimlerini kesintiye uğratır.

Aşağıdaki durumlarda kullanın awaitTermination() :

Kullanım örneği Davranış
Tüm amaçlı hesaplamada etkileşimli not defterleri awaitTermination() hücreyi çalışır durumda tutar, sorgu durumunu gözlemlemenizi sağlar ve not defteri çıkışında hataların ortaya çıkarılmasını sağlar.
Yerel ve geliştirme ortamları Spark programını yerel olarak çalıştırırken, ana iş parçacığı tamamlandığında işlemden çıkılır. Akış sorgusu bitene veya başarısız olana kadar programı canlı tutmak için çağırın awaitTermination() .
Hatanın sürücüye iletilmesi awaitTermination() olmadan, iş bağlamı olmayan bir akış sorgusu hatası, çağıran iş parçacığına aktarılamayabilir. Sorgu sessizce başarısız olabilir ve hataların algılanıp tanılanabilmesini zorlaştırır. awaitTermination() çağrısı, sürücüde sorgu istisnasını yeniden oluşturur.