Not
Bu sayfaya erişim yetkilendirme gerektiriyor. Oturum açmayı veya dizinleri değiştirmeyi deneyebilirsiniz.
Bu sayfaya erişim yetkilendirme gerektiriyor. Dizinleri değiştirmeyi deneyebilirsiniz.
Yapılandırılmış akış, Spark üzerinde oluşturulmuş ölçeklenebilir, hataya dayanıklı bir akış işleme altyapısıdır. Canlı veri akışını, yeni satırların sürekli eklendiği bir tablo olarak ele alır. Yapılandırılmış Akış, CSV, JSON, ORC ve Parquet gibi yerleşik dosya kaynaklarının yanı sıra Kafka ve Azure Event Hubs gibi mesajlaşma hizmetlerini destekler.
Bu makale, Azure Event Hubs gibi bir akış kaynağı ayarlamayı, akış verilerini bir lakehouse Delta tablosuna almayı, bölümleme ve olay toplu işleme ile yazma performansını iyileştirmeyi ve akış işlerini üretimde güvenilir bir şekilde çalıştırmayı kapsar.
Akış kaynağı ayarlama
Bir göl evinde veri akışı yapmak için önce akış kaynağınıza bir bağlantı yapılandırın. Azure Event Hubs yaygın bir seçenektir. Spark uygulamanızı Azure Event Hubs'a bağlamak için Apache Spark için Azure Event Hubs Bağlayıcısı'nı kullanın.
Temel bir Event Hubs yapılandırması için Event Hubs ad alanı adı, hub adı, paylaşılan erişim anahtarı adı ve tüketici grubu gerekir.
Tüketici grubu, olay hub'ının tamamının görünümüdür. Tüketici grupları, her biri için birden çok tüketen uygulamanın olay akışının ayrı bir görünümüne sahip olmasını ve akışı kendi hızlarında ve kendi uzaklıklarıyla bağımsız olarak okumasını sağlar.
Event Hubs'daki bölümler, büyük hacimli olayları paralel olarak işlemenize olanak sağlar. Tek bir işlemci, saniye başına olayları işlemek için sınırlı kapasiteye sahipken, birden çok işlemci bölümler arasında paralel olarak çalışabilir.
Düşük alım hızıyla çok fazla bölüm kullanılıyorsa, bölüm okuyucular verilerin küçük bir kısmını işler ve bu da optimal olmayan işlemeye neden olur. İdeal bölüm sayısı, istenen işleme hızına bağlıdır. Ad alanınızdaki aktarım hızı birimi sayısını artırdıkça, eş zamanlı okuyucuların maksimum aktarım hızına ulaşmasına izin vermek için ek bölümler isteyebilirsiniz.
Aktarım hızı senaryonuz için en iyi bölüm sayısını test edin. Yüksek aktarım hızına sahip senaryolarda genellikle 32 veya daha fazla bölüm kullanılır.
Akış havuzu olarak Delta tablosu
Delta Lake, data lake storage üzerinde ACID (atomiklik, tutarlılık, izolasyon ve dayanıklılık) işlemleri sağlayan bir açık kaynak depolama katmanıdır. Doku Veri Mühendisliği'nde Delta Lake, upsert'leri, veri sıkıştırmayı, zaman yolculuğu, şema evrimi ve açık biçimli depolamayı destekler.
writeStream içinde delta çıkış formatı olarak, akış verileri doğrudan bir Delta tablosuna akar. Aşağıdaki örnek Event Hubs'dan okur, ileti gövdesini ayrıştırır ve Delta tablosuna yazar:
import pyspark.sql.functions as f
from pyspark.sql.types import *
df = (
spark.readStream
.format("eventhubs")
.options(**ehConf)
.load()
)
Schema = StructType([
StructField("<column_name_01>", StringType(), False),
StructField("<column_name_02>", StringType(), False),
StructField("<column_name_03>", DoubleType(), True),
StructField("<column_name_04>", LongType(), True),
StructField("<column_name_05>", LongType(), True),
])
rawData = (
df
.withColumn("bodyAsString", f.col("body").cast("string"))
.select(f.from_json("bodyAsString", Schema).alias("events"))
.select("events.*")
.writeStream
.format("delta")
.option("checkpointLocation", "Files/checkpoint")
.outputMode("append")
.toTable("deltaeventstable")
)
Kodda Delta'yı format("delta") çıkış biçimi olarak ayarlar, outputMode("append") tabloya yalnızca yeni satırlar yazar ve toTable("deltaeventstable") akış verilerini yönetilen delta tablosuna kalıcı hale döndürür.
Akış performansını iyileştirme
Temel akış alımı işe yaradıktan sonra, aşağıdaki bölümlerdeki iyileştirme teknikleriyle aktarım hızını ve dosya düzenlemesini geliştirebilirsiniz.
Yazma işlemleri için verileri bölümleme
Aktarım hızını iyileştirmek için verilerinizi etkili bir şekilde bölümleme. Bölümleme hem yazma aktarım hızını hem de aşağı akış sorgu performansını geliştirir. Verileri bellekte, diskte veya her ikisinde de bölümleyebilirsiniz.
Diskte — Verileri sütun değerlerine göre alt dizinler halinde düzenlemek için kullanın partitionBy() . En iyi şekilde boyutlandırılmış dosyalar üreten iyi kardinaliteye sahip sütunları seçin. Çok fazla küçük bölüm veya çok az büyük bölüm oluşturan sütunlardan kaçının.
Bellekte — Verileri yazmadan önce çalışan düğümler arasında dağıtmak için repartition() veya coalesce() kullanın.
-
repartition()tam karışıklığı olan bölümleri artırır veya azaltır, verileri eşit şekilde dengeler. -
coalesce()yalnızca bölümleri azaltır ve veri taşımayı en aza indirir.
Her iki yaklaşımın da birleştirilmesi, yüksek aktarım hızı senaryolarında iyi çalışır. Aşağıdaki örnek verileri bellekte 48 bölüme (kullanılabilir CPU çekirdekleriyle eşleşen) ve ardından diskteki bölümleri iki sütuna böler:
rawData = (
df
.withColumn("bodyAsString", f.col("body").cast("string"))
.select(f.from_json("bodyAsString", Schema).alias("events"))
.select("events.*")
.repartition(48)
.writeStream
.format("delta")
.option("checkpointLocation", "Files/checkpoint")
.outputMode("append")
.partitionBy("<column_name_01>", "<column_name_02>")
.toTable("deltaeventstable")
)
İyileştirilmiş Yazma kullanma
El ile bölümlemeye alternatif olarak, Optimize Edilmiş Yazma, disk aktarım hızını en üst düzeye çıkararak, el ile repartition() veya coalesce() çağrılar olmadan bölümleri yazmadan önce birleştirir veya bölümlere ayırır. Spark yapılandırmasıyla etkinleştirin:
spark.conf.set("spark.databricks.delta.optimizeWrite.enabled", True)
İyileştirilmiş Yazma etkinleştirildiğinde, kodunuzdan repartition() veya coalesce() kaldırabilir ve Spark'ın bölüm boyutlandırmayı işlemesine izin verebilirsiniz. Yine de disk düzeyinde kuruluş için kullanabilirsiniz partitionBy() .
Tetikleyicilerle toplu işleme olayları
Yazma performansını daha da iyileştirmek için, olayları diske yazmadan önce toplu işleyin. Varsayılan olarak, Spark her mikrobatch'i bir önceki tamamlandıktan hemen sonra işler. Tetikleyici zaman aralığı ayarlamak, belirli bir zaman dilimindeki verileri biriktirir ve onları daha az sayıda, daha büyük işlemlerle yazar. Daha büyük toplu işlemler daha büyük Delta dosyaları oluşturur ve küçük dosya ek yükünü azaltır.
Aşağıdaki örnek olayları bir dakikalık aralıklarla işler:
rawData = (
df
.withColumn("bodyAsString", f.col("body").cast("string"))
.select(f.from_json("bodyAsString", Schema).alias("events"))
.select("events.*")
.writeStream
.format("delta")
.option("checkpointLocation", "Files/checkpoint")
.outputMode("append")
.partitionBy("<column_name_01>", "<column_name_02>")
.trigger(processingTime="1 minute")
.toTable("deltaeventstable")
)
Gelen verilerin hacmini analiz edin ve Delta tablosunda iyi boyutlandırılmış Parquet dosyaları üreten bir işlem aralığı seçin.
Üretimde streaming işlerini çalıştırma
Spark notebook'lar, akış mantığını geliştirmek ve test etmek için etkili bir araçtır. Ancak sürekli çalışması gereken üretim iş yükleri için bunun yerine Spark iş tanımlarını kullanın. Spark iş tanımları, Bir Spark kümesinde çalıştırılan ve daha fazla sağlamlık ve kullanılabilirlik sağlayan etkileşimli olmayan, kod odaklı görevlerdir.
Akış işini çalıştıran altyapı, donanım arızaları veya altyapı yamaları gibi işi durduran sorunlarla karşılaşabilir. Yeniden deneme ilkesi, beklenmedik bir şekilde durduğunda işi otomatik olarak yeniden başlatır. Spark iş tanımında yeniden deneme ilkesini, işin kaç kez yeniden başlatılacağını (sonsuz yeniden denemelere kadar) ve yeniden denemeler arasındaki zaman aralığını belirtmek için yapılandırın. Yeniden deneme ilkesi etkinleştirildiğinde, siz açıkça durdurana kadar akış işiniz çalışmaya devam eder.
Doku izleme hub'ı Giriş Hızı, İşlem Hızı, Giriş Satırları, Toplu İş Süresi ve İşlem Süresi gibi ölçümleri içeren bir Yapılandırılmış Akış sekmesi içerir.