Üretim iş yükleri için Otomatik Yükleyici'yi yapılandırma

Dosya bulma modu seçimi, şema yönetimi ve veri kalitesi işleme dahil olmak üzere Otomatik Yükleyici'yi ayarlamaya yönelik kapsamlı en iyi yöntemler için bkz. Otomatik Yükleyici en iyi yöntemleri.

Databricks, artımlı veri alımı için Lakeflow işlem hatlarında Otomatik Yükleyici kullanılmasını önerir. Lakeflow işlem hatları, Apache Spark Structured Streaming’in işlevselliğini genişletir ve bildirimsel Python veya SQL’in yalnızca birkaç satırını yazarak aşağıdakilere sahip üretim kalitesinde bir veri işlem hattını devreye almanıza olanak tanır:

Databricks, üretimde Otomatik Yükleyici'yi çalıştırmak için en iyi akış yöntemlerini izlemenizi de önerir. Bkz Yapılandırılmış Akış için Üretimle İlgili Dikkat Edilmesi Gerekenler.

Note

Lakeflow boru hatları, çoğu üretim ortamındaki veri alımı için Otomatik Yükleyici'yi çalıştırmanın önerilen yöntemidir. İş yükünüz düşük gecikme süresi gereksinimlerine sahip değilse ve önceliğiniz işlem maliyetini en aza indiriyorsa, bunun yerine Otomatik Yükleyici'yi kullanan Trigger.AvailableNowtetiklenmiş bir toplu iş olarak zamanlayabilirsiniz. Bkz. Maliyetle ilgili dikkat edilmesi gerekenler.

Otomatik Yükleyiciyi İzleme

Aşağıdaki bölümlerde ölçümler, günlükler, uyarılar ve yaygın sorun giderme iş akışları dahil olmak üzere üretimde Otomatik Yükleyici'nin nasıl izleneceği açıklanmaktadır. Pano desenlerini, gecikme analizini ve şema kayması algılamayı kapsayan kapsamlı bir başvuru için bkz. Otomatik Yükleyici'yi izleme ve gözlemleme.

Otomatik Yükleyici tarafından bulunan dosyaları sorgulama

Otomatik Yükleyici, bir akışın durumunu incelemek için bir SQL API'sini sağlar. cloud_files_state işlevini kullanarak, Otomatik Yükleyici akışı tarafından bulunan dosyalar hakkındaki meta verileri bulabilirsiniz. Otomatik Yükleyici akışıyla ilişkili denetim noktası konumunu sağlayarak sorgulayın cloud_files_state.

Note

cloud_files_state işlevi Databricks Runtime 11.3 LTS ve üzeri sürümleriyle kullanılabilir.

SELECT * FROM cloud_files_state('path/to/checkpoint');

Akış güncellemelerini dinleyin

Databricks, Otomatik Yükleyici akışlarını daha fazla izlemek için Apache Spark'ın Akış Sorgu Dinleyicisi arabiriminin kullanılmasını önerir. Bkz. Azure Databricks üzerinde Yapılandırılmış Akış sorgularını izleme.

Otomatik Yükleyici, ölçümleri her toplu işlemde Akış Sorgusu Dinleyicisi'ne bildirir. Bekleyen işlerin arasında kaç dosya olduğunu ve beklemenin ne kadar büyük olduğunu, akış sorgusu ilerleme panosundaki numFilesOutstanding sekmesinde, numBytesOutstanding ve ölçümlerinde görüntüleyebilirsiniz.

{
  "sources": [
    {
      "description": "CloudFilesSource[/path/to/source]",
      "metrics": {
        "numFilesOutstanding": "238",
        "numBytesOutstanding": "163939124006"
      }
    }
  ]
}

Databricks Runtime 10.4 LTS ve üzerinde dosya bildirim modunu kullanırken ölçümler, aws ve Azure için approximateQueueSize olarak bulut kuyruğundaki yaklaşık dosya olayı sayısını da içerir.

Maliyetle ilgili dikkat edilmesi gerekenler

Otomatik Yükleyici'yi çalıştırırken, ana maliyet kaynaklarınız işlem kaynakları ve dosya bulmadır.

İş yükünüzün düşük gecikme gereksinimleri yoksa, Auto Loader’ı sürekli çalıştırmak yerine Trigger.AvailableNow kullanarak toplu işler olarak zamanlamak için Lakeflow Jobs’u kullanıp işlem maliyetlerinizi azaltabilirsiniz. Bkz . Yapılandırılmış Akış tetikleyici aralıklarını yapılandırma. Bu toplu işler , dosya varış ve işleme arasındaki gecikme süresini daha da azaltmak için dosya varış tetikleyicileri kullanılarak tetiklenebilir.

Dosya bulma maliyetleri, depolama hesaplarınızda dizin listeleme modunda LIST işlemleri, abonelik hizmetinde API istekleri ve dosya bildirim modunda kuyruk hizmeti biçiminde gelebilir. Gibi sürekli tetikleyiciler Trigger.ProcessingTime özellikle dizin listeleme modunda pahalıdır, çünkü Otomatik Yükleyici yeni dosyaları bulmak için dizinin tamamını sürekli listeler. İş yükünüz sürekli tetikleyiciler gerektiriyorsa Databricks, gecikme süresi gereksinimlerinize göre bir dosya bulma modu seçmenizi önerir:

  • Düşük gecikme süresi ve basitlik: Dosya olaylarıyla Otomatik Yükleyici'yi kullanın. Dosya olayları demet başına yalnızca bir kuyruk gerektirir ve sonraki çalıştırmalarda artımlı bulma kullanır. Daha fazla bilgi için bkz. Dosya olaylarıyla otomatik yükleyiciye genel bakış.
  • Gecikme süresine duyarlı uygulamalar: Klasik dosya bildirim modunu kullanın. Klasik mod, dosya olaylarının neden olduğu ek önbellekleme adımı olmadan doğrudan bulut kuyruğundan okur. Bu modda, kaynak etiketlerini kullanarak maliyetlerinizi izlemek için Otomatik Yükleyici tarafından oluşturulan kaynakları etiketleyebilirsiniz. Ayrıntılar için bkz. Dosya bildirimi.

Kaynak veri saklama

Note

Databricks Runtime 16.4 LTS ve üzerinde kullanılabilir.

Dosyalar kaynak dizininizde biriktikçe depolama maliyetleri artar ve özellikle dizin listeleme modunda dosya bulma yavaşlar. Otomatik Yükleyici, dosyaları işlendikten cloudFiles.cleanSource sonra arşivleyerek veya silerek dosya saklamayı otomatik olarak yönetme seçeneği sunar.

Maliyetleri düşürmek için kaynak dizindeki dosyaları arşivleme

Warning

  • Ayarı cloudFiles.cleanSource , kaynak dizindeki dosyaları siler veya taşır.
  • Veri işleme için kullanıyorsanız foreachBatch , işleminiz toplu işlemdeki dosyaların yalnızca bir alt kümesini kullansa bile, işleminiz foreachBatch başarıyla döndürüldüğü anda dosyalarınız taşıma veya silme adayları haline gelir.

Databricks, bulma maliyetlerini azaltmak için dosya olaylarıyla Otomatik Yükleyici'nin kullanılmasını önerir. Bu, keşif sürecinin artımlı olması nedeniyle hesaplama maliyetlerini de azaltır.

Dosya olaylarını kullanamıyorsanız ve dosyaları bulmak için dizin listesini kullanmanız gerekiyorsa, Otomatik Yükleyici bunları işledikten sonra bulma maliyetlerini düşürmek için dosyaları otomatik olarak arşivleme veya silme seçeneğini kullanabilirsiniz cloudFiles.cleanSource . Otomatik Yükleyici işlendikten sonra kaynak dizininizdeki dosyaları temizlediğinden, bulma sırasında daha az dosyanın listelenmesi gerekir.

cloudFiles.cleanSource ve MOVE seçeneğini kullanırken, aşağıdaki gereksinimleri göz önünde bulundurun:

  • Hem kaynak dizin hem de hedef taşıma dizini aynı demette veya kapsayıcıda bulunmalıdır. Çapraz demetler ve kapsayıcılar arası taşımalar desteklenmez ve hataya neden olur.
  • Taşıma hedefi bir birim yolu olabilir (örneğin, /Volumes/my_catalog/my_schema/my_volume/archive/).
  • Kaynak ve hedef dizininiz aynı dış konumdaysa, yönetilen depolama alanı (örneğin, yönetilen birim veya katalog) içeren eşdüzey dizinlere sahip olmamalıdır. Bu gibi durumlarda, Otomatik Yükleyici hedef dizine yazmak için gerekli izinleri alamaz.

Databricks aşağıdaki durumlarda bu seçeneğin kullanılmasını önerir:

  • Kaynak dizininiz zaman içinde çok sayıda dosya biriktirir.
  • İşlenen dosyaları uyumluluk veya denetim için tutmanız gerekir (cloudFiles.cleanSource öğesini MOVE olarak ayarlayın).
  • Alımdan sonra dosyaları kaldırarak depolama maliyetlerini azaltmak istiyorsunuz (cloudFiles.cleanSource'i DELETE olarak ayarlayın). Databricks, DELETE modunu kullanırken, Auto Loader silmelerinin yumuşak silme olarak davranabilmesi ve hatalı yapılandırma durumunda erişilebilir olması için kova üzerinde sürüm oluşturmayı etkinleştirmeyi önerir. Ayrıca Databricks, kurtarma gereksinimlerinize göre belirli bir yetkisiz kullanım süresi (60 veya 90 gün gibi) sonra eski, geçici olarak silinen sürümleri temizlemek için bulut yaşam döngüsü ilkeleri ayarlamanızı önerir.

Tam referans için cleanSource seçenekler ve varsayılanları hakkında bkz. cloudFiles.cleanSource.

İşlenen dosyaları soğuk depolama yoluna taşıma

Aşağıdaki örnek, 14 gün sonra işlenen dosyaları aynı demet içinde bir arşiv dizinine taşımak için Otomatik Yükleyici'yi yapılandırılır. Dosyaları daha ucuz depolama katmanlarına (örneğin AWS S3 Glacier, Azure Cool/Archive veya GCS Coldline/Archive) aktarmak için arşiv yoluna bir bulut yaşam döngüsü ilkesi uygulayabilirsiniz.

Python

# Step 1: Configure Auto Loader to move processed files to an archive path.
checkpoint = "/Volumes/my_catalog/my_schema/my_volume/checkpoints/ingest_stream"
archive_path = "s3://my-bucket/archive/landing/"

df = (spark.readStream.format("cloudFiles")
  .option("cloudFiles.format", "json")
  .option("cloudFiles.cleanSource", "MOVE")
  .option("cloudFiles.cleanSource.moveDestination", archive_path)
  .option("cloudFiles.cleanSource.retentionDuration", "14 days")
  .option("cloudFiles.schemaLocation", checkpoint)
  .load("s3://my-bucket/landing/")
)

# Step 2: Write to a Delta table.
(df.writeStream
  .option("checkpointLocation", checkpoint)
  .trigger(availableNow=True)
  .toTable("my_catalog.my_schema.raw_events")
)

# Step 3 (outside Databricks): Set up a cloud lifecycle policy on the
# archive path to transition files to cold storage after a grace period.
# For example, in AWS you can configure an S3 Lifecycle rule to move
# objects under s3://my-bucket/archive/landing/ to S3 Glacier after
# 30 days.

SQL

-- Step 1: Configure Auto Loader to move processed files to an archive path
-- using a Lakeflow Declarative Pipeline.
CREATE OR REFRESH STREAMING TABLE raw_events
AS SELECT * FROM STREAM read_files(
  's3://my-bucket/landing/',
  format => 'json',
  cleanSource => 'MOVE',
  cleanSourceMoveDestination => 's3://my-bucket/archive/landing/',
  cleanSourceRetentionDuration => '14 days'
);

-- Step 2 (outside Databricks): Set up a cloud lifecycle policy on the
-- archive path to transition files to cold storage.
-- For example, in AWS configure an S3 Lifecycle rule to move objects
-- under s3://my-bucket/archive/landing/ to S3 Glacier after 30 days.

Trigger.AvailableNow ve hız sınırlamayı kullanma

Note

Databricks Runtime 10.4 LTS ve üzerinde kullanılabilir.

Otomatik Yükleyici, Trigger.AvailableNow kullanılarak Lakeflow İşleri'nde toplu iş olarak çalıştırılmak üzere zamanlanabilir. Tetikleyici, AvailableNow Otomatik Yükleyici'ye sorgu başlangıç saatinden önce gelen tüm dosyaları işlemesini emreder. Akış başladıktan sonra gelen yeni dosyalar, bir sonraki tetikleyici anına kadar göz ardı edilir.

ile Trigger.AvailableNow, dosya bulma veri işleme ile zaman uyumsuz olarak gerçekleşir ve veriler hız sınırlaması olan birden çok mikro toplu işlemde işlenebilir. Otomatik Yükleyici varsayılan olarak her mikro toplu iş için en fazla 1000 dosya işler. Bir mikro toplu işlemde kaç dosya veya kaç bayt işlenmesi gerektiğini yapılandırabilir cloudFiles.maxFilesPerTrigger ve cloudFiles.maxBytesPerTrigger yapılandırabilirsiniz. Dosya sınırı katı bir sınırdır, ancak bayt sınırı esnek bir sınırdır; yani sağlanan maxBytesPerTrigger'dan daha fazla bayt işlenebilir. Her iki seçenek de birlikte sağlandığında, Otomatik Yükleyici sınırlardan birine ulaşmanız için gereken sayıda dosyayı işler.

Denetim noktası konumu

Denetim noktası konumu, akışın durum ve ilerleme bilgilerini depolamak için kullanılır. Databricks, denetim noktası konumunu bulut nesnesi yaşam döngüsü ilkesi olmayan bir konuma ayarlamanızı önerir. Denetim noktası konumundaki dosyalar ilkeye göre temizlenirse akış durumu bozulur. Böyle bir durumda akışı sıfırdan yeniden başlatmanız gerekir.

Dosya olayı takibi

Otomatik Yükleyici, tam olarak bir kez alma garantisi sağlamak için RocksDB kullanarak bulunan dosyaları kontrol noktası konumunda izler. Yüksek hacimli veya uzun ömürlü alım akışları için Databricks, Databricks Runtime 15.4 LTS veya üzeri sürümlere yükseltmenizi önerir. Bu sürümlerde Otomatik Yükleyici, akış başlamadan önce RocksDB durumunun tamamının indirilmesi için beklemez ve bu da akış başlatma süresini hızlandırabilir. Dosya durumlarının sınırsız olarak büyümesini önlemek istiyorsanız, belirli bir yaştan cloudFiles.maxFileAge daha eski dosya olaylarının süresinin dolmasını sağlama seçeneğini de kullanabilirsiniz. cloudFiles.maxFileAge için ayarlayabileceğiniz en düşük değer "14 days"değeridir. RocksDB’deki silme işlemleri, mezar taşı girdileri olarak görünür. Bu nedenle, olayların süresi dolmadan önce depolama kullanımının geçici olarak arttığını görebilirsiniz, ardından düzey kaybolmaya başlayacaktır.

Warning

cloudFiles.maxFileAge yüksek hacimli veri kümeleri için bir maliyet denetimi mekanizması olarak sağlanır. Çok agresif bir şekilde ayarlamak cloudFiles.maxFileAge , yinelenen veri alımı veya eksik dosyalar gibi veri kalitesi sorunlarına neden olabilir. Bu nedenle Databricks, karşılaştırmalı veri alımı çözümlerinin önerdiğine benzer şekilde 90 gün gibi muhafazakar bir ayar cloudFiles.maxFileAgeönerir.

Seçeneği ayarlamaya cloudFiles.maxFileAge çalışmak, işlenmemiş dosyaların Otomatik Yükleyici tarafından yoksayılmasını veya zaten işlenmiş dosyaların süresinin dolmasına ve sonra yinelenen verilere neden olarak yeniden işlenmesine neden olabilir. Bir seçim cloudFiles.maxFileAgeyaparken göz önünde bulundurmanız gereken bazı şeyler şunlardır:

  • Akışınız uzun bir süre sonra yeniden başlatılırsa, kuyruktan çekilen ve cloudFiles.maxFileAge tarihinden daha eski olan dosya bildirim etkinlikleri yoksayılır. Benzer şekilde, dizin listeleme kullanıyorsanız, çalışma dışı süre sırasında ortaya çıkmış olabilecek ve cloudFiles.maxFileAge tarihinden daha eski olan dosyalar yoksayılır.
  • Dizin listeleme modunu kullanıyorsanız ve cloudFiles.maxFileAgekullanıyorsanız , örneğin "1 month"olarak ayarlanmışsa, akışınızı durdurur ve cloudFiles.maxFileAgeolarak ayarlanmış "2 months" ile akışı yeniden başlatırsınız, ancak 1 aydan eski, ancak 2 aydan daha yeni olan dosyalar yeniden işlenir.

Akışı ilk kez başlattığınızda bu seçeneği ayarlarsanız, cloudFiles.maxFileAge'den eski verileri almayacaksınız. Bu nedenle, eski verileri almak istiyorsanız akışınızı ilk kez başlatırken bu seçeneği ayarlamamalısınız. Ancak, sonraki çalıştırmalarda bu seçeneği ayarlamanız gerekir.

CloudFiles.backfillInterval kullanarak normal geri doldurmaları tetikleme

Nadir durumlarda, yalnızca bildirim sistemlerine bağlı olunduğunda, bildirim iletisi saklama sınırlarına ulaşıldığında dosyalar kaybolabilir veya gecikebilir. Veri eksiksizliği ve SLA ile ilgili katı gereksinimleriniz varsa, belirli bir aralıkta zaman uyumsuz geri doldurmaları tetikleme ayarını cloudFiles.backfillInterval göz önünde bulundurun. Örneğin, günlük doldurmalar için bir gün veya haftalık geri doldurmalar için bir hafta olarak ayarlayın. Normal geri doldurmaların tetiklenmesi yinelemelere neden olmaz.

Dosya olaylarını kullanırken, akışınızı en az 7 günde bir çalıştırın

Dosya olaylarını kullanırken tam dizin listelemesini önlemek için Otomatik Yükleyici akışlarınızı en az 7 günde bir çalıştırın. Otomatik Yükleyici akışlarınızı bu sıklıkla çalıştırmak, dosya bulma işleminin artımlı olmasını sağlar.

Kapsamlı yönetilen dosya olayları en iyi yöntemleri için bkz. Dosya olaylarıyla Otomatik Yükleyici için en iyi yöntemler.