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
İşlem hatları (Pipelines), dönüşümler (transforms), PCollections, çalıştırıcılar (runners):
- Bir Beam işlem hattı, PCollections (sınırlı veya sınırsız) üzerinde yönlendirilmiş döngüsel olmayan bir PTransforms grafiği uygular.
- Çalıştırıcılar (Dataflow, Spark, Flink, Direct) işlem hattını yürütür; Dataflow yönetilen otomatik ölçeklendirme, kontrol noktası oluşturma (checkpointing) ve operasyonel görünürlük sağlar.
- Dönüşümler arasında öğe bazında (ParDo), gruplama ve birleştirme (GroupByKey, Combine), birleştirmeler (CoGroupByKey) ve G/Ç işlemleri (PubSubIO, BigQueryIO, FileIO) bulunur.
Pencereler (Windows):
- Sabit pencereler (Fixed windows): periyodik toplamalar için birbiriyle kesişmeyen zaman dilimleri (ör. 1 dakikalık takla atan pencereler - tumbling windows).
- Kayan pencereler (Sliding windows): düzgün hareketli metrikler için birbiriyle kesişen pencereler (ör. her 1 dakikada bir kayan 5 dakikalık pencereler).
- Oturum pencereleri (Session windows): belirli bir süre etkinlik olmadığında kapanan, kullanıcı oturumları veya cihaz etkinlik patlamaları için ideal olan dinamik pencereler.
- Genel pencere (Global window): tüm sınırsız akışın varsayılan, penceresiz görünümü; genellikle periyodik somutlaştırma (materialization) için tetikleyicilerle birlikte kullanılır.
Olay zamanı ve işleme zamanı karşılaştırması:
- Olay zamanı (Event time): olayın kaynakta meydana geldiği zaman; değişken aktarım gecikmelerine rağmen mantıksal olarak tutarlı toplamalar yapılmasını sağlar.
- İşleme zamanı (Processing time): olayın işlem hattı tarafından gözlemlendiği zaman; operasyonel tetikleyiciler için kullanışlıdır ancak anlamsal doğruluk için uygun değildir.
Filigranlar (Watermarks):
- Bir filigran, olay zamanı bütünlüğünü tahmin eder (çalıştırıcının T zamanına kadar olan tüm olayları gördüğüne dair tahmini).
- Filigranlar, geri basınç (backpressure) veya kaynak gecikmeleri altında düzensiz ilerleyebilir veya duraksayabilir; geç gelen veri, zaman damgası < filigran olan her şeydir.
Tetikleyiciler (Triggers) ve gecikme:
- Varsayılan: Filigran pencere sonunu geçtiğinde ateşlenen AfterWatermark tetikleyicisi; izin verilen gecikme = 0 olduğunda, geç gelen veriler atılır.
- Erken ateşlemeler (işleme zamanı veya sayı tabanlı), düşük gecikmeli ön sonuçlar verir.
- Geç ateşlemeler, geç gelen veriler ulaştığında düzeltmelere olanak tanır; biriktirme modu (accumulation mode), bölmelerin (panes) sonuçları biriktirip biriktirmeyeceğini veya önceki çıktıyı atıp atmayacağını yönetir.
- İzin verilen gecikmeyi, iş toleransına ve depolama/hesaplama ödünleşimlerine göre seçin; daha fazla gecikme, durum (state) saklama süresini ve maliyeti artırır.
Durum bilgili işleme (Stateful processing), zamanlayıcılar (timers), oturumlara ayırma (sessionization), tekilleştirme (deduplication):
- Durum bilgili DoFns, anahtar başına durumu (ör. son görülen olay, anlık toplamalar) tutar ve durumu yaymak veya temizlemek için zamanlayıcılar ayarlar.
- Oturumlara ayırma, SessionWindows aracılığıyla doğal olarak ifade edilir; özel mantık için anahtarlı durum (keyed state) ve işleme/olay zamanı zamanlayıcıları kullanın.
- Tekilleştirme: her olay için kararlı bir kimlik (ID) kullanın ve pencere başına Distinct/Combine veya anahtar başına durum (ör. Bloom filtresi veya TTL’li bir küme) kullanın. Bellek ve yanlış pozitifler ile katı doğruluk arasında bir ödünleşim yapın.
Hata modları ve ödünleşimler:
- İş metrikleri için işleme zamanı pencerelerini kullanmak, ani yükselişler veya yeniden denemeler sırasında kaymalara neden olur; olay zamanı pencerelerini tercih edin.
- Sık erken tetikleyicilere sahip çok küçük pencereler, aşırı bölme (pane) yayımına ve hedef (sink) yazma çoğaltılmasına neden olur.
- Sınırsız izin verilen gecikme, durumu (state) şişirebilir; her zaman durum TTL’sini sınırlayın ve etkin olmayan anahtarları temizlemek için zamanlayıcılar ayarlayın.
Akış tabanlı iş yükleri için Dataflow’u işletme
Worker boyutlandırma ve otomatik ölçeklendirme:
- Yatay otomatik ölçeklendirme; birikime, watermark gecikmesine, CPU’ya ve iş hacmine (throughput) göre worker ekler/kaldırır; ani artışları karşılamak için makul bir maxWorkers değeri belirleyin.
- Darboğazlara göre makine türleri seçin: CPU-yoğun (daha fazla vCPU), bellek-yoğun (yüksek bellekli türler), ağ-yoğun (daha büyük VM’ler shuffle ek yükünü azaltır).
- Ağır shuffle işlemleri veya dosya tabanlı sink’ler için önyükleme diskini artırın. Sistem gecikmesini (system lag) ve birikim saniyelerini (backlog seconds) izleyin.
Streaming Engine ve shuffle:
- Streaming Engine, state (durum) ve shuffle işlemlerini hizmetin arka ucuna taşıyarak esnekliği artırır, worker’lar üzerindeki bellek baskısını azaltır ve daha hızlı güncellemelere olanak tanır.
- Toplu iş (batch) ağırlıklı aşamalar veya çok büyük anahtar gruplamaları için, shuffle G/Ç (I/O) yükünü worker’lardan almak üzere Dataflow Shuffle’ı kullanın. Her ikisi de aşırı yüklenmiş worker (hot-worker) hatalarını ve disk thrashing’i azaltır.
Geri basınç (backpressure), hot key’ler ve veri dengesizliği (skew):
- Dataflow, geri basıncı dinamik iş yeniden dengeleme yoluyla yönetir; yine de, uygun olduğunda kaynak akış kontrolünü (örneğin, Pub/Sub’da bekleyen mesaj/bayt sayısı) ayarlayın.
- Hot key’ler (sık erişilen anahtarlar, ör. popüler kimlikler) gecikmelilere (straggler) neden olur. Anahtar parçalama (key sharding, ör. key#N), kısmi ön toplama ve ardından yeniden anahtarlama veya taslak tabanlı (sketch-based) yaklaşımlarla bu durumu hafifletin.
- Aykırı kayıtlardan (çok büyük veri yükleri) veya anlık yoğunluk yaratan yayıncılardan (bursty publishers) kaynaklanan veri dengesizliği, yayıncı başına bölümleme, toplu işleme (batching) veya sıkıştırma gerektirebilir.
Pub/Sub entegrasyonu:
- Veri alımı (ingestion) için Pub/Sub topic’lerini kullanın; meta veriler (ör. deviceId, olay zaman damgası) için mesaj niteliklerini (message attributes) etkinleştirin.
- PubSubIO ile veri alın; olay zaman damgalarını niteliklerden veya veri yükünden (payload) çıkarın, aksi takdirde yayınlanma zamanını (publish time) kullanın.
- Sıralama anahtarları (ordering keys) anahtar başına sıralama sağlar; en az bir kez teslimat (at-least-once delivery) garantisi nedeniyle Dataflow’un alt sistemlerde (downstream) yine de idempotent davranışa ihtiyacı vardır.
BigQuery’ye akış (streaming) desenleri:
- Yüksek iş hacmi, düşük gecikme süresi ve bir akış içinde “tam olarak bir kez” (exactly-once) semantiği için akış ofsetleri (stream offsets) ve otomatik yeniden denemeler aracılığıyla Storage Write API ile BigQueryIO’yu tercih edin.
- Düşük oranlı basit pipeline’lar için akış tabanlı eklemeler (streaming inserts) kabul edilebilirdir; istemci yeniden denemelerini tekilleştirmek için insertId’yi ayarlayın.
- Akış tamponları (streaming buffers) üzerindeki sorgular nihai tutarlıdır (eventually consistent); zamana duyarlı analizler için, bir tampon gecikmesinden sonra (ör. gözlemlenen kullanılabilirlik gecikmesinin ~2 katı kadar bekle) sorgulayın veya mikro-toplu iş pencereleri ve Storage Write API’nin kararlı modu (committed mode) aracılığıyla veriyi somutlaştırın.
Tam olarak bir kez (exactly-once) etkileri, idempotency, yeniden oynatma (replay) ve sink’ler:
- Beam, en az bir kez işlemeyi (at-least-once processing) garanti eder; “tam olarak bir kez” (exactly-once) etkisi, sink tarafında idempotent (tekrarlanabilir) yazmalar, işlemler (transactions) veya tekilleştirme anahtarları kullanılarak sağlanmalıdır.
- BigQuery: Bir akış içinde tam olarak bir kez (exactly-once) garantisi için Storage Write API’nin varsayılan akışlarını (default streams) veya kararlı akışlarını (committed streams) kullanın; akış tabanlı eklemelerle (streaming inserts) kararlı bir insertId ayarlayın.
- Dosyalar: Geçici dosyaları benzersiz adlarla yazın, pencere tamamlandığında sonlandırın ve atomik yeniden adlandırmalar sağlayın; kısmi kopyaları önlemek için üzerine yazmaktan kaçının.
- Harici veritabanları: Kararlı bir kimliğe (id) göre anahtarlanmış upsert (ekle veya güncelle) işlemleri kullanın veya tekilleştirme pencereleri uygulayın.
- Yeniden oynatma (replay) için tasarım yapın: Deterministik dönüşümleri koruyun; sink’lerin yeniden deneme durumunda tekilleştirme yaptığından emin olun.
İşlenemeyen mesajların yönetimi (dead-letter handling), hata yönlendirme ve gözlemlenebilirlik:
- Riskli ayrıştırma/zenginleştirme işlemlerini ParDo içinde try/catch bloğuna alın ve hataları bir TupleTag aracılığıyla işlenemeyen mesajlar için ayrılmış bir PCollection’a (dead-letter PCollection) gönderin; veri yükünü (payload), hata kodunu ve bağlamı dahil edin.
- DLQ’ları (Dead-Letter Queue) analiz için BigQuery’ye veya Cloud Storage’a yönlendirin; yeniden işleme için ayrı bir Pub/Sub topic’i kullanmayı düşünün.
- Gözlemlenebilirlik: Dataflow iş metriklerini (watermark gecikmesi, sistem gecikmesi, iş hacmi), özel sayaçları, dağılım metriklerini ve Cloud Logging’deki adım başına logları kullanın. Cloud Monitoring’de gecikme ve hata oranları için uyarılar oluşturun. İstisnaları (exceptions) toplamak için Error Reporting’i kullanın.
Performans ayarlama desenleri:
- Verimli okuma: BigQuery kaynakları için, yalnızca gerekli alanları ve filtreleri seçen Storage Read API’yi veya sorgu tabanlı okumaları tercih edin.
- Combiner’ları erken kullanma (Combine lifting): GroupByKey’den önce shuffle hacmini azaltmak için combiner’ları kullanın.
- Yan girdiler (Side inputs): Küçük referans verilerini bellekte önbelleğe alın; fanout (yayılım) ve güncelleme sıklığını gözlemleyin.
- Serileştirme: Kompakt şemalar (Avro/Proto) kullanın ve sık kullanılan yollarda (hot paths) aşırı JSON ayrıştırmasından kaçının.
Dağıtım, şablonlar ve yükseltme stratejileri
Flex Templates:
- İş hatlarını (pipeline), tekrarlanabilir dağıtımlar için konteynerize edilmiş, parametreli şablonlar halinde paketleyin. Flex Templates, özel bağımlılıkları, GPU imajlarını ve ortam izolasyonunu destekler.
- Ortama özgü dağıtımları etkinleştirmek için çalışma zamanı parametrelerini (ör. girdi aboneliği, çıktı tablosu, işlenemeyen mesaj havuzu, maxWorkers) dışsallaştırın.
İş hattı güncellemeleri ve uyumluluk:
- Dataflow, dönüşüm (transform) adları, durum (state) özellikleri ve çıktı türleri uyumlu kaldığı sürece birçok akış (streaming) iş hattı için yerinde güncellemeyi (in-place update) destekler. Kararlı PTransform adları kullanın.
- Uyumsuz grafik (graph) veya durum (state) değişiklikleri için kontrollü bir geçiş (cutover) yapın: yeni işi başlatın, ardından devam eden işi bitirmesi ve yeni öğeleri okumayı durdurması için eski işi boşaltın (drain).
Boşaltma (Draining) ve anlık görüntüler (snapshots):
- Boşaltma (drain), işlemeyi düzgün bir şekilde tamamlar, kalan çıktıyı yazar ve sonlanır; veri boşluklarını önlemek için Pub/Sub saklama (retention) süresi veya anlık görüntüler (snapshots) ile koordine edin.
- Sürekliliği sağlamak için bir Pub/Sub anlık görüntüsü (snapshot) oluşturabilir, yeni iş hattını anlık görüntüye veya uygun bir zaman damgasına gidecek şekilde başlatabilir, çıktıyı doğrulayabilir ve ardından eski işi boşaltabilirsiniz (drain).
Yapılandırma örnekleri:
- Erken/geç tetikleyiciler ve biriktirme (accumulation) ile örnek pencereleme (windowing):
undefined
- Storage Write API ile örnek BigQueryIO:
undefined
- Sık karşılaşılan tuzaklar:
- Akış modunda, pencereli yazmalar (windowed writes) olmadan dosya tabanlı havuzlara (sinks) yazmak, sonlandırmayı (finalization) durdurabilir; pencereli yazmaları ve tetikleyicileri etkinleştirin.
- Sınırsız Büyüme: durumu (state) veya izin verilen gecikmeyi (allowed lateness) sınırlamayı unutmak, bellek sızıntılarına ve ölçeklendirme hatalarına neden olabilir.
- Eksik zaman damgaları: olay zaman damgalarını atamamak, iş hattının varsayılan olarak işleme zamanını (processing time) kullanmasına ve değişken gecikmeler altında doğruluğunu kaybetmesine neden olur.
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:
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.
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.
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.
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.
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.
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.
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.
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.
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 →