Spark #5 - Resilient Distributed Dataset (RDD)
RDD’ler oluşturulduktan sonra durumu güncellenemeyen (immutable) dağıtık nesneler koleksiyonudur. Her RDD parçalara ayrılarak farklı hesaplama düğümleri tarafından işlenecek şekilde otomatik olarak dağıtılır. RDD’ler Python, Java, Scala veya kullanıcı tanımlı objeleri içerebilir ve bunlar üzerinde çalışabilir. RDD’ ler herhangi bir kaynaktan veri yükleyerek veya obje koleksiyonları (liste, set, küme) kullanılarak oluşturulabilir.
Örneğin:
lines=sc.textFile(“C:\Spark\README.md”)
Kullanımı ile text dosyası olunarak lines adlı bir RDD oluşturulur. RDD’ ler oluşturulduktan sonra dönüşümler ve aksiyonlar ile (Transformation, action) işlenirler. Dönüşümler ile kaynak RDD’ den yeni bir RDD oluşturulur. Daha önce kullandığımız filter operasyonu yeni bir RDD üreten bir dönüşümdür:
Sparklines=lines.filter(lambda satir: “Spark” in satir)
Aksiyonlar ise RDD’ yi kullanarak bir sonuç üretirler. Bu sonuç ana programa geri döndürülebilir veya HDFS gibi bir dış saklama sistemine kayıt edilebilir.
sparklines.first() ise bir aksiyondur. Hesaplamalar dönüşümler ve aksiyonlar için farklı gerçekleşir. RDD oluşturuyor gibi görünsek de, RDD’ nin hesaplanması bir aksiyon çağırıldığında gerçekleştirilir. (Lazy evaluation) Böylece ardı ardına oluşturulan RDD zincirinden sadece gerekli olan, doğru zamanda hesaplandığı için boş yere hafıza kaybı yaşanmaz. Spark core engine tüm dönüşümü bir bütün olarak görür ve sadece son aksiyon için ihtiyaç duyulan hesap gerçekleştirilir.
Örneğin:
lines=sc.textFile(“C:\Spark\README.md”)
sparklines=lines.filter(lambda line: “Spark” in line)
sparklines.first()
kod bloğunda tüm dosya okunmaz. İlk “Spark” bulunana satırda işlem sonlandırılır. Spark RDD’ leri aksiyonlar kullanıldığında yeniden hesaplar. Eğer RDD’ nin birden fazla aksiyonda kullanılmasını planlıyorsanız
Rdd.persist() veya Rdd.cache() diyerek kalıcı olarak saklanmasını sağlayabilirsiniz
.RDD’ leri hafıza yerine (ram) Diskte saklamak da mümkündür. Aslında Spark’ ın RDD’ leri yeniden hesaplaması “Resilient” kavramı ile ilgilidir. RDD veya herhangi bir düğümdeki bir RDD yi saklayan makine de bir sorun yaşandığında, Spark kullanıcının haberi olmadan eksik bölümü yeniden hesaplar. RDD oluşturmanın diğer bir yöntemi ise spark.paralelize() metodunu kullanmaktır. Bu metoda parametre olarak kendi veri yapılarımızı (array, list vb.) vererek RDD oluşturabiliriz.
sc.paralelize([“To be”,”or”,”not to be”]) Python
sc.paralelize(Arrays.asList(“To be”,”or”,”not to be”)) Java
Bir sonraki yazımızda dönüşüm ve aksiyonları daha detaylı anlatmayı ve RDD bağımlılıklarını saklayan lineage graf’ ından bahsetmeyi son not olarak ekleyelim.
Görüşmek Üzere.












