Google PDE: Spark, Dataproc ve Dağıtık Veri İşleme — Çalışma kılavuzu
Şunun bir parçası: Google Professional Data Engineer — Çalışma kılavuzu. Doğrulanmış cevaplarla şurada pratik yapın: Google sınav merkezi, veya şurada süreli deneme sınavları çözün: ExamRoll.io.
Genel Bakış
Google Cloud Dataproc üzerindeki Apache Spark, dağıtık veri işleme için yönetilen, elastik bir platform sağlar. Kontrol ihtiyaçlarına, çalışma zamanı değişkenliğine ve yönetimsel ek yüke bağlı olarak, uzun süreli veya geçici Dataproc kümeleri ile Spark için Dataproc Serverless arasında seçim yapabilirsiniz. Spark, dayanıklı soyutlamalar (RDD’ler), ilişkisel API’ler (DataFrames ve Spark SQL) ve büyük ölçekte yinelemeli ve toplu ETL için optimize edilmiş, hataya dayanıklı bir DAG yürütme motoru sunar. Google Cloud’da, Cloud Storage dayanıklı, düşük maliyetli depolama için HDFS’in yerini alır; BigQuery bağlayıcısı doğrudan analitik yük boşaltmayı mümkün kılar; ve Dataproc Metastore şema yönetimini merkezileştirir. Etkili çözümler, depolama ve işlem yaşam döngülerini uyumlu hale getirir, Spark’ı iş yüküne göre ayarlar, gözlemlenebilirliği araçlarla donatır ve en az ayrıcalık ilkesi ve ağ izolasyonu ile güvenliği uygular.
Dataproc Mimarisi: Kümeler, Sunucusuz, Depolama ve Metastore
- Küme türleri ve düğüm rolleri
- Birincil (ana) düğümler YARN, HDFS NameNode (kullanılıyorsa) ve Spark sürücü arayüzlerini barındırır; HA modu birden çok birincil düğüm kullanır.
- Çalışan düğümler, yürütücüleri (executor) ve HDFS DataNode’larını (kullanılıyorsa) çalıştırır.
- İkincil/yardımcı çalışanlar, HDFS rolleri olmaksızın elastik, daha düşük maliyetli kapasite için genellikle kesintiye uğratılabilir/spot (preemptible/spot) düğümlerdir.
- İmajlar, işletim sistemi ve bileşen sürümlerini bir araya getirir (örneğin, 2.1-debian11, 2.2-ubuntu20); Spark/Hadoop uyumluluğunu kontrol etmek ve yükseltmeleri bilinçli bir şekilde yapmak için imaj sürümlerini sabitleyin.
- Component Gateway, kullanıcı arayüzlerini (Spark History Server, YARN RM) HTTPS üzerinden güvenli bir şekilde yayınlar.
- Spark için Dataproc Serverless
- Küme oluşturma yok, otomatik ölçeklendirme ve yürütücüler (executor) ile sürücüler (driver) için saniye bazında faturalandırma. Düzensiz veya anlık yoğunluklu işler ya da operasyonel ek yükü en aza indirmek istendiğinde idealdir.
- Dezavantajları: kümelere göre daha az alt seviye ayar imkanı; iş başlangıç gecikmesi sıcak bir kümeden daha yüksek olabilir; sorun giderme için sunucusuz metrikleri ve olay günlüklerini kullanın.
- Otomatik ölçeklendirme
- Küme otomatik ölçeklendirme politikaları, YARN/Spark metriklerine ve soğuma sürelerine göre çalışan ekler/kaldırır ve birincil ile ikincil çalışan gruplarını ayrı ayrı ayarlar.
- Sunucusuz otomatik ölçeklendirme hizmet tarafından yönetilir; en iyi ölçeklendirme için bölüm-paralel (partition-parallel) olacak şekilde tasarlayın ve serileştirilmiş darboğazlardan kaçının.
- Depolama ve bağlayıcılar
- Kayıt sistemi (system-of-record) olarak Google Cloud Storage’ı (GCS) tercih edin; bu, işlemi depolamadan ayırır, kalıcı disk maliyetini azaltır ve küme yaşam döngülerinden bağımsız olarak varlığını sürdürür.
- GCS bağlayıcısı (gs://), Hadoop/Spark ile entegre olur. Nesne depolarına yazma işlemleri, işleme (commit) protokollerini kullanır; GCS üzerinde yeniden adlandırma ek yükünü azaltmak ve işin tamamlanmasını (commit) hızlandırmak için FileOutputCommitter algoritmasını v2 olarak ayarlayın:
--conf mapreduce.fileoutputcommitter.algorithm.version=2
```
- Sütun ayıklama (column pruning) ve koşul itme (predicate pushdown) ile Parquet/ORC kullanın. Verimli tarama için dosya başına 128–512 MiB hedefleyerek küçük dosyaları sıkıştırma (compaction) yoluyla yönetin.
- Hive metastore
- Katalogları kümeler arasında paylaşmak için şemaları ve tablo meta verilerini Dataproc Metastore (yönetilen Apache Hive Metastore) veya Cloud SQL destekli bir metastore'da merkezileştirin.
- Dayanıklılık için GCS'yi işaret eden harici tablolar kullanın; tarama maliyetini sınırlamak için tarihe/saate göre bölümleyin.
- İşler, başlatma ve iş akışları
- spark, pyspark, spark-sql veya hadoop işleri gönderin. Başlatma eylemleri (Initialization actions), küme oluşturma sırasında ek kütüphaneler veya aracılar kurar (örneğin, bağlayıcılar, Python kütüphaneleri).
- İş akışı şablonları (Workflow templates), çok adımlı işlem hatlarını parametrelendirir; her iş akışı için geçici kümeler oluşturup ardından bunları kaldırabilirler. Bu, izolasyonu iyileştirir ve boşta kalma maliyetini azaltır.
- Toplu ETL için geçici (ephemeral) kümeler önerilir; veri ve metastore, kümenin dışında yaşar (GCS, Dataproc Metastore, BigQuery).
- BigQuery entegrasyonu
- Spark BigQuery bağlayıcısı, BigQuery'ye doğrudan okuma/yazma yapar; yüksek veri işleme hacmi (throughput) için BigQuery Storage Read API'yi ve daha düşük gecikmeli, tam olarak bir kez (exactly-once) akış eklemeleri için Write API'yi değerlendirin.
- Tablo bakımı için, yüklemeleri atomik olarak sonlandırmak amacıyla BigQuery'de alt akış (downstream) MERGE/bölüm üzerine yazma (partition overwrite) işlemleri gerçekleştirin.
### Spark Modeli, Performans Ayarlaması ve Güvenilirlik
- API'ler ve yürütme
- RDD'ler: düşük seviyeli, değişmez (immutable), Scala/Java'da tip güvenli; bölümlemeyi (partitioning) ve kalıcılığı (persistence) siz kontrol edersiniz.
- DataFrames/Datasets: ilişkisel, Catalyst ile optimize edilmiş; sorgu optimizasyonu ve kod üretimi nedeniyle ETL için bunları tercih edin.
- Dönüşümler (transformations) tembeldir (lazy) (map, filter, join); eylemler (actions) yürütmeyi tetikler (count, collect, save). Spark, shuffle'larla ayrılmış aşamalardan (stages) oluşan bir DAG oluşturur; görevler (tasks) her bölüm (partition) için çalışır.
- Bölümleme ve shuffle
- Giriş bölümlemesi: tüm çekirdekleri kullanmak için yeterli sayıda bölüm; toplam executor çekirdek sayısının 2-4 katı ile başlayın. RDD'ler için
--conf spark.sql.shuffle.partitions= {total_executor_cores * 2 to 3}
```
ve DataFrames için okuyucu seçenekleri aracılığıyla kontrol edin.
- Shuffle bölümleri: varsayılan 200 değeri genellikle yetersiz veya aşırı kaynak sağlar. Ayarlayın:
--conf spark.sql.shuffle.partitions= {total_executor_cores * 2 to 3}
```
- Geniş dönüşümlerden (wide transforms) sonra bölüm başına ~100–256 MiB hedefleyin; çok küçük olması zamanlayıcı (scheduler) ek yüküne neden olur ve çok büyük olması executor OOM riskini artırır.
- Shuffle, `join`, `groupBy` ve `orderBy` işlemleri için baskın maliyettir. Yeterli executor belleği ve disk olduğundan emin olun; yoğun shuffle içeren kümelerde (cluster) yerel SSD'leri değerlendirin.
- Veri Eğriliği (skew) ve join stratejisi
- Veri eğriliğini (uzun süren görevler, büyük bölüm boyutları) tespit edin. Azaltma yöntemleri:
- Shuffle'ları önlemek için küçük tabloları yayınlayın (broadcast):
--conf spark.sql.autoBroadcastJoinThreshold=64m
```
- Yoğun kullanılan (hot) bölümler için anahtarlara tuzlama (salting) yapın; harita tarafında (map-side) ön toplama uygulayın; filtrelemeyi erken yapın.
- Shuffle sonrası bölümleri birleştirmek ve eğri join'leri yönetmek için Adaptive Query Execution (AQE) özelliğini etkinleştirin:
--conf spark.sql.adaptive.enabled=true
```
- Önbellekleme (caching), denetim noktası oluşturma (checkpointing) ve soy (lineage)
- Yeniden kullanıldığında sık erişilen (hot) ara DataFrame'leri idareli bir şekilde önbelleğe alın; OOM'dan kaçınmak için `MEMORY_AND_DISK` tercih edin.
- Hatalarda yeniden hesaplamayı sınırlamak için uzun soyları (lineage) GCS veya HDFS'e denetim noktası olarak kaydedin.
- Executor'lar ve dinamik ayırma
- Paralellik ve GC ek yükünü dengelemek için executor'ları doğru boyutlandırın:
- Executor başına çekirdek sayısı: Dengeli G/Ç (I/O) ve CPU görevleri için 2–5; daha az çekirdek GC duraklamalarını azaltır.
- Bellek ek yükü: geniş shuffle'lar için `spark.yarn.executor.memoryOverhead` ayarını yapın.
- İş yüküne göre executor'ları ölçeklendirmek için kümelerde harici shuffle hizmeti ile dinamik ayırmayı etkinleştirin:
--conf spark.dynamicAllocation.enabled=true
--conf spark.shuffle.service.enabled=true
--conf spark.dynamicAllocation.minExecutors=0
--conf spark.dynamicAllocation.maxExecutors=200
```
- Toplu ETL için hataya dayanıklılık (fault-tolerance) kalıpları
- Bir kez etkili (idempotent) yazmalar: geçici/hazırlık (temp/staging) yoluna yazın, ardından dizin düzeyinde bir commit ile atomik olarak yükseltin; BigQuery için, hazırlık tablosuna yazın ve
MERGEkullanın:
- Bir kez etkili (idempotent) yazmalar: geçici/hazırlık (temp/staging) yoluna yazın, ardından dizin düzeyinde bir commit ile atomik olarak yükseltin; BigQuery için, hazırlık tablosuna yazın ve
MERGE target t USING staging s
ON t.id = s.id
WHEN MATCHED THEN UPDATE SET ...
WHEN NOT MATCHED THEN INSERT (...)
```
- Artımlı işleme: `ingestion_date` bölümlerinde damga (watermark) tabanlı filtreleme kullanın; yeniden işlemeyi önlemek için GCS'de işlenmiş bir bildirim (processed-manifest) dosyası tutun.
- İşlenemeyen mesajların yönetimi (Dead-letter handling): ayrıştırma/doğrulama hatalarında, hatalı kayıtları tanılama bilgileriyle birlikte bir karantina yoluna/tablosuna yönlendirin. Katı şema zorunluluğu ve yerleşik DLQ'lar için Dataflow'u değerlendirin; Spark ile kayıt başına `try/catch` ve ayrı bir hedef (sink) uygulayın.
### Güvenlik, Gözlemlenebilirlik ve Maliyet
- Kimlik ve erişim
- Kümeleri ve işleri, en az ayrıcalık ilkesine sahip (least-privilege) IAM ile atanmış özel hizmet hesapları altında çalıştırın. Yalnızca gerekli rolleri atayın, örneğin:
- örnek hizmet hesaplarına roles/dataproc.worker
- GCS G/Ç yolları için roles/storage.objectViewer veya objectAdmin
- hedef veri kümelerinde roles/bigquery.dataEditor
- Dataproc Serverless için, erişim kapsamını belirlemek amacıyla iş başına hizmet hesapları kullanın.
- Ağ yalıtımı ve şifreleme
- Bir VPC alt ağında özel IP'li kümeler kullanın, ana (master) kullanıcı arayüzlerini güvenlik duvarıyla kısıtlayın ve genel çıkış (egress) olmadan GCS/BigQuery için Private Google Access'i etkinleştirin.
- Merkezi kontrol için kümeleri Shared VPC projelerine yerleştirin. İsteğe bağlı olarak, küme içi kimlik doğrulama için Dataproc üzerinde Kerberos'u etkinleştirin.
- CMEK ile bekleme durumundaki verileri (at rest) şifreleyin: GCS bucket'ları, Persistent Disk'ler, Dataproc Metastore ve BigQuery üzerinde CMEK'i yapılandırın; aktarım sırasındaki veriler (in transit) için varsayılan olarak TLS kullanın.
- Günlük kaydı, geçmiş ve metrikler
- GCS'e Spark olay günlüklerini etkinleştirin ve History Server'ı dağıtın:
--conf spark.eventLog.enabled=true
--conf spark.eventLog.dir=gs://bucket/spark-events/
```
- Dataproc, sürücü (driver) ve YARN günlüklerini Cloud Logging’e akıtır; saklama/adli analiz (retention/forensics) için havuzlara (sinks) aktarın.
- Cloud Monitoring metrikleriyle izleyin: YARN bekleyen container’lar, CPU, bellek, HDFS sağlığı (kullanılıyorsa), GCS iş hacmi (throughput). Uzun süren aşama (stage) yeniden denemeleri, yürütücü (executor) kaybı ve spekülatif yürütme (speculative execution) ani artışları için uyarılar ayarlayın.
- Hata analizi: yaygın nedenler arasında veri çarpıklığının (skew) neden olduğu yavaş görevler (stragglers), shuffle sırasında yürütücü OOM (bellek yetersizliği), nesne depolama (object-store) commit hataları ve preemptible/spot düğüm kaybı bulunur. Yeniden deneme sayılarını dikkatli bir şekilde artırın; aşırı yeniden denemeler maliyeti ve gecikmeyi artırabilir.
- Maliyet optimizasyonu
- Boşta kalma maliyetinden kaçınmak için geçici (ephemeral) kümeler veya Dataproc Serverless kullanın; kalıcı disk (persistent disk) kullanımını en aza indirmek için verileri GCS’te tutun.
- Yoğun talebi karşılamak için preemptible/spot ikincil çalışanlar (secondary workers) ekleyin; kaybolan düğümlerdeki görevler yeniden denendiği için yeniden hesaplamaya (recomputation) uygun tasarım yapın. Ana (master) düğümleri preemptible düğümlere yerleştirmeyin.
- Makine türlerini doğru boyutlandırın ve kuyruklar boşaldığında kapasiteyi küçültmek için otomatik ölçeklendirme (autoscaling) kullanın. Tarama maliyetini ve CPU kullanımını azaltmak için bölüm budama (partition pruning) ile Parquet/ORC’yi tercih edin.
- Çıktıları birleştirerek (compacting) küçük dosyalardan kaçının; daha az sayıda ve daha büyük dosyalar, meta veri ek yükünü ve iş çalışma süresini azaltır.
- Kısa, periyodik işler için (örneğin, haftalık 30 dakikalık Spark ETL), preemptible çalışanlar veya sunucusuz (serverless) çözümler genellikle en iyi maliyet profilini sunar.
Pratik Problem Senaryosu
Acme Retail, alt sistemlerdeki (downstream) analitikleri besleyen, gecelik Spark ve Hive ETL işlerini çalıştıran 30 düğümlü şirket içi (on-prem) Hadoop kümesini taşıyor. Mevcut işleri minimum değişiklikle yeniden kullanmak, kümeleri tam zamanlı yönetmekten kaçınmak, verileri küme ömrünün ötesinde kalıcı kılmak ve depolama maliyetini düşürmek istiyorlar.
Yaklaşım:
Verileri ve meta verileri yönetilen hizmetlere yerleştirin
- Tüm ham ve işlenmiş verileri, bölümleme (partitioning) kullanarak (örneğin, dt=YYYY-AA-GG) Parquet formatında Cloud Storage’da saklayın.
- Gerekçe: GCS dayanıklı, düşük maliyetlidir ve işlem (compute) ile depolamayı (storage) birbirinden ayırır, böylece geçici (ephemeral) kümeler ve sunucusuz (serverless) işler kalıcı diskler olmadan çalışabilir. Bölümlenmiş Parquet, koşul itme (predicate pushdown) ve verimli taramalar sağlar.
Kataloğu Dataproc Metastore ile merkezileştirin
- Hive metastore’unu Dataproc Metastore’a taşıyın. GCS yollarına referans veren harici (external) Hive tabloları oluşturun ve mevcut şema/bölümleme mantığını koruyun.
- Gerekçe: Yönetilen bir metastore, birden çok geçici kümenin ve sunucusuz işin, yüksek erişilebilirlikli (HA) bir MySQL/PostgreSQL örneği çalıştırmadan tablo tanımlarını paylaşmasına olanak tanır.
Toplu ETL için geçici Dataproc kümeleri ve orkestrasyon için iş akışı şablonları kullanın
- Gerekli imajla (örneğin, 2.1-debian11) bir küme oluşturan, Spark işlerini (spark-sql ve pyspark) çalıştıran ve tamamlandığında kümeyi silen bir iş akışı şablonu (workflow template) tanımlayın. Özel kütüphaneleri yüklemek için başlatma eylemleri (initialization actions) ekleyin.
- Gerekçe: Geçici kümeler boşta kalma maliyetini ortadan kaldırır ve iş bağımlılıklarını yalıtır. İş akışı şablonları, tekrarlanabilirlik ve parametreleştirme (tarihler, girdi yolları) sağlar.
Otomatik ölçeklendirmeyi ve preemptible çalışanları etkinleştirin
- Küçük bir çekirdek çalışan (core worker) grubu ve daha büyük bir preemptible ikincil çalışan (secondary worker) havuzu içeren bir otomatik ölçeklendirme politikası (autoscaling policy) ekleyin; iş bittikten sonra hızla küçülmek için bekleme sürelerini (cooldowns) ayarlayın.
- Gerekçe: Çekirdek çalışanlar küme kararlılığını korur; preemptible çalışanlar, shuffle ve geniş dönüşümleri (wide transformations) daha düşük maliyetle karşılar. Spark/YARN yeniden denemeleri, preemption (görevin geri alınması) durumunda kaybolan görevleri yönetir.
Spark BigQuery bağlayıcısı aracılığıyla BigQuery ile entegre edin
- Boyut/olgu (dimension/fact) yüklemeleri için, Spark sonuçlarını hazırlık (staging) BigQuery tablolarına yazın, ardından hedefleri atomik olarak güncellemek için MERGE ifadelerini çalıştırın. Doğrudan üzerine yazmanın (overwrite) güvenli olduğu durumlarda, bölümlenmiş tabloları bölüm üzerine yazma modu (partition overwrite mode) kullanarak yazın.
- Gerekçe: BigQuery, büyük ölçekte analitik ve iş zekası (BI) hizmeti sunar; hazırlık+MERGE, toplu Spark’tan işlemsel benzeri (transactional-like) upsert’ler üreterek alt sistemlerdeki tutarsızlığı azaltır.
Performans ve güvenilirlik için Spark’ı ayarlayın
- Shuffle bölümlerini (partitions) yürütücü çekirdeklerine (executor cores) göre ayarlayın ve AQE’yi etkinleştirin:
--conf spark.sql.shuffle.partitions=600
--conf spark.sql.adaptive.enabled=true
```
- Küçük boyutlar (dimensions) için broadcast join'leri kullanın ve kararlılık için uzun soy ağaçlarını (long lineages) GCS'e checkpoint alın.
- Gerekçe: Doğru bölümleme, veri çarpıklığını (skew) ve zamanlayıcı (scheduler) ek yükünü azaltır; AQE, çalışma zamanında veri profillerine uyum sağlar; checkpointing, hatalardan sonra yeniden hesaplamayı sınırlar.
7) Güvenliği ve ağı sağlamlaştırın
- Kümeleri, yalnızca GCS yolları, metastore ve BigQuery veri kümeleri için gereken rolleri veren özel hizmet hesaplarıyla çalıştırın. Kısıtlı bir alt ağda Private Google Access ile özel IP'li kümeler oluşturun ve kullanıcı arayüzü erişimini güvenlik duvarı kurallarıyla sınırlayın.
- Gerekçe: En az ayrıcalık ilkesi ve ağ yalıtımı, saldırı yüzeyini azaltır; özel kontrol düzlemi çıkışı (control-plane egress), genel kullanıma maruz kalmayı önler.
8) Günlük kaydı, geçmiş ve uyarıları yapılandırın
- GCS'e Spark olay günlüklerini etkinleştirin ve History Server'ı dağıtın; sürücü/YARN günlüklerini saklama (retention) ile Cloud Logging'e yönlendirin. Uzun süre bekleyen container'lar, tekrarlanan görev hataları veya aşırı iş süresi için Monitoring uyarıları ekleyin.
- Gerekçe: Merkezi günlükler kök neden analizini destekler; proaktif uyarılar veri çarpıklığını, OOM'ları veya düşen G/Ç performansını erken tespit eder.
9) Anlık (ad hoc) ve esnek ani artışlar için Dataproc Serverless ile seçici olarak modernleştirin
- Düzensiz veya keşif amaçlı Spark SQL iş yüklerini Dataproc Serverless'a taşıyın; gecelik işlem hatlarını (pipelines) sunucusuz ortamda tamamen doğrulanana kadar geçici kümelerde tutun.
- Gerekçe: Sunucusuz (Serverless) çözüm, küme operasyonlarını ortadan kaldırır ve otomatik olarak ölçeklenir, bu da onu öngörülemeyen yükler için ideal kılar; mevcut iş akışları minimum kod değişikliği ile devam eder.
10) Nesne depolama committer'larını ve küçük dosya yönetimini doğrulayın
- FileOutputCommitter algoritmasını v2 olarak ayarlayın ve yazma işlemlerinden önce repartition/coalesce aracılığıyla çıktıları dosya başına 256–512 MiB olacak şekilde birleştirin.
- Gerekçe: Nesne depoları atomik yeniden adlandırma (atomic rename) özelliğinden yoksundur; optimize edilmiş committer'lar kopyalama/yeniden adlandırma ek yükünü azaltır. Birleştirme (Compaction), performans ve maliyet açısından küçük dosyalar sorununu hafifletir.
Bu tasarım, mevcut Spark ve Hive işlerini minimum yeniden düzenleme (refactoring) ile yeniden kullanır, GCS'te veri dayanıklılığını sağlar, şemaları merkezileştirir, güvenlik etki alanını (blast radius) sınırlar, sağlam bir gözlemlenebilirlik sunar ve geçici kümeler, otomatik ölçeklendirme, preemptible kapasite ve sunucusuz yürütmenin hedeflenmiş kullanımı yoluyla maliyeti optimize eder.
---
← [Mesajlaşma](/tr/posts/pde-messaging-ingestion/) · [Tüm alanlar](/tr/posts/google-pde-study-guide/) · [Veri Alımı](/tr/posts/pde-ingestion-migration/) →
**[Bu soruları çözün →](/tr/kb/google/)** · **[ExamRoll.io'da süreli pratik →](https://www.examroll.io/?utm_source=guide&utm_medium=referral&utm_campaign=PDE)**
Pass the whole exam — not just this question
You found this answer. Get every verified question and explanation in one place, and save hours of prep. Free to start.
Sınavınızı geçin →