Memory hook: Event time says when; watermark estimates completeness.
Must remember
Pub/Sub buffers messages between independent producers and consumers. A partition/key strategy, acknowledgments, retries and dead-letter handling affect processing correctness. Separate transport delivery guarantees from the final business effect: a retried database write still needs deduplication or a transactional/idempotent design.
Dataflow runs Apache Beam pipelines for batch and streaming. Event time is the timestamp of the real event; processing time is when the worker handles it. Windows group records, watermarks estimate event-time progress, and triggers decide when to emit results. Allowed lateness and accumulation choices determine whether delayed events revise earlier answers. A watermark is an estimate, not proof that no older message will arrive.
Choose Managed Service for Apache Spark/Dataproc when preserving Spark/Hadoop libraries and existing jobs matters. Ephemeral clusters or serverless execution reduce idle cost; persistent clusters may suit repeated tightly scheduled workloads. Data Fusion offers visual integration; Dataform organizes SQL transformations and assertions in BigQuery. ELT loads before transformation; ETL transforms before loading. Neither excuses storing prohibited raw sensitive data.
Design a schema contract, malformed-record path, idempotency key and reconciliation counts before increasing throughput. Handle evolving fields compatibly, monitor freshness and backlogs, and avoid logging raw secrets. For stream joins, reason about state size, time bounds and late data. Hot keys can bottleneck an otherwise well-provisioned pipeline; redistribute or pre-aggregate where semantics permit.
AI enrichment is another processing dependency: bound timeouts/cost, retain model/version provenance, validate structured output and isolate failed records. A model response is not automatically correct or safe to publish.
Review details
Kafka uses partitioned retained logs, offsets and consumer groups; order is within a partition, not universal across a topic. Pub/Sub hides more broker operations and uses subscriptions/acknowledgments; do not import every Kafka partition-management assumption. Select by ecosystem compatibility, delivery/access patterns and operating effort.
Beam fixed windows divide time into distinct intervals; sliding windows overlap; session windows group activity separated by a gap. Triggers control output timing, and accumulation modes control whether subsequent panes include earlier results. “Exactly once” must be scoped to the documented processing/storage boundary; an external email or payment still needs idempotency.
Choose under exam pressure
| Requirement | Choice and reason |
|---|---|
| Portable batch and streaming transformations | Beam on Dataflow. |
| Preserve existing Spark jobs | Managed Spark/Dataproc. |
| Late transactions change hourly totals | Event-time windows with a deliberate late-data/trigger policy. |
Traps
- Exactly-once transport does not guarantee exactly-once external side effects.
- Adding workers does not fix a single hot key.
Active recall
1. Event time versus processing time?
When the event occurred versus when a worker handles it.
2. What does a watermark estimate?
Progress/completeness of event-time data.
3. Why retain a rejected-record stream?
To investigate, correct and replay malformed records without losing the whole pipeline.
4. When choose Dataform?
For managed SQL transformation workflows and assertions around BigQuery.
5. Why version model enrichment?
To reproduce and audit output when model behavior changes.