← All articles

Data Engineering

Your Streaming Pipeline Will Break in Four Places

Pub/Sub, Dataflow and Cloud Scheduler, explained through the things that go wrong.

Every tutorial shows you the happy path: a publisher sends a message, a subscriber receives it, Dataflow writes it to BigQuery. Nothing in that diagram tells you why your row count is 3% higher than the source system, or why yesterday's events landed in today's window.

So here are the same three services, taught backwards. Four failure modes, what causes each one, and the specific setting or pattern that fixes it. If you understand these, the happy path explains itself.

Four places a streaming pipeline breaks, with the fix for each

First, the shape of the thing

Pub/Sub has five parts. A publisher sends a message to a topic. A subscription attached to that topic holds messages for a subscriber until they are acknowledged. One topic can feed many subscriptions, and each subscription gets its own independent copy of the stream. That last sentence is the whole fan-out model.

Underneath, Pub/Sub splits into a control plane and a data plane. The control plane decides which servers a given publisher or subscriber talks to. The data plane moves the actual messages: the servers facing publishers are called routers, and the ones facing subscribers are called forwarders. You will never configure either, but knowing the split is why the service scales globally without you provisioning anything.

Dataflow is a managed runner for Apache Beam pipelines. Beam gives you one programming model, Read then Transform then Write, that compiles to both batch and streaming from the same codebase. That is the actual selling point: not speed, but the fact that your backfill and your live pipeline are the same code.

Cloud Scheduler is managed cron. It fires an HTTP call, a Pub/Sub message or an App Engine request on a schedule, and it is usually the thing that starts everything else.

One Pub/Sub topic feeding three subscriptions, pull and push

Break 1. The same message arrives twice

Pub/Sub's default guarantee is at-least-once delivery. Read that literally: it promises a message will not be lost, and promises nothing about it arriving only once.

Duplicates come from the acknowledgement deadline. When a subscriber receives a message it gets a lease, 10 seconds by default, and the client libraries extend that automatically as work continues, up to an hour. If the lease expires before the subscriber acknowledges, Pub/Sub assumes the subscriber died and redelivers.

So a duplicate almost always means one of four things: the subscriber is slower than the deadline, the volume being pushed exceeds what it can handle, the application crashed mid-processing, or processing genuinely takes longer than the lease and the extension is not configured.

The cost is not theoretical. Duplicate messages become duplicate rows, duplicate rows become wrong aggregates, and every reprocessed message is compute you pay for twice.

The fix, in order of how much work it is
Monitor acknowledgement operations for the expired response code, which is how you detect the condition at all. Then make the subscriber idempotent, so processing the same message twice leaves the same result: key the write on the message id, or MERGE rather than INSERT. Then, if you genuinely need it, enable exactly-once delivery.

Exactly-once is worth a paragraph of its own, because the constraints decide whether you can use it. It works on pull subscriptions only. Push and export subscriptions do not support it. The guarantee holds only when subscribers connect in the same region, and it costs significantly higher publish-to-subscribe latency. Dataflow adds its own exactly-once processing semantics on top, which is a different guarantee from delivery and is why Beam is the usual answer to this problem.

Timeline of an expired acknowledgement deadline causing a duplicate delivery

Break 2. The messages arrive out of order

By default, Pub/Sub does not order anything. Messages published in sequence can arrive in any sequence, because the whole architecture is optimised for throughput across many servers rather than a single ordered log.

This is the most commonly repeated outdated claim about Pub/Sub. Ordering does exist, and has since 2020, through ordering keys.

How ordering keys actually work
Set the same ordering key string on every message that must stay in sequence, typically a customer id or a row primary key. Messages sharing a key are delivered in order. Messages with different keys have no ordering relative to each other, and messages with an empty key are not ordered at all.

Three constraints matter more than the feature itself. All messages with a given key must be published in the same region. Publishing throughput is capped at 1 MBps per ordering key, so a hot key becomes your bottleneck. And ordering must be enabled when the subscription is created and cannot be changed afterwards.

Combine ordering with exactly-once and throughput drops to the order of thousands of messages per second, because acknowledgements must also be processed in order. Which leads to the design answer worth remembering: order per entity, never globally. You almost never need every event in the system ordered. You need each customer's events ordered.

Messages without an ordering key arrive in any order; with a key they arrive in sequence

Break 3. The data is late, and lands in the wrong window

This one is not a Pub/Sub problem. It is the hardest idea in streaming, and it is where Dataflow earns its keep.

Every event has two timestamps. Event time is when the thing happened: when the user tapped the button. Processing time is when your pipeline saw it. The gap between them is skew, and it is never zero: a phone goes through a tunnel, a network retries, a batch of IoT readings uploads when the device reconnects.

Streaming pipelines aggregate over windows, say five minutes of data at a time. If you window by processing time, an event that took twenty minutes to arrive gets counted in the wrong window, and both windows are now wrong. If you window by event time, you need to know when to stop waiting.

Watermarks and triggers
The watermark is Dataflow's estimate of how far event time has progressed: effectively "I believe I have seen everything up to 10:05". When the watermark passes the end of a window, the window fires and emits a result. Triggers control that firing: emit early for a provisional answer, emit again when late data arrives, and decide whether the late element updates the old result or is discarded.

The practical version: decide how late is too late, set allowed lateness to that, and pick whether a late arrival corrects history or is dropped. Getting this wrong does not throw an error. It silently produces numbers that do not match the source, which is the worst failure mode there is.

A late event counted in the wrong window by processing time, and the right one by event time

Break 4. The half nobody schedules

Streaming gets the attention and the monitoring. The batch half runs on Cloud Scheduler, and it fails quietly.

The usual shape is Scheduler firing a Cloud Function, which then acts on something else: starting and stopping Compute Engine instances on a working-hours schedule, backing up buckets in Cloud Storage, kicking off a load. On the warehouse side, BigQuery scheduled queries and the BigQuery Data Transfer Service do the same job for SQL that runs on a cadence.

None of this is complicated. It is simply the part where a job stops running on a Tuesday and nobody notices until the month-end numbers are short. If you build it, alert on the absence of a successful run, not only on failures. A job that never started produces no error to alert on.

Cloud Scheduler triggering a Cloud Function that acts on Compute Engine, Cloud Storage and BigQuery

Why Dataflow rather than a cluster you maintain

Spark and Flink both run well, and both run on infrastructure you install, patch, size and pay for while it idles. Dataflow is fully managed, with horizontal autoscaling, automated resource management and flexible resource scheduling that trades batch latency for a lower price.

The rest of the feature list matters only when procurement or security asks: templates, inline monitoring, customer-managed encryption keys, VPC Service Controls, private IPs, Dataflow SQL.

Its lineage explains the design. MapReduce handled batch. FlumeJava added pipeline composition. MillWheel solved streaming with watermarks. Dataflow is those three ideas unified, and Beam is that model handed to everyone else.

One Apache Beam codebase running as both a batch and a streaming pipeline

The access control you will be asked about

Pub/Sub permissions are granted per topic or per subscription, or project-wide when you are being lazy. The roles are roles/pubsub.publisher, roles/pubsub.subscriber, roles/pubsub.viewer, roles/pubsub.editor and roles/pubsub.admin. Dataflow has Admin, Developer, Viewer and Worker.

Cross-project publishing works and is common: project A publishes into a topic owned by project B by granting A's service account the publisher role on that topic. Granting project-level editor to make it work is the mistake to avoid.

What to take from this

Pub/Sub guarantees delivery, not uniqueness, and not order unless you ask for it by key. Dataflow's job is to make late data land in the right window. Scheduler runs the half of your pipeline that nobody watches.

Design for duplicates, order per entity, window on event time, and alert on jobs that did not run. Those four decisions prevent most of what goes wrong in production.

Watch the pipeline get built

This article covers one stage of a larger pipeline. Edukast AI builds the whole thing on YouTube, stage by stage: Stage 2, Ingest covers where event streams like Pub/Sub and Kafka sit in the pipeline, and the full Data Engineering series runs from collection through to consumption. More at edukastai.com.