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

    --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.

    --conf spark.sql.shuffle.partitions= {total_executor_cores * 2 to 3}
    ```

  - Geniş dönüşümlerden (wide transforms) sonra bölüm başına ~100256 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
  ```
    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/
```

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:

  1. 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.
  2. 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.
  3. 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.
  4. 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.
  5. 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.
  6. 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  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ı  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  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  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 256512 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 →

Google'a göz atın →

Related guides

Hepsi bir arada erişim

Tek abonelik. Her sınav.

Her plan, sınırsız cevap aramayı, pratik testlerini, AI açıklamalarını ve tam kaynak kütüphanesini — 20'den fazla dilde — açar.

Aylık
24.87
Just €0.83/day
Her şey dahil:
  • Sınırsız cevap arama
  • Sınırsız pratik testi
  • AI destekli açıklamalar
  • Tam kaynak kütüphanesi
  • 20+ dil
  • Haftalık içerik güncellemeleri
  • Ödüller ve yönlendirmeler
  • Öncelikli destek
Ücretsiz denemeyi başlat

Kredi kartı gerekmez*

En iyi değer
12 ay
179.87
Just €0.49/daySave 40%
Her şey dahil:
  • Sınırsız cevap arama
  • Sınırsız pratik testi
  • AI destekli açıklamalar
  • Tam kaynak kütüphanesi
  • 20+ dil
  • Haftalık içerik güncellemeleri
  • Ödüller ve yönlendirmeler
  • Öncelikli destek
Ücretsiz denemeyi başlat

Kredi kartı gerekmez*

✓ Ücretsiz plan dahil · ✓ İstediğiniz zaman iptal edin · ✓ Tüm planlar tam ürünü açar