In a Dataflow pipeline, what does PubsubIO's withIdAttribute change about deduplication?
answer
- there are two kinds of duplicate here
- the service-assigned identifier is the default key
- a retried publish gets a fresh identifier
- the publisher has to stamp something stable
- dedup memory is bounded, not forever
basics
~20 sBy default the Dataflow runner deduplicates Pub/Sub messages using the service-assigned message ID. withIdAttribute tells PubsubIO to dedupe on a publisher-set message attribute instead, so two separate publishes of the same logical event collapse into one.
solid answer
~40 sWhen a Beam pipeline reads Pub/Sub on the Dataflow runner, the runner already removes duplicates caused by redelivery of the *same* message, keyed on the service-assigned message ID. That does not help when a publisher retried and published the same logical event twice: those are two distinct messages with two distinct IDs, and the pipeline sees both. `PubsubIO.readMessages().withIdAttribute("event_id")` (in Python, `ReadFromPubSub(..., id_label='event_id')`) tells the source to use a publisher-supplied attribute as the dedup key instead, so republished copies carrying the same value are collapsed. The publisher must actually set that attribute on every message, and the deduplication is scoped to a bounded time window rather than being remembered forever, so it protects against near-term republishing rather than acting as a permanent idempotency store. Sinks still need their own idempotency.
code
java · 6 linesPCollection<PubsubMessage> events = pipeline.apply(
"ReadEvents",
PubsubIO.readMessages()
.fromSubscription("projects/p/subscriptions/events")
.withIdAttribute("event_id")
.withTimestampAttribute("occurred_at"));go deeper
Know that the Pub/Sub source in a Beam pipeline can be told to deduplicate on a message attribute you choose instead of on the service-assigned message identifier.
Explain why the two cases differ: redelivery of one message shares an ID and is handled by the runner, while a retried publish creates a new message that only a stable publisher-supplied key can collapse.
Treat it as a contract, not a flag — the publisher must stamp a key computed once before the first attempt, the dedup horizon is bounded, and the sink still has to be idempotent across restarts and cutovers.
Own the end-to-end idempotency design: where the business key originates, who is accountable for stamping it across every producer, and why sink-side keyed upserts remain the durable guarantee regardless of source configuration.
## Two different duplicate problems "Duplicates from Pub/Sub" covers two distinct situations, and the whole value of this question is telling them apart. **Redelivery of the same message.** Pub/Sub may deliver a message more than once — for instance when an acknowledgement is delayed. Each delivery carries the *same* service-assigned message ID. **Republication of the same logical event.** An upstream service publishes an event, does not get a confirmation, retries, and publishes again. Now there are two genuinely different messages in the topic with two different message IDs that happen to describe the same real-world event. When a Beam pipeline reads Pub/Sub on the Dataflow runner, the runner handles the first case for you: it deduplicates on the Pub/Sub message ID, which is why Dataflow can offer exactly-once *processing* of records within the pipeline. It cannot handle the second case, because from the runner's point of view those are two unrelated messages. ## What withIdAttribute does `withIdAttribute` changes the dedup key from "whatever the service assigned" to "an attribute the publisher set": ```java PCollection<PubsubMessage> events = p.apply( PubsubIO.readMessages() .fromSubscription("projects/p/subscriptions/events") .withIdAttribute("event_id")); ``` In the Python SDK the same idea is `id_label`: ```python events = p | beam.io.ReadFromPubSub( subscription='projects/p/subscriptions/events', id_label='event_id', with_attributes=True) ``` Now the source treats two messages carrying the same `event_id` attribute as the same element and emits only one of them. If the publisher stamps a stable idempotency key — an order ID, an event UUID generated once before the first publish attempt — then publisher retries stop reaching your pipeline as duplicates. There is a companion knob worth knowing in the same breath: `withTimestampAttribute` (Java) / `timestamp_attribute` (Python), which tells the source to take each element's event time from a message attribute rather than from the Pub/Sub publish time. It is a different concern — event time versus identity — but it is configured the same way and is the other reason to set attributes on the publisher side. ## The three conditions this depends on 1. **The publisher must set the attribute, on every message.** A message without the attribute has nothing to dedupe on. This is a contract between the publishing service and the pipeline, and it is the part that breaks in practice: one producer in a fleet forgets, and duplicates start appearing. 2. **The value must be stable across retries.** Generating a fresh UUID inside the retry loop defeats the whole mechanism — the key must be computed once, before the first publish attempt, and reused by every retry of that same logical event. 3. **Deduplication is bounded in time, not permanent.** The source keeps a record of recently-seen IDs; it does not maintain a lifetime registry of every key it has ever observed. A copy republished long after the original will pass through. Treat the feature as protection against near-term republishing, not as a general-purpose idempotency store. ## What it does not give you This is source-side deduplication. It says nothing about what happens at the *other* end of the pipeline. Dataflow's exactly-once story is scoped to processing within a running job: if you cancel and restart, if you drain and relaunch, or if your sink is not idempotent, records can still be written twice downstream. The durable pattern is unchanged by anything on the source side: give each record a business key and make the sink an upsert or a keyed merge, so re-writing a record is a no-op. Equally, if what you need is deduplication over a long, explicitly-controlled horizon, or deduplication on a key computed from the message *body* rather than an attribute, that belongs in the pipeline as a deliberate stateful step (Beam has a `Deduplicate` transform with a configurable duration) rather than as a source option. ## How this shows up in an interview Usually as a scenario: "our downstream table has duplicate rows and we read from Pub/Sub — why?" The strong answer separates the layers out loud. Redelivery of one message is already handled by the runner. Duplicate *publishes* are not, and `withIdAttribute` with a publisher-supplied stable key is the source-side fix, subject to the publisher honouring the contract and to the bounded dedup horizon. And whatever you do at the source, the sink still has to be idempotent, because job restarts and replays exist.
- Why does the default message-ID deduplication not stop duplicates caused by a publisher retry?Because each publish attempt creates a genuinely new Pub/Sub message with its own service-assigned ID. The runner correctly sees two distinct messages and passes both through. Only a publisher-supplied key that is stable across retries — stamped once before the first attempt and reused — lets the source recognise them as the same logical event.
- Does configuring an id attribute mean the pipeline's sink no longer needs to be idempotent?No. This is source-side deduplication with a bounded horizon, and Dataflow's exactly-once guarantee is scoped to processing within one running job. Restarts after a cancel, drain-and-relaunch cutovers, and non-idempotent sinks all still produce duplicate writes, so keyed upserts or merge-on-business-key remain the durable protection.
- What is the related timestamp attribute option used for?It tells the Pub/Sub source to take each element's event time from a message attribute instead of the publish time, so windowing reflects when the event actually happened rather than when it reached the topic. It is configured alongside the id attribute and shares the same requirement: the publisher must set it consistently on every message.
saying these in an interview costs you the question
- Assumes reading from Pub/Sub guarantees no duplicates ever reach the sink
- Thinks the id attribute dedupes forever rather than over a bounded window
- Generates a new idempotency key inside the publisher's retry loop
- Confuses the identity attribute with the event-time timestamp attribute
- Drops sink idempotency because source deduplication is configured