← Professional Data Engineer: pipelines and data decisions
10 / 18 · 135 MIN

Ingestion, late events and replay

Define event identity and time; practice deduplication, conflicts, time closure and controlled recovery.

1. Define the ingestion contract

Before selecting services, describe what a row represents and when it may be used. An event can be an immutable instruction, an entity version or a correction to an earlier movement. These models need different rules. Under an immutable-event contract, repeating an identifier with matching content is a repeat; reusing that identifier with a different amount is a conflict. Under a versioned-change contract, one entity can legitimately have several versions. Do not remove versions merely because they share an entity key. Record event identity, entity identity, schema version, business time, publication time and value units. In a fictional APS example, a team receives portfolio movements to prepare daily close. The source can become unavailable and send its batch hours later. The business needs to know whether a movement entered the correct period, was counted twice or corrected an accepted report. Prepare three observable outputs: accepted records, rejected records with reasons, and period reconciliation. Define who decides on conflicts and corrections. An empty backlog can coexist with wrongly discarded movements. Acceptance must measure correct data use as well as ingestion speed. Examples in this lesson are fictional and do not describe BNP Paribas procedures.

2. Make startup reproducible

Separate submission identity, worker identity and resource access. A successful read under your own account does not prove workers can read the same source. When an error identifies a service account, confirm the principal and resource before increasing permissions. At submission, also check whether the authorized principal can attach the dedicated account to the job. During diagnosis, distinguish IAM rejection, missing resources, DNS errors and connection timeouts. Requesting more permissions for a timeout without investigating networking can broaden access while leaving the failure intact. Keep worker region, selected subnet and address capacity in the plan. For Shared VPC, the subnet reference must identify the host project. A team planning to grow from 20 to 100 workers should check free addresses with network owners. Record assumptions and capacity exercise results. For releases, retain identity of the image, Flex Template specification and approved parameters. A mutable tag is insufficient evidence of which executable passed testing. The change package should let the team reconstruct what was submitted and connect it to the data used in validation. This makes investigation of a failed launch a comparison of known conditions rather than a series of permission changes without a hypothesis.

3. Separate repeats, conflicts and legitimate changes

Transport identity and business identity answer different questions. A message ID helps recognize redelivery of that message. If a producer publishes the same instruction twice and receives different IDs, business identity is still needed. Define key scope: it may need source system, instruction identifier and version. Document how long deduplication memory lasts and what happens when replay exceeds that period. Keeping only the last observed key does not handle interleaved repeats or distributed execution. Retain incompatible content under one identity as evidence. If A contains quantity 7 and later quantity 9, dropping the second arrival as a duplicate hides a difference. Automatically replacing the first payload would also introduce a new rule. Route the conflict with sufficient context for a decision, avoiding unnecessary sensitive data exposure. At the destination, inspect the scope of the protection being used. Storage Write API offsets are positions within a stream, not global business keys across streams. After a lost response, retain stream, offset and payload for a consistent retry. The team should distinguish write acknowledgement from the movement’s functional state in the report. A technically accepted write can still contain a movement assigned to the wrong portfolio or period.

4. Decide how to handle late data

Event time describes when the business occurrence happened; processing time describes when a stage handles it. Define which one drives reporting. Yesterday’s movement published today can still belong to yesterday’s close. In Beam, windows group events and triggers determine emissions. A watermark represents estimated progress in event time, while lateness policy affects inclusion of later arrivals. Exactly-once does not prove completeness when events are dropped for lateness. The consumer contract must distinguish provisional results, technical closure and business acceptance. In our case, published reporting must retain history. The team agrees an allowance for normal operation and a separate correction process beyond that allowance. Retain excluded events, reasons, affected periods and decisions. Decide how a correcting version is identified and how consumers are notified. Do not increase the allowance without assessing retained state, latency and cost. Equally, do not interpret a closed window as permission to delete evidence. A useful exercise includes out-of-order delivery, arrivals inside and outside the allowance, approved corrections and repeats of old events. Compare the business expectation with what each policy actually includes. Before releasing a report, identify which evidence comes from the source and which conclusions merely follow from a configured timer.

5. Prepare replay without repeating effects

A replay plan starts with data availability and recovery scope. Check retention, starting boundary, selection and destination. A Pub/Sub snapshot holds acknowledgement state; it neither restores tables nor reverses notices already sent. Timestamp seek follows publication time, which can differ from business time. Also allow for a delivery transition after the request. To compare a calculation fix, choose an isolated destination and explicitly control external effects. An API call notifying a customer does not become idempotent merely because internal processing uses exactly-once. Build a recovery manifest: reason, code version, intended interval, available data, deduplication key, comparison output and promotion criterion. In the application, separate transient failures from invalid data. An invalid message does not become valid by retrying the same transform indefinitely. Implement failure routing inside the pipeline where needed; a subscription dead-letter topic does not automatically cover errors after acknowledgement at the first stage. Recovery should permit correction and reprocessing of selected records while retaining their link to the initial error. At the change review, show differences by key and period, rejection counts and effects blocked during the exercise. An accepted launch request is only the beginning of this evidence.

6. Lab: a small observable contract

The file content/labs/pde-event-replay/run.py implements an original Python model without networking or cloud dependencies. Run python3 content/labs/pde-event-replay/run.py < content/labs/pde-event-replay/case.json from the project root. Input contains windowSeconds, allowedLatenessSeconds and an events sequence. Each operation is either an event with id, eventTime and value, using integers for numeric fields, or an explicit watermark advance. Times are nonnegative seconds on a fictional scale. Windows are half-open: with width 60, time 60 belongs to [60,120). The program rejects unexpected fields, invalid IDs, decreasing watermarks and out-of-range numbers. Local policy is intentionally explicit. For a new ID, watermark >= window end + allowance produces too-late. Before that boundary, the event is accepted; if watermark >= window end, it receives accepted-late. For previously accepted IDs, equal content produces duplicate and different content produces conflict, including after closure. The first accepted payload remains. The model remembers IDs during this run, bounded to 10,000 operations, without persisting them between processes. It implements no triggers, panes, distributed state, checkpoints or Beam runner details. Use it to discuss a contract; it neither simulates nor validates Dataflow.

7. Predict results and seek counterexamples

Before running the example, trace its nine steps manually. A(10,7) is accepted and its repeat is duplicate. Watermark advances to 70. B(20,5) enters the first window as accepted-late, still within allowance 30. A(10,9) is conflict and does not replace 7. Watermark becomes 90. C(30,4) is too-late at the inclusive boundary. D(60,3) belongs to the second window and is accepted. The final repeat of A remains duplicate because ID memory remains. Expect sum 12 and count 2 in the first window, sum 3 and count 1 in the second. Policy closure applies only to the first. Now change one condition at a time. Use watermark 89 before C: the event becomes accepted-late. Use zero allowance: at watermark 60, new IDs from the first window are already rejected. Reverse A’s conflicting payload order: the retained value changes because the rule is first accepted. Repeat all events in a new process: earlier memory is gone. Explain in writing why these results do not prove exactly one write at an external destination. Finally, identify missing evidence for authorizing real replay: persistence, concurrency control, source validity and destination behavior when records repeat.

8. Hand over with testable criteria

The runbook should connect each symptom to a decision. Growing backlog calls for capacity and dependency investigation; conflicts call for identity and payload analysis; more too-late records call for comparison of source availability with time policy. Define owners, minimum evidence and intervention limits. For planned changes, describe what stopping reads and finishing pending work does. In Dataflow, drain can emit windows still incomplete relative to the business period. Cancel can leave in-flight data needing reconciliation. Record the boundary between executions and how consumers handle output emitted during shutdown. At RUN handover, ask a colleague to explain how they would detect an incomplete report despite successful job completion. Give them one repeated publication, one payload conflict and one arrival beyond the allowance. Their answer should identify available evidence, permitted decisions and follow-up validation. Summarize the path in an acceptance record: executed version, covered data, accepted and excluded records, differences, external effects and approval. The central lesson is that a technical guarantee needs a scope. Identity, time and destination need compatible contracts before a processing result can support a business decision. Keep unresolved discrepancies visible and assign an owner instead of treating a successful job status as their resolution.

"""Original bounded teaching model. No Beam runner, network or durable state."""
import json
import re
import sys


def integer(value, low, high):
 if type(value) is not int or not low <= value <= high:
 raise ValueError('integer outside the declared contract')
 return value


def exact(value, keys):
 if type(value) is not dict or set(value)!= set(keys):
 raise ValueError('unexpected or missing fields')


def evaluate(document):
 exact(document, ['windowSeconds', 'allowedLatenessSeconds', 'events'])
 width = integer(document['windowSeconds'], 1, 86400)
 grace = integer(document['allowedLatenessSeconds'], 0, 86400)
 events = document['events']
 if type(events) is not list or len(events) > 10000:
 raise ValueError('events must be a list of at most 10000 operations')
 # Validate the whole input before producing a result.
 last_watermark = 0
 for event in events:
 if type(event) is dict and set(event) == {'watermark'}:
 last_watermark = integer(event['watermark'], last_watermark, 10**9)
 else:
 exact(event, ['id', 'eventTime', 'value'])
 if type(event['id']) is not str or not re.fullmatch(r'[A-Za-z0-9_-]{1,64}', event['id']):
 raise ValueError('invalid logical event id')
 integer(event['eventTime'], 0, 10**9)
 integer(event['value'], -10**6, 10**6)
 watermark = 0
 accepted = {}
 decisions = []
 for index, event in enumerate(events):
 if 'watermark' in event:
 watermark = event['watermark']
 decisions.append({'index': index, 'outcome': 'watermark', 'value': watermark})
 continue
 identifier = event['id']
 payload = (event['eventTime'], event['value'])
 end = (event['eventTime'] // width + 1) * width
 if identifier in accepted:
 outcome = 'duplicate' if accepted[identifier] == payload else 'conflict'
 elif watermark >= end + grace:
 outcome = 'too-late'
 else:
 accepted[identifier] = payload
 outcome = 'accepted-late' if watermark >= end else 'accepted'
 decisions.append({'index': index, 'id': identifier, 'outcome': outcome})
 windows = {}
 for event_time, value in accepted.values:
 start = event_time // width * width
 row = windows.setdefault(start, {'start': start, 'end': start + width, 'count': 0, 'sum': 0})
 row['count'] += 1
 row['sum'] += value
 for row in windows.values:
 row['closedByLocalPolicy'] = watermark >= row['end'] + grace
 return {'accepted': [{'id': key, 'eventTime': value[0], 'value': value[1]} for key, value in sorted(accepted.items)],
 'windows': [windows[key] for key in sorted(windows)], 'decisions': decisions,
 'watermark': watermark, 'stateLifetime': 'this invocation only',
 'sourceCompletenessProven': False, 'cloudSemanticsValidated': False,
 'productionReplayApproved': False}


def unique_object(pairs):
 result = {}
 for key, value in pairs:
 if key in result:
 raise ValueError('duplicate JSON field')
 result[key] = value
 return result


if __name__ == '__main__':
 try:
 raw = sys.stdin.read(2_000_001)
 if len(raw) > 2_000_000:
 raise ValueError('input exceeds local limit')
 result = evaluate(json.loads(raw, object_pairs_hook=unique_object))
 print(json.dumps(result, sort_keys=True))
 except (ValueError, TypeError, RecursionError) as error:
 print(json.dumps({'error': str(error)}), file=sys.stderr)
 sys.exit(2)
IN PRACTICE

A(10,7) and B(20,5) sum to 12; A(10,9) is a conflict and C at the cutoff is too-late.

Common pitfalls

Confuse messages with events; replace conflicts; infer completeness from watermarks; repeat notices during replay.

Related topics: Migration reconciliation · Data modeling · Observability and SLOs

Take this idea with you

Specify identity, time policy and destination effects; reconcile the result before accepting it.

Create account

Reference: Professional Data Engineer standard exam guide · Current linked standard guide (document title v4.2); edition date unconfirmed (2026-09-30 inspection)

Google Cloud is a trademark of Google LLC. bigsavant.com is an independent preparation platform and is not affiliated with, associated with, sponsored, authorised or endorsed by Google. Content and questions are original, are not official exam questions, and completing our tests does not award or guarantee any certification. Names are used only to identify the subject. All other trademarks belong to their respective owners.