Google PDE: Dataflow ve Apache Beam ile Akış İş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’da akış işleme, Dataflow çalıştırıcısı tarafından yürütülen Apache Beam’in birleşik programlama modelini merkez alır. Beam, kodunuzu paralellik, otomatik ölçeklendirme ve hataya dayanıklılık gibi yürütme ayrıntılarından ayıran mantıksal bir soyutlama (PCollections üzerindeki dönüşümlerden oluşan işlem hatları) sağlar. Akış işlemede doğruluk; zaman semantiğine (olay zamanı ve işleme zamanı karşılaştırması), pencerelere (sabit, kayan, oturum, genel), filigranlara (watermarks), tetikleyicilere (triggers) ve geç gelen verilerin işlenmesine bağlıdır. Dataflow’da operasyonel mükemmellik; doğru çalışan boyutlandırması, otomatik ölçeklendirme politikası, akış motoru (streaming engine), karıştırma (shuffle) seçenekleri, bir kez etkili (idempotent) hedef (sink) tasarımı, işlenemeyen mesajların yönetimi (dead-letter handling) ve sağlam bir gözlemlenebilirlik gerektirir.

Apache Beam modeli ve zaman semantiği

Hata modları ve ödünleşimler:

Akış tabanlı iş yükleri için Dataflow’u işletme

Dağıtım, şablonlar ve yükseltme stratejileri

undefined

undefined

Pratik Problem Senaryosu

NovaTrack Inc., 50.000 sıcaklık sensöründen küresel IoT telemetrisi toplamakta ve dakika düzeyinde birleştirilmiş veriler (aggregates) sunmalı, ham verileri kalıcı olarak saklamalı ve gerçek zamanlı bir gösterge panosu (dashboard) sağlamalıdır. Ara sıra hatalı biçimlendirilmiş mesajlar ve sırası bozuk teslimat beklenmektedir. Çözüm, otomatik olarak ölçeklenmeli, hatalı kayıtları inceleme için yüzeye çıkarmalı ve sıfır kesintiyle yükseltmeleri desteklemelidir.

Yaklaşım:

  1. Veri alımı (Ingestion) ve zaman semantiği

    • deviceId ve eventTs (RFC3339) özniteliklerine (attributes) sahip bölgesel bir Pub/Sub konusu (topic) ve bölge başına yayıncılar (publishers) oluşturun. Mümkün olduğunda deviceId’ye göre sıralama anahtarlarını (ordering keys) etkinleştirin.
    • Gerekçe: Pub/Sub, en az bir kez teslimat (at-least-once delivery) ile dayanıklı, elastik bir giriş (ingress) sağlar. Olay zaman damgalarını en uç noktada (edge) eklemek, gerçek olay zamanını korur; cihaz başına sıralama, merkezi darboğazlar olmadan cihaz içi yeniden sıralamayı azaltır.
  2. Olay zamanı pencereleri (event-time windows) ile Dataflow akış iş hattı

    • PubSubIO aracılığıyla özel bir abonelikten (subscription) okuyun, eventTs’yi Beam zaman damgası olarak çıkarın, eksikse publishTime’a geri dönün.
    • 30 saniyede bir erken tetikleyici (early trigger) ve her geç öğede geç tetiklemeler (late firings) ile 1 dakikalık FixedWindows uygulayın; izin verilen gecikmeyi (allowed lateness) 10 dakika olarak ayarlayın ve panelleri biriktirin (accumulating panes).
    • Gerekçe: Olay zamanı pencereleri, doğru dakika agregalarını sağlar; erken tetiklemeler, gösterge panosunu dakika altı tazelikle besler; geç tetiklemeler, gecikmiş veriler geldikçe agregaları düzeltir. Gecikme sınırı, durum (state) boyutunu ve maliyeti sınırlar.
  3. Doğrulama, zenginleştirme ve işlenemeyen mesaj (dead-letter) yönlendirmesi

    • JSON’u ayrıştıran (parse), şemayı ve aralıkları doğrulayan ve iş başlangıcında BigQuery’den yüklenen bir yan girdi (side input) aracılığıyla küçük statik referans verileriyle zenginleştiren bir ParDo uygulayın.
    • Geçerli kayıtları ana çıktıya ve hataları, payload, hata, deviceId ve ayrıştırma zaman damgasını içeren bir işlenemeyen mesaj (dead-letter) PCollection’ına göndermek için TupleTags kullanın; DLQ’yu bölümlenmiş (partitioned) bir BigQuery tablosuna yazın.
    • Gerekçe: Yan girdiler (side inputs), düşük gecikme için referans verilerini bellekte tutar. İşlenemeyen mesajların yakalanması, ana akışı engellemeden hatalı satırların incelenmesine ve hedeflenmiş olarak yeniden işlenmesine olanak tanır.
  4. Birleştirme (Aggregation) ve hot-key (sık erişilen anahtar) azaltma

    • deviceId’ye göre anahtarlayın ve CombineFns ile dakika başına ort/min/maks hesaplayın. En iyi N bölgesel metrikler için, hot-key’leri önlemek amacıyla region#N’e göre parçalara ayırın (shard), ardından yeniden birleştirin (re-aggregate).
    • Gerekçe: Combiner’lar, shuffle hacmini ve maliyetini en aza indirir; anahtar parçalama (key sharding), bölgesel fan-in sırasında tek anahtar darboğazlarını önler.
  5. Havuzlar (Sinks) ve tam olarak bir kez (exactly-once) etkileri

    • Ham doğrulanmış olayları ve dakika agregalarını, Storage Write API ile BigQueryIO kullanarak BigQuery’ye yazın. Herhangi bir özel yeniden denemede idemopotans (idempotency) için deviceId + eventTs’ye dayalı kararlı bir ekleme kimliği (insert id) ayarlayın.
    • Gerekçe: Storage Write API, bir akış içinde tam olarak bir kez (exactly-once) semantiği ile yüksek verimli (high-throughput), düşük gecikmeli veri alımı sağlar. Kararlı kimlikler, yeniden oynatmalar (replays) meydana gelirse aşağı akışta (downstream) yinelenenleri önlemeyi (dedup) sağlar.
  6. Gösterge panosu tutarlılık stratejisi

    • Gösterge panosu, filigrana (watermark) göre 2 dakikalık bir geriye bakışla veya akış verileri için gözlemlenen kullanılabilirlik gecikmesinin 2 katı sabit bir gecikmeyle bölümlenmiş agrega tablolarını sorgular.
    • Gerekçe: BigQuery akış görünürlüğü nihayetinde tutarlıdır (eventually consistent); okumaları hafifçe ertelemek, neredeyse gerçek zamanlı davranışı korurken devam eden satırların kaçırılmasını önler.
  7. Operasyonlar: otomatik ölçeklendirme ve akış motoru (streaming engine)

    • Streaming Engine’i etkinleştirin; beklenen zirveye göre (ör. ortalamanın 3 katı) maxWorkers’ı ayarlayın, CPU-yoğun ayrıştırma ve şifreleme için boyutlandırılmış bir makine türü seçin ve geçici shuffle’ı barındırmak için önyükleme diskini (boot disk) artırın.
    • Filigran gecikmesini (watermark lag), birikim saniyelerini (backlog seconds), CPU’yu ve adım başına verimliliği (throughput) izleyin; sürekli gecikme ve DLQ oranı artışlarında uyarı ayarlayın.
    • Gerekçe: Streaming Engine, esneklik ve daha basit yükseltmeler için durumu (state)/shuffle’ı dışsallaştırır; doğru boyutlandırma ve izleme, sessiz SLO ihlallerini önler.
  8. Flex Templates ile dağıtım ve yükseltmeler

    • İş hattını; girdi aboneliği, çıktı tabloları, DLQ tablosu, maxWorkers ve bölge gibi parametrelerle bir Flex Template olarak paketleyin. Uyumsuz bir değişiklik için, yeni iş hattını aynı konuyu (topic) hedefleyen yeni bir abonelikle başlatın, çıktıları doğrulayın, ardından eski işi boşaltın (drain). İsteğe bağlı olarak bir Pub/Sub anlık görüntüsü (snapshot) oluşturun ve boşluk olmamasını garanti etmek için yeni aboneliği anlık görüntüye yönlendirin (seek).
    • Gerekçe: Flex Templates, tekrarlanabilir, parametreli dağıtımları mümkün kılar. Boşaltma (drain) ile yapılan doğrulanmış bir mavi/yeşil geçiş (blue/green cutover), sıfır veri kaybı ve minimum kesinti süresi sağlar.
  9. Yeniden işleme ve toplu (batch) geriye dönük doldurma (backfills)

    • Ham olayların sıkıştırılmış Avro dosyalarını bir yan çıktı (side output) aracılığıyla Cloud Storage’da saklayın; modeller veya şemalar değiştiğinde BigQuery’ye geriye dönük doldurma veya yeniden işleme yapmak için bir toplu (batch) Dataflow iş hattı çalıştırın.
    • Gerekçe: Dayanıklı ham arşivler, sıcak yolu (hot path) etkilemeden tekrarlanabilirliği ve şema evrimini destekler.

Bu tasarım, küresel ölçekte sırası bozuk ve geç verileri işlerken sınırlı maliyetle doğru, düşük gecikmeli agregalar, net hata izolasyonu, güçlü gözlemlenebilirlik ve güvenli yükseltme yolları sağlar.


BigQuery Analitiği ve Veri Ambarı Mühendisliği · Tüm alanlar · Mesajlaşma

Bu soruları çözün → · ExamRoll.io’da süreli pratik →

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