Google PDE: Mesajlaşma, Olay Alımı ve Gerçek Zamanlı Hizmetler — Ç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’daki mesajlaşma, olay alımı (event ingestion) ve gerçek zamanlı hizmetler; birbirinden bağımsız (decoupled) ve dayanıklı (durable) aktarım için Cloud Pub/Sub ve Eventarc, durum bilgisi olan (stateful) akış işleme için Dataflow ve BigQuery, Cloud Storage gibi havuzlar (sinks) ile operasyonel veritabanları etrafında şekillenir. En az bir kez teslimat (at-least-once delivery), etkisiz (idempotent) tüketim ve gözlemlenebilirlik için tasarım yapmak; hata, geri basınç (backpressure) ve şema evrimi durumlarında doğruluğu korurken elastik olarak ölçeklenen dayanıklı sistemler sağlar.
Pub/Sub ile Temel Mesajlaşma
- Konular (Topics) ve abonelikler (subscriptions)
- Yayıncılar (Publishers) bir konuya mesaj gönderir; aboneler (subscribers) abonelikler aracılığıyla bağlanır (birden fazla abone aynı mesajları bağımsız olarak tüketebilir).
- Abonelik türleri:
- Pull: İstemciler mesajları açıkça çeker; en yüksek verim (throughput) ve daha az gidiş-dönüş (round trip) için streaming pull kullanın.
- Push: Pub/Sub, HTTPS üzerinden teslimat yapar; uç noktanızın (endpoint) onaylamak (acknowledge) için 2xx döndürmesi gerekir.
- BigQuery’ye Aktarma (Export to BigQuery): Bir BigQuery aboneliği, mesajları kod olmadan bir BigQuery tablosuna teslim eder; yüklerin (payloads) beyan edilen şemayla eşleştiği ve analitik için düşük gecikmeli alımın (ingestion) gerekli olduğu durumlarda en iyisidir.
- Sıralama anahtarları (Ordering keys)
- Sıralama anahtarı başına sıralı teslimat almak için konu ve abonelikte mesaj sıralamasını etkinleştirin. Anahtar başına verim serileştirilir: anahtar başına uçuş halindeki (in-flight) bir mesaj sonrakileri engelleyebilir; ölçeklenmek için çok sayıda anahtar kullanın (örneğin, hash(device_id)).
- Yelpazeleme (Fan-out) ve yeniden oynatma (replay)
- İş yüklerini ve saklama (retention) sürelerini izole etmek için farklı tüketiciler için ayrı abonelikler oluşturun.
- Kurtarma (recovery) ve geçmişe dönük doldurma (backfills) için bir zaman damgasından veya anlık görüntüden (snapshot) yeniden oynatmak için seek veya snapshot kullanın.
Artıları ve eksileri:
- Sıralama, paralelliği ve anahtar başına verimi azaltır; kesinlikle gerekli olmadıkça sıralamayı devre dışı bırakın.
- Push, istemci kodunu basitleştirir ancak HTTP uç noktası ölçeklendirme, güvenlik ve geri çekilme (backoff) endişeleri doğurur; pull, yüksek verimde daha fazla kontrol ve kararlılık sağlar.
Teslimat Semantiği, Onaylama, Saklama ve İşlenemeyen Mesajlar
- Onaylama (Acknowledgment) ve son tarihler (deadlines)
- En az bir kez teslimat (At-least-once delivery): yinelenen mesajlar (duplicates) oluşabilir.
- Her teslimatın bir onay son tarihi (ack deadline) vardır (varsayılan 10 saniye). Uzun süren işleri işlerken süreyi uzatın (ModifyAckDeadline); son tarihten önce onaylanmaması (ack), yinelenen push teslimatlarının en yaygın nedenidir.
- Olumsuz onay (Nack) veya son tarihin dolması, mesajı yeniden teslimat için uygun hale getirir.
- Saklama (Retention)
- Onaylanmamış mesajlar, aboneliğin onay son tarihi boyunca saklanır ve yeniden denenir; onaylanmış mesajlar, yeniden oynatma için konunun mesaj saklama süresine kadar saklanabilir. Saklama süresini, maksimum kesinti sürenizi ve kurtarma sürenizi kapsayacak şekilde yapılandırın.
- Yeniden denemeler (Retries)
- Pull: yeniden teslimat, onay son tarihi dolduktan sonra gerçekleşir; akış kontrolü (flow control) limitleriyle eşzamanlılığı kontrol edin.
- Push: üssel geri çekilme (exponential backoff); yalnızca HTTP 2xx başarılı sayılır. 3xx/4xx/5xx yeniden denemeleri tetikler. Tekrarları tolere etmek için etkisiz (idempotent) işleyiciler (handlers) uygulayın.
- İşlenemeyen mesaj konuları (Dead-letter topics - DLT’ler)
- Zehirli (poison) mesajları karantinaya almak için bir DL konusu ve abonelik başına maksimum teslimat denemesi sayısı yapılandırın.
- DLQ hacmini izleyin; triyaj iş akışları oluşturun ve düzeltmeden sonra ana konuya yeniden yayınlayın.
Örnek:
undefined
Teslimat semantiği özeti:
- Pub/Sub: en az bir kez (at-least-once), etkinleştirilmişse bir sıralama anahtarı içinde en iyi çaba (best-effort) sıralama.
- Havuzlar (Sinks): BigQuery ekleme API’leri yinelenenleri azaltma (insertId veya Storage Write API akış ofsetleri) sağlar, ancak yine de tüketicileri ve yazıcıları etkisiz (idempotent) olacak şekilde tasarlayın.
Şemalar, Uyumluluk ve Doğrulama
- Pub/Sub şemaları
- Merkezi olarak depolanan şemalarla Avro ve Protocol Buffers için yerel destek.
- Konu düzeyinde şema ayarları: kodlama (encoding) (Avro veya Protobuf) ve zorlama (enforcement) (yok, yalnızca doğrula veya zorunlu kıl).
- Üretici (Producer), kodlanmış yükleri (payloads) yayınlar; zorlama etkinleştirildiğinde Pub/Sub, mevcut şemaya göre doğrulama yapar.
- Evrim ve uyumluluk
- Geriye dönük uyumlu değişiklikler kullanın (isteğe bağlı alanlar ekleyin, Avro’da varsayılan değerli alanlar ekleyin, Protobuf’ta etiketleri asla yeniden kullanmayın, alanları kaldırmaktan veya yeniden adlandırmaktan kaçının).
- Şemaları açıkça versiyonlayın. Kırıcı değişiklikler (breaking changes) için, v1 ve v2 konularına çift yayın yapın veya bir sürüm alanı ekleyip buna göre yönlendirme yapın.
- Üretici-tüketici sözleşmeleri
- Tüketiciler bilinmeyen alanları göz ardı etmeli ve eksik olanları varsayılan değerleriyle almalıdır.
- Üretime (prod) almadan önce tüm tüketicilerde şema uyumluluğunu test edin; hazırlık (staging) aboneliklerinde prod ile aynı şema zorlamasıyla doğrulama yapın.
Kısa Avro örneği (alıntı):
undefined
Olay Güdümlü Entegrasyon, Eventarc ve Kafka ile Birlikte Çalışabilirlik
- Eventarc ve CloudEvents
- Eventarc, Google Cloud hizmetlerinden (ve Pub/Sub aracılığıyla özel kaynaklardan) gelen olayları CloudEvents spesifikasyonunu kullanarak Cloud Run, GKE veya Workflows’a yönlendirir.
type,source,subjectgibi nitelikler, ayrıntılı filtreleme ve denetlenebilirlik sağlar. - Fan-out’u en aza indirmek ve aşağı akış (downstream) yükünü azaltmak için nitelik filtreleri kullanın.
- Teslimat en az bir kez (at-least-once) yapılır; işleyicileri (handler) mümkün olduğunca bir kez etkili (idempotent) ve durumsuz (stateless) yapın.
- Eventarc, Google Cloud hizmetlerinden (ve Pub/Sub aracılığıyla özel kaynaklardan) gelen olayları CloudEvents spesifikasyonunu kullanarak Cloud Run, GKE veya Workflows’a yönlendirir.
- Eventarc tetikleyici (trigger) örneği:
gcloud eventarc triggers create gcs-finalize-to-run
–destination-run-service=ingestor
–event-filters=“type=google.cloud.storage.object.v1.finalized”
–event-filters=“bucket=my-data-bucket”
–service-account=eventarc-sa@PROJECT_ID.iam.gserviceaccount.com - Kafka ile birlikte çalışabilirlik ve yönetilen geçiş (migration)
- Dataflow şablonları, aşamalı geçiş için Kafka <-> Pub/Sub arasında bağlantı kurar. Konuları (topic) anahtarları (key) korunarak yansıtın (mirror); önce tüketicileri (consumer), sonra üreticileri (producer) taşıyın veya geçiş sırasında çift yazma (dual-write) yapın.
- Pub/Sub Lite, anahtar tabanlı yönlendirme ve daha düşük maliyet ile bölümlenmiş (partitioned), kapasite-sağlanmış (capacity-provisioned) akış sunar; bölgesel/zonaldir ve öngörülebilir kapasite ile bölüm başına sıralamanın (per-partition ordering) birincil öncelik olduğu Kafka benzeri iş yükleri için uygundur.
- Geçişle ilgili dikkat edilmesi gerekenler:
- Sıralama: Kafka anahtarlarını Pub/Sub sıralama anahtarlarına (ordering key) veya Lite bölümlerine (partition) eşleyin.
- Offset’ler: Teşhis (diagnostics) için offset’leri mesaj nitelikleri olarak taşıyın; tüketiciler geçişten sonra Kafka offset’lerine güvenemez.
- Teslimat: En az bir kez teslimatı (at-least-once) kabul edin; aşağı akışta (downstream) bir kez etkililiği (idempotency) zorunlu kılın.
- Şemalar: Confluent Schema Registry tanımlarını Pub/Sub şemalarına taşıyın veya uyumlu evrim kurallarıyla Protobuf/Avro üzerinde standartlaşın.
Akış Halinde Veri Alım Desenleri, Aktarım Hızı, Ölçeklendirme, Güvenlik ve Operasyonlar
- Gerçek zamanlı veri alım desenleri
- Pub/Sub -> Dataflow -> BigQuery: yüksek aktarım hızı ve akış ofsetleri ile idempotency (tekrar etkisizlik) için BigQuery Storage Write API havuzunu kullanın; hataları incelemek üzere bir dead-letter (işlenemeyen mesaj) tablosuna yönlendirin.
- Pub/Sub -> Dataflow -> Cloud Storage: ham olayları yeniden işlemek üzere arşivleyin; maliyet ve gecikmeyi dengelemek için pencereli, sıkıştırılmış yazma işlemleri kullanın.
- Pub/Sub -> operasyonel veri depoları: düşük gecikmeli aramalar için Bigtable’a, güçlü tutarlılığa sahip işlemler için Spanner’a veya iş yükü ihtiyaçlarına göre Cloud SQL/Firestore’a yazın. Benzersiz bir olay kimliği (event ID) ile anahtarlanmış idempotent upsert (güncelleme veya ekleme) işlemleri sağlayın.
- En az bir kez teslim, yinelenenleri önleme ve idempotency
- Her mesajda benzersiz bir event_id ve event_time taşıyın; üretici tarafında UUID’leri zorunlu kılın.
- BigQuery akışında yinelenenleri temizleme (de-dup): insertId ayarlayın veya sıralı akışlarla Storage Write API kullanın; yine de sorguları yinelenenleri temizleme mantığıyla koruyun.
- Sorgu zamanında yinelenenleri temizleme örneği: WITH ranked AS ( SELECT t.*, ROW_NUMBER() OVER (PARTITION BY event_id ORDER BY event_time DESC) AS rn FROM dataset.events t ) SELECT * EXCEPT(rn) FROM ranked WHERE rn = 1;
- Push uç noktaları için, yalnızca başarılı işlemden sonra 2xx kodu döndürün; aksi takdirde yeniden teslimat bekleyin.
- Mesaj aktarım hızı, kotalar ve ölçeklendirme
- Yayıncılar (Publishers): mesajları toplu halde gönderin ve bağlantıları yeniden kullanın; birden çok istemci arasında paralelleştirin. Sıralı iş yüklerini ölçeklendirmek için çok sayıda sıralama anahtarı (ordering key) kullanın.
- Aboneler (Subscribers): akış kontrolü (maksimum bekleyen bayt/mesaj) ile streaming pull’u tercih edin. Onay (ack) son tarihlerini işlem süresine göre boyutlandırın ve gerektiğinde uzatın.
- Hacimler arttıkça yayınlama ve abone olma aktarım hızı için kota artışlarını izleyin ve talep edin; ani artışları karşılamak için bir pay (örneğin, beklenen zirvenin 2 katı) bırakarak tasarım yapın.
- Tutarlılık ve kullanılabilirlik
- BigQuery akışı, sorgu görünürlüğü için nihai tutarlıdır (eventually consistent); akışla gelen satırları içermesi gereken etkileşimli sorgular için, gözlemlenen gecikmeye dayalı olarak bekleyin (örneğin, P50 kullanılabilirlik gecikmesinin 2 katı) veya Dataflow’da filigran (watermark) hizalı toplamalar kullanarak tasarım yapın ve somutlaştırılmış (materialized) sonuçları sorgulayın.
- Güvenlik
- IAM: en az ayrıcalık ilkesine uygun roller atayın (konu üzerinde üreticilere pubsub.publisher; abonelik üzerinde tüketicilere pubsub.subscriber). Her iş yükü için adanmış hizmet hesapları kullanın.
- Push kimlik doğrulaması: bir hizmet hesabından OIDC token’ları eklemek için push aboneliklerini yapılandırın; uç noktada hedef kitle (audience) doğrulamasını zorunlu kılın. Dahili kimlik doğrulama ve TLS için Cloud Run özel uç noktalarını tercih edin.
- Şifreleme: Pub/Sub, aktarım sırasında ve bekleme durumundaki verileri şifreler; müşteri tarafından yönetilen anahtarlar için konularda CMEK kullanın. Veri sızdırma riskini azaltmak için VPC Service Controls uygulayın. Gerekirse hassas yük (payload) alanları için istemci tarafı şifreleme kullanın.
- Gecikme, yeniden teslimat ve abone hatalarının operasyonel teşhisi
- Cloud Monitoring ile izleyin:
- birikme (backlog) için subscription/num_undelivered_messages ve oldest_unacked_message_age.
- yinelenenlere neden olan kaçırılmış onayları (ack) tespit etmek için expired_ack_deadline_count.
- aktarım hızı için publish_request_count ve pull_request_count.
- Eksik pano olaylarını araştırmak için, bilinen bir veri setini işlem hattından yeniden geçirin ve hatalı dönüşümü (transform) veya havuzu (sink) izole etmek için aşama aşama çıktıları karşılaştırın.
- Dataflow akışı için:
- Birçok kaynaktan gelen yükü karşılamak için uygun bir maxWorkers ile otomatik ölçeklendirme kullanın.
- Uyumsuz güncellemeler için işlem hatlarını boşaltın (drain), böylece devam eden işlerin tamamlanmasına izin verin ve veri kaybını önleyin.
- BigQuery ekleme bildirimleri için, belirli tablolara göre filtrelenmiş Cloud Logging denetim girdilerini, uyarı amacıyla bir Pub/Sub konusuna bir havuz (sink) aracılığıyla yönlendirin.
- Cloud Monitoring ile izleyin:
Pratik Problem Senaryosu
Contoso Freight’in, kamyonlardan dakikada 10.000 IoT telemetri mesajı almak, olayları zenginleştirmek, etkileşimli analizleri güçlendirmek ve harici ortaklardan gelen dosya bırakmalarında iş akışlarını tetiklemek için küresel, gerçek zamanlı bir olay platformuna ihtiyacı var. Bazı ortak CSV’leri bozuk biçimli satırlar içeriyor ve analiz ekibinin akışı engellemeden hataları incelemesi gerekiyor.
- Temel mesajlaşma ve şema katmanını oluşturun
- Eylem: Telemetri için bir Avro şeması tanımlayın ve bunu, şema zorunluluğu (schema enforcement) ‘require’ olarak ayarlanmış bir Pub/Sub konusu olan telemetry’ye ekleyin. Mesaj sıralamayı etkinleştirin ve ordering_key = hash(device_id) ile yayınlayın.
- Gerekçe: Konu düzeyinde şema zorunluluğu, bozuk biçimli olayları erken bir aşamada reddeder. Cihaz başına sıralama, gerektiğinde sıralı işlemeyi desteklerken, hash’leme anahtarları yayarak aktarım hızını korur.
- İzolasyon ve dead-lettering ile abonelikleri sağlayın
- Eylem: Dataflow için bir dead-letter konusu telemetry-dlt ve max_delivery_attempts=10 ile bir pull aboneliği olan telemetry-stream-sub oluşturun. Ham olayları soy takibi (lineage) ve yeniden oynatma için zaman bölümlü bir tabloya aktarmak üzere bir BigQuery aboneliği olan telemetry-raw-bq ekleyin.
- Gerekçe: DLQ (Dead-Letter Queue), ‘zehirli’ mesajları araştırma için izole eder. Ayrı bir BigQuery aboneliği, ham olayların işlem hattından bağımsız olarak korunması için düşük operasyonlu bir dışa aktarma yolu sağlar.
- Zenginleştirme ve havuzlar için bir Dataflow akış işlem hattı oluşturun
- Eylem: Akış kontrolü ile streaming pull kullanarak telemetry-stream-sub’dan veri alın. Şemaya göre doğrulayın, referans verilerle zenginleştirin ve pencereli toplamaları hesaplayın. Adlandırılmış bir akış (named stream) ve insertId = event_id ile Storage Write API kullanarak BigQuery’ye yazın; ham yedekleri saatlik olarak Cloud Storage’a yazın; kötü/başarısız kayıtları bir dead-letter BigQuery tablosuna yönlendirin.
- Gerekçe: Storage Write API, insertId/akış ofsetleri aracılığıyla idempotency ile yüksek aktarım hızlı, düşük gecikmeli yazma işlemleri sağlar. Bir dead-letter tablosu, akışı engellemeden incelemeyi destekler ve Cloud Storage arşivleri yeniden oynatmayı mümkün kılar.
- Analitikte yinelenenleri ve nihai tutarlılığı yönetin
- Eylem: Yinelenenleri hariç tutması gereken etkileşimli sorgular için, her kayıtta event_id ve event_time yayınlayın ve bir yinelenenleri temizleme (dedup) görünümü kullanın: CREATE OR REPLACE VIEW analytics.latest_events AS SELECT * EXCEPT(rn) FROM ( SELECT e.*, ROW_NUMBER() OVER (PARTITION BY event_id ORDER BY event_time DESC) rn FROM analytics.events e ) WHERE rn = 1; Gözlemlenen BigQuery akış kullanılabilirliğine dayalı olarak kısa bir sorgu gecikmesi ekleyin (örneğin, medyan gecikmenin iki katı).
- Gerekçe: En az bir kez teslimat, idempotent yazma işlemleri ve sorgu zamanında yinelenenleri temizleme gerektirir. Beklemek, akış görünürlüğü gecikmesi göz önüne alındığında, devam eden satırlardan kaynaklanan ıskalamaları azaltır.
- Ortak dosyası bırakmalarını Eventarc ile entegre edin
- Eylem: partner-drops bucket’ı için Cloud Storage object.finalized olaylarını, CSV’leri BigQuery’ye yüklemek üzere bir toplu Dataflow işi başlatan bir Cloud Run hizmetine yönlendirmek için Eventarc’ı yapılandırın ve ayrıştırma (parse) hatalarını bir dead-letter tablosuna gönderin.
- Gerekçe: Eventarc, bucket ve nesne önekine göre CloudEvents filtrelemesi ile olay güdümlü düzenleme (orchestration) sağlar. Bir toplu Dataflow işi, iyi verileri anında yüklerken analiz için bozuk biçimli satırları ayırır.
- Platformu güvenli hale getirin
- Eylem: Farklı hizmet hesapları kullanın: üreticiler telemetry üzerinde pubsub.publisher rolünü alır; Dataflow çalışanı SA’sı (hizmet hesabı) telemetry-stream-sub üzerinde pubsub.subscriber rolünü ve hedef BigQuery veri setleri ile Cloud Storage’a yazma erişimini alır; Eventarc tetikleyicisi, Cloud Run üzerinde invoker rolüne sahip adanmış bir SA kullanır. telemetry konusu ve BigQuery veri setlerinde CMEK’i etkinleştirin. Varsa, push uç noktalarını OIDC ve hedef kitle (audience) denetimleriyle yapılandırın.
- Gerekçe: En az ayrıcalık ilkesine uygun IAM ve CMEK, güvenlik ve uyumluluk gereksinimlerini karşılar; kimliği doğrulanmış teslimat, sahtekarlığı (spoofing) önler.
- Güvenilir bir şekilde işletin ve ölçeklendirin
- Eylem: Zirveleri karşılamak için Dataflow otomatik ölçeklendirmeyi cömert bir maxWorkers ile ayarlayın. subscription/oldest_unacked_message_age ve expired_ack_deadline_count metriklerini izleyin; eşikler aşıldığında uyarı verin. Uyumluluğu bozan işlem hattı değişiklikleri için, mesaj kaybını önlemek amacıyla boşaltma (drain) ile dağıtım yapın. Gecikme artarsa, abone paralelliğini artırın ve onay (ack) son tarihlerini işlem süresiyle orantılı olarak uzatın.
- Gerekçe: Proaktif izleme, gecikmeyi ve yeniden teslimatları erken tespit eder. Otomatik ölçeklendirme ve ayarlanmış onay son tarihleri, yinelenen fırtınalarını önler. Boşaltma (Draining), yükseltmeler sırasında devam eden mesajları korur.
Bu tasarım, olay güdümlü toplu iş entegrasyonu ile dayanıklı, güvenli ve gözlemlenebilir gerçek zamanlı veri alımı sunar, yinelenenlere karşı toleransı ve şema evrimini destekler ve hedeflenmiş düzeltme için kötü verileri izole ederken hızlı analizler sunar.
← Dataflow ve Apache Beam ile Akış İşleme · Tüm alanlar · Spark →
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 →