
Welcome to Part 4 of Pub/Sub AI Bytes! Over the past three weeks, we explored in-flight AI inference with SMTs, proto-first streaming event telemetry, and generating and editing SMTs with Gemini in seconds.
This week, we dive into a mission-critical AI Architecture Pattern: how frontier AI labs build high-throughput, zero-ops training telemetry firehoses to monitor, safeguard, and optimize large-scale model training running across GPUs and TPUs with Cloud Pub/Sub. This AI Architecture pattern has been generalized so that it can be used for general-purpose model training use cases as well.
The challenge: Model training events at scale
Training frontier foundation models across thousands of accelerators — such as GPU clusters or Cloud TPU — is one of the most expensive workloads in computing. Because distributed training frameworks (like PyTorch FSDP, JAX) synchronize workers at every step, the entire cluster moves only as fast as its slowest node. This introduces two fundamental infrastructure challenges:
- Zero tolerance for I/O blocking: Experiment tracking and hardware logging must be strictly decoupled from the core training loop. If a single training node hangs, thousands of accelerators sit idle.
- Traffic burstiness: Training telemetry is cyclical. Network logging drops to near-zero during checkpoint loads, then explodes to million+ events per minute at step boundaries.
To maximize the ratio of productive training time to total cluster time, ML platform teams must continuously ingest various telemetry streams without stalling compute such as –
- Algorithmic & model state metrics: Training loss, validation perplexity, learning rate schedules, layer-wise gradient norms, token throughput, diagnostic tensors (embeddings, attention maps).
- Hardware & infrastructure telemetry: GPU/TPU utilization, junction temperatures, memory errors.
- System & lifecycle events: Collective communication timeouts, OOM triggers, shard-level checkpoint commit confirmations.
Why serverless Pub/Sub over self-managed brokers?
Historically, engineering teams turned to self-managed Apache Kafka clusters to buffer training telemetry. However, operating fixed brokers for bursty training workloads forces lean infrastructure teams to constantly provision for peak spikes, monitor broker disk, and manually rebalance partitions as clusters grow.
By adopting Google Cloud Pub/Sub as their unified telemetry backbone, one leading frontier AI lab scaled their training metrics pipeline from a prototype to millions of events per day without making any changes to the service:
- Elastic shock absorption: Pub/Sub scales horizontally and automatically to absorb synchronized step-boundary spikes with zero operational overhead.
- Experiment isolation: Rather than spinning up a dedicated topic for every training experiment, teams can publish to shared telemetry topic and isolate experiments downstream using either Single Message Transforms or Pub/Sub Subscription Message Filtering on attributes like run_id.
The reference architecture
Here is the end-to-end architecture handle billions of daily training events across multi-cloud accelerator clusters:

In this architecture, telemetry emitted from the GPU / TPU Training cluster is split across two complementary ingestion paths to ensure logging never blocks synchronized training steps. During normal operation, high-frequency training metrics (loss, learning rate, gradient norms) and hardware health signals under the 10 MB limit flow directly along the fast path into Cloud Pub/Sub (<10 MB). When payloads exceed Pub/Sub’s 10 MB message limit — such as high-dimensional embeddings, attention maps, or local flow-control overflow bursts — the publisher diverts them along the dotted offload path into Cloud Storage (>10 MB) and publishes a lightweight URI reference to Pub/Sub so downstream consumers know a large blob is ready for processing.
From Cloud Pub/Sub, the unified event stream fans out in real time across three specialized consumption paths:
- Zero-ETL analytics: A Zero-ETL BigQuery subscription streams structured training and hardware telemetry directly into Cloud BigQuery for historical run comparison, Model FLOPs Utilization (MFU) profiling, and ad-hoc SQL diagnostics with no intermediate pipeline to manage.
- Sub-second cluster automation: A Push subscription delivers critical health and anomaly events via low-latency HTTPS requests directly to services running on GKE, enabling the cluster orchestrator to trigger immediate actions such as predictive node eviction, checkpoint recovery, or automated rollbacks.
- Large-payload path: For events that spilled over to Cloud Storage (>10 MB), Cloud Dataflow asynchronously fetches, decompresses, and processes the referenced objects before forwarding the enriched signals to the downstream GKE operational services.
Turning telemetry into action
Because large training clusters are expensive to run, telemetry streams passing through Pub/Sub are converted into immediate automated actions across three consumer paths:
1. High-fidelity infrastructure profiling
Using Pub/Sub BigQuery Subscriptions (alongside Dataflow cold-path routing), raw telemetry streams into partitioned BigQuery tables to help platform engineers with isolating network & I/O bottlenecks. Engineers can run SQL queries cross-referencing wait times against packet drop metrics to isolate faulty network switches throttling the cluster, or compare GPU utilization against storage read queues to tune data pre-fetching caches.
2. Automated failure mitigation and predictive node eviction
At hundreds of GPU scales, hardware failures are expected. Instead of relying on slow polling loops, teams use Pub/Sub Push Subscriptions to deliver critical infrastructure events via low-latency HTTPS webhooks directly to cluster orchestrators (e.g., GKE operators or Slurm controllers) the millisecond they occur:
- Predictive node eviction: When a consumer detects rising memory errors on a specific host, the orchestrator flags the node and schedules a graceful migration before a hard crash stalls the job.
- Instant checkpoint resumes: On an unrecoverable hardware fault, the controller immediately cordons the broken node, swaps in a hot-spare node, and instructs the scheduler to resume training from the latest checkpoint.
3. Model health auditing
Training runs can also fail algorithmically. Large language models frequently experience sudden loss spikes or gradient explosions caused by combinations of particular data batches and optimizer states. Cloud Dataflow can support with automated checkpoint rollbacks. Dataflow tracks sliding-window moving averages of gradient norms and loss values. If it detects an exponential gradient spike or NaN (Not a Number) loss anomaly, it triggers an automated webhook to pause training, roll back model parameters to a healthy checkpoint ~N steps prior, adjust the learning rate, and skip the problematic data batches that triggered the instability.
Operational trade-offs
When designing a model training telemetry firehose on Cloud Pub/Sub, evaluate these three operational trade-offs:
- Retention & replayability: Unlike self-managed brokers that retain logs on disk by default, Pub/Sub removes messages once acknowledged by subscribers. To reprocess historical training telemetry (for example, re-running an anomaly detector over yesterday’s steps), leverage Pub/Sub message storage offerings (e.g. Topic Message Retention (up to 31 days))
- Chronological node ordering: Pub/Sub optimizes for massive horizontal scale and unordered delivery. If a downstream state machine requires logs in the exact chronological sequence emitted by a specific worker rank, enable Pub/Sub Ordering Keys.
- Cross-cloud data ingress: If your compute cluster runs in an external cloud, factor in network egress costs from the source environment. Combining compression, client-side batching, and larger blob offloading to object store such as Cloud Storage will help with cost optimization.
Tune in next Friday for Part 5 of Pub/Sub AI Bytes, where we explore Zero-ETL Streaming: Direct to Database with NEW Bigtable Subscriptions!!
Pub/Sub AI Bytes: Part 4 — High-throughput ingestion for large-scale model training was originally published in Google Cloud – Community on Medium, where people are continuing the conversation by highlighting and responding to this story.
Source Credit: https://medium.com/google-cloud/pub-sub-ai-bytes-part-4-high-throughput-ingestion-for-large-scale-model-training-b628232d85f9?source=rss—-e52cf94d98af—4
