A data engineer is designing a new data ingestion pipeline for a global IoT platform. Millions of devices send telemetry data every second to a regional Cloud Pub/Sub topic. The data needs to be processed by a Dataflow streaming job and then stored in BigQuery. The engineer is concerned about potential message loss or duplicates during ingestion and processing, especially under high load or network interruptions. They need to ensure that each unique telemetry event is processed exactly once in BigQuery. How can this be achieved in the Dataflow pipeline?
- AUse Cloud Storage as an intermediate buffer to deduplicate messages before Dataflow reads them.
- BConfigure the Pub/Sub subscription to use 'at-least-once' delivery and implement custom duplicate detection in Dataflow.
- CSet the BigQuery sink to use 'insertAll' with best-effort deduplication.
- DEnable Dataflow's exactly-once processing guarantee by using unique event IDs and stateful processing.
Show answer & explanationAnswer & explanation
Correct answer: D. Enable Dataflow's exactly-once processing guarantee by using unique event IDs and stateful processing.
Dataflow's exactly-once processing guarantee for streaming pipelines, when properly configured with unique event IDs and stateful processing (e.g., using `StatefulDoFn` or `Combine.Globally` with state), ensures that each unique event is processed and its effects are reflected in the output exactly once, even in the face of failures or retries.
Why the other options are wrong
- A. Using Cloud Storage as an intermediate buffer would introduce significant latency, complexity, and cost, and would not inherently provide exactly-once processing without further complex mechanisms.
- B. At-least-once delivery from Pub/Sub means duplicates are possible. Custom duplicate detection in Dataflow is possible but less robust and performant than Dataflow's built-in exactly-once features for this scale.
- C. BigQuery's `insertAll` with best-effort deduplication is not a strong 'exactly-once' guarantee and relies on client-provided insert IDs which might not cover all scenarios or provide the same robustness as Dataflow's end-to-end guarantee.
Dataflow Exactly-Once Processing
A Dataflow streaming pipeline guarantee that ensures each unique input record affects the output exactly one time, even in the presence of system failures or retries, typically achieved using unique event IDs and stateful processing.
- Critical for financial transactions, IoT, and accurate aggregations.
- Requires unique event identifiers.
- Leverages Dataflow's fault tolerance and state management.
Memory trick: Unique IDs and Dataflow's state ensure it's done exactly once.