Apache Spark ile değişiklik besleme

Azure Cosmos DB Spark Bağlayıcısı, Apache Spark kullanarak değişiklik akışını uygun ölçekte işlemek için güçlü bir yol sağlar. Bağlayıcı altında Java SDK'sını kullanır ve spark yürütücüleri arasında işlemeyi saydam bir şekilde dağıtan bir çekme modeli uygulayarak büyük ölçekli veri işleme senaryoları için idealdir.

Spark Bağlayıcısı nasıl çalışır?

Azure Cosmos DB için Spark Bağlayıcısı, Azure Cosmos DB Java SDK'sının üzerine oluşturulur ve değişiklik akışını okumak için bir çekme modeli yaklaşımı uygular. Temel özellikler şunlardır:

  • Java SDK'sı temeli: Güvenilir değişiklik akışı işleme için altındaki güçlü Azure Cosmos DB Java SDK'sını kullanır
  • Çekme modeli uygulaması: İşlem hızını kontrol etmenizi sağlayan değişiklik akışı çekme modeli modelini izler
  • Dağıtılmış işleme: Paralel işleme için değişiklik akışı işlemeyi birden çok Spark yürütücüslerine otomatik olarak dağıtır
  • Saydam ölçeklendirme: Bağlayıcı, bölümleme ve yük dağıtımını manuel müdahale gerektirmeden işler

Benzersiz denetim noktası oluşturma özelliği

Spark Bağlayıcısı'nı değişiklik akışı işleme için kullanmanın temel avantajlarından biri, yerleşik denetim noktası mekanizmasıdır. Bu özellik şunları sağlar:

  • Otomatik kurtarma: Değişiklik akışını büyük ölçekte işlerken kurtarma için kullanıma hazır mekanizma
  • Hataya dayanıklılık: Hata durumunda işlemeyi son denetim noktasından sürdürme olanağı
  • Durum yönetimi: Spark oturumları ve küme yeniden başlatmaları boyunca işlem durumunu korur.
  • Ölçeklenebilirlik: Dağıtılmış Spark ortamlarında denetim noktası oluşturmayı destekler

Bu denetim noktası oluşturma özelliği Spark Bağlayıcısı'na özgüdür ve SDK'ları doğrudan kullanırken kullanılamaz ve yüksek kullanılabilirlik ve güvenilirlik gerektiren üretim senaryoları için özellikle değerlidir.

Warning

Denetim spark.cosmos.changeFeed.startFrom noktası konumunda mevcut yer işaretleri varsa yapılandırma yoksayılır. Bir denetim noktasından devam ederken, bağlayıcı belirtilen başlangıç noktası yerine son işlenen konumdan devam eder.

Değişiklik akışı işleme için Spark ne zaman kullanılır?

Bu senaryolarda değişiklik akışı işleme için Spark Bağlayıcısı'nı kullanmayı göz önünde bulundurun:

  • Büyük ölçekli veri işleme: Tek makineli özellikleri aşan yüksek hacimli değişiklik akışı verilerini işlemeniz gerektiğinde
  • Karmaşık dönüştürmeler: Değişiklik akışı işlemeniz karmaşık veri dönüştürmeleri, toplamalar veya diğer veri kümeleriyle birleştirmeler içeriyorsa
  • Dağıtılmış analiz: Dağıtılmış bir ortamda değişiklik akışı verileri üzerinde gerçek zamanlı veya neredeyse gerçek zamanlı analiz gerçekleştirmeniz gerektiğinde
  • Veri hatları ile tümleştirme: Değişiklik akışı işleme süreci, zaten Spark kullanan daha büyük ETL/ELT veri hatlarının bir parçası olduğunda
  • Hataya dayanıklılık gereksinimleri: Üretim iş yükleri için sağlam denetim noktası oluşturma ve kurtarma mekanizmalarına ihtiyacınız olduğunda
  • Çok kapsayıcılı işleme: Birden çok kapsayıcıdan gelen değişiklik akışlarını aynı anda işlemeniz gerektiğinde

Daha basit senaryolar için veya tek tek belge işleme üzerinde ayrıntılı denetime ihtiyacınız olduğunda , değişiklik akışı işlemcisini veya çekme modelini doğrudan SDK'larla kullanmayı göz önünde bulundurun.

Kod örnekleri

Aşağıdaki örneklerde Spark Bağlayıcısı'nı kullanarak değişiklik akışından nasıl okunduğu gösterilmektedir. Daha kapsamlı örnekler için örnek not defterlerinin tamamına bakın:

# Configure change feed reading

changeFeedConfig = {
    "spark.cosmos.accountEndpoint": "https://<account-name>.documents.azure.com:443/",
    "spark.cosmos.accountKey": "<account-key>",
    "spark.cosmos.database": "<database-name>",
    "spark.cosmos.container": "<container-name>",
    # Start from beginning, now, or specific timestamp (ignored if checkpoints exist)
    "spark.cosmos.changeFeed.startFrom": "Beginning",  # "Now" or "2020-02-10T14:15:03"
    "spark.cosmos.changeFeed.mode": "LatestVersion",  # or "AllVersionsAndDeletes"
    # Control batch size - if not set, all available data processed in first batch
    "spark.cosmos.changeFeed.itemCountPerTriggerHint": "50000",
    "spark.cosmos.read.partitioning.strategy": "Restrictive"
}

# Read change feed as a streaming DataFrame
changeFeedDF = spark \
    .readStream \
    .format("cosmos.oltp.changeFeed") \
    .options(**changeFeedConfig) \
    .load()

# Configure output settings with checkpointing
outputConfig = {
    "spark.cosmos.accountEndpoint": "https://<target-account>.documents.azure.com:443/",
    "spark.cosmos.accountKey": "<target-account-key>",
    "spark.cosmos.database": "<target-database>",
    "spark.cosmos.container": "<target-container>",
    "spark.cosmos.write.strategy": "ItemOverwrite"
}

# Process and write the change feed data with checkpointing
query = changeFeedDF \
    .selectExpr("*") \
    .writeStream \
    .format("cosmos.oltp") \
    .outputMode("append") \
    .option("checkpointLocation", "/tmp/changefeed-checkpoint") \
    .options(**outputConfig) \
    .start()

# Wait for the streaming query to finish
query.awaitTermination()

Anahtar yapılandırma seçenekleri

Spark'ta değişiklik akışıyla çalışırken bu yapılandırma seçenekleri özellikle önemlidir:

  • spark.cosmos.changeFeed.startFrom: Değişiklik akışını okumaya nereden başlayacağınızı denetler
    • "Beginning" - Değişiklik akışının başından başlayın
    • "Now" - Geçerli saatten başlayarak
    • "2020-02-10T14:15:03" - Belirli bir zaman damgasından başlayın (ISO 8601 biçimi)
    • Not: Denetim noktası konumunda mevcut yer işaretleri varsa bu ayar yok sayılır
  • spark.cosmos.changeFeed.mode: Değişiklik akışı modunu belirtir
    • "LatestVersion" - Değiştirilen belgelerin yalnızca en son sürümünü işleme
    • "AllVersionsAndDeletes" - Silmeler de dahil olmak üzere değişikliklerin tüm sürümlerini işleme
  • spark.cosmos.changeFeed.itemCountPerTriggerHint: Toplu işlem boyutunu denetler
    • Her mikro toplu iş/tetikleyici için değişiklik akışından okunan yaklaşık en fazla öğe sayısı
    • Örnek: "50000"
    • Önemli: Ayarlanmadıysa, değişiklik akışındaki tüm kullanılabilir veriler ilk mikro toplu işte işlenir
  • checkpointLocation: Hataya dayanıklılık ve kurtarma için denetim noktası bilgilerinin depolandığı yeri belirtir
  • spark.cosmos.read.partitioning.strategy: Verilerin Spark yürütücüleri arasında nasıl bölümlendiğini denetler

Sonraki Adımlar