1. Define what must be recovered
An ingestion failure leaves at least three questions open: which changes still exist at the source, which state the pipeline saved, and which effects already reached the destination. A process that runs again does not automatically answer all three. Before choosing replay or backfill, record the last confirmed point at each boundary and the evidence supporting it. Distinguish an observation timestamp from a replication position: the first tells you when you looked; the second participates in the source continuity contract. In a fictional fund-position example, the team needs current balances for an operational query and every intraday change for a separate control. Backfill can help rebuild balances without recovering every intermediate operation. Create two acceptance criteria and two owners. If history cannot be reconstructed, retain an explicit exception for that requirement. Equal balances do not authorize declaring that no operation was lost. Also define who can stop handover and where reconciliation results are stored, preventing closing pressure from turning an assumption into evidence.
2. Resume the source without hiding a missing interval
In Datastream, recovering a permanently failed stream requires selecting a position appropriate to the source type. If a required MySQL binlog can be recovered, evaluate restoring that file and retrying the current position before jumping to the latest one. With a new PostgreSQL slot, changes between the lost position and the first available LSN do not reappear merely because the stream runs again. Record that interval as a recovery or reconciliation obligation, including the limits of available history. Do not compare positions from two sources merely because their database names match. After failover, verify sequence correspondence on the server now producing changes. A connection test confirms connectivity but does not demonstrate log continuity. When planning the intervention, record available retention, affected objects, writes continuing during the change and the destination where overlap will be rehearsed. The starting-position decision must be explainable to whoever validates the data afterwards. A fast recovery that skips an interval needs an explicit plan for that interval.
3. Control historical reading and application
Completed backfill status confirms that the object has been read; destination loading may continue. The handover owner should therefore not accept only a screenshot of status. Request destination evidence at a comparable cut, including relevant keys, values and deletions. Where concurrent writes exist, document how CDC events and historical images are ordered. The last file received can contain an older image than the state already applied. Stopping a backfill can be a valid way to protect the source from excessive load. However, when restarting a stopped object, the plan must allow for its existing data to be transmitted again. Time and load budgets should not assume exact continuation at the row where reading stopped. Rehearsing an object of representative size helps select concurrency and the change window. Do not extrapolate to every table from a small table with no ongoing changes. Record throughput, source impact, application lag and criteria for reducing concurrency or stopping the operation.
4. Separate product ordering from exercise policy
Datastream supplies event metadata, including identity and source-dependent ordering information. A custom consumer must not deduplicate every change to a key as though all were one delivery. An update and a deletion of the same row are different legitimate events. In BigQuery CDC, custom sequence values can order changes to the same key; hexadecimal sections are compared numerically and in order. Equal sequences are resolved by ingestion time. Mixing writes with and without the custom sequence is not a predictable policy. The lab below uses a different, deliberately more conservative rule. Each key has an integer version invented for the exercise, without claiming to represent LSN, SCN or the BigQuery format. Different contents for the same key and version block that key, even when the conflict predates the highest version. This choice exposes disagreement for investigation. Do not present it as a vendor guarantee. When adapting the reasoning to a pipeline, define who produces the version, how it is compared and what happens when two sources disagree.
5. Run local reconciliation
The exercise accepts baseline, events and expected. Every baseline and reference row has key, version, deleted and value. A deletion retains its version and uses value=null; an active row uses an integer. Events add id and op, either upsert or delete. Each list has a maximum of one thousand entries. The program validates the complete input, rejects repeated keys in state lists and rejects reuse of an event id with different content. Two identical deliveries count as a duplicate delivery without creating a new change. In case.json, A starts at version 4 with value 100; B is deleted at version 7. Versions 6 and 5 of A arrive, followed by version 6 of B, a new key C and a repeat of e1. Run python3 content/labs/pde-cdc-reconcile/run.py < content/labs/pde-cdc-reconcile/case.json. A ends at 120 in version 6, B remains deleted and C is 30. The result identifies one duplicate delivery and matches the declared reference. That reference is supplied by the exercise user; the program does not confirm that it represents the real source.
6. Try to disprove the result
Change A’s reference value to 80 and C’s to 70. The total remains 150, but per-key comparison should produce two differences. Then add another event for A, version 5, with value 111: the key becomes blocked by a conflict despite having a version 6. Finally, add version 8 of B with value 55. Under this contract, a later change can reactivate the row; retaining a tombstone does not mean permanently preventing every new version. Reorder the events and confirm that states and conflicts remain equal. This property follows from version selection and conflict collection, rather than arrival order. Remove a key from the reference to observe unexpected-state; add a nonexistent key to obtain missing-state. The result also compares versions even when values match. None of these tests proves atomicity of a transaction across several keys, durability after restart or completeness of the entire history. State exists only during this invocation. Those limits must accompany any conclusion presented in the rehearsal.
7. Coordinate snapshot, replay and sink
A Dataflow snapshot can preserve job state and, when configured, Pub/Sub source snapshots. It does not include an automatic sink snapshot. If a table received results after the saved point, restoring the pipeline does not undo those results. The plan must control overlap, for example through idempotent writing demonstrated for the contract or an isolated destination reconciled before promotion. Do not declare that guarantee merely because message delivery has an option named exactly-once. Check region, new-version compatibility and the validity of every associated snapshot. Source validity can end before Dataflow-state validity. Seek also affects delivery control: the seeking subscription’s filter limits redelivery and the dead-letter delivery-attempt count is reset. Preserve historical evidence outside those counters. During the rehearsal, explain which consumers remain active, how earlier work is distinguished from repeated work and how an old acknowledgment that is no longer valid after seek is recognized.
8. Hand RUN a verifiable decision
Handover should connect the incident to a small set of verifiable evidence: affected interval, restart position, pipeline version, restoration resources and destination differences. In an international setting, prepare a short English note containing confirmed state, remaining uncertainty and the next owner. Avoid writing recovered without saying whether that means an active process, reconciled current state or complete history. Each meaning requires different evidence and can have a different completion time. In BigQuery CDC, max_staleness tolerance can permit reading an earlier baseline. Choose a comparable reconciliation cut and verify change application before interpreting differences as loss. Final review should include deletions, conflicts, unexpected keys and values that cancel in global totals. A green lab report means only a match with the supplied reference and absence of modeled conflicts. Summarize the learning this way: recover reading, control ordering, preserve deletion evidence, reconcile the destination and separate current state from history. Production authorization remains part of the organization’s actual process. A handover note can state: current state reconciled at the agreed cut; intraday history still under investigation; replay limited to the rehearsal destination; next owner identified. Add links to queries and results, script versions and collection time. Another person should be able to repeat the comparison without relying on the incident responder’s memory. If an assumption changes, such as available retention or source continuity, revisit the decision and record the new evidence. This sequence helps distinguish a temporary mitigation from a recovery accepted by the business. It also provides a concrete agenda for the next support handover.
"""Original bounded versioned-row teaching model; no cloud or durable writes."""
import json
import re
import sys
def exact(value, fields):
if type(value) is not dict or set(value)!= set(fields):
raise ValueError('missing or unexpected fields')
def identifier(value):
if type(value) is not str or not re.fullmatch(r'[A-Za-z0-9_-]{1,48}', value):
raise ValueError('invalid identifier')
return value
def integer(value, lo, hi):
if type(value) is not int or not lo <= value <= hi:
raise ValueError('integer outside contract')
return value
def row(value):
exact(value, ['key', 'version', 'deleted', 'value'])
identifier(value['key'])
integer(value['version'], 0, 10**9)
if type(value['deleted']) is not bool:
raise ValueError('deleted must be boolean')
if value['deleted']:
if value['value'] is not None:
raise ValueError('tombstone value must be null')
else:
integer(value['value'], -10**9, 10**9)
return dict(value)
def row_list(values):
if type(values) is not list or len(values) > 1000:
raise ValueError('at most 1000 rows')
result = {}
for value in values:
r = row(value)
if r['key'] in result:
raise ValueError('duplicate row key')
result[r['key']] = r
return result
def evaluate(document):
exact(document, ['baseline', 'events', 'expected'])
baseline, expected = row_list(document['baseline']), row_list(document['expected'])
events = document['events']
if type(events) is not list or len(events) > 1000:
raise ValueError('at most 1000 events')
unique, duplicate_deliveries = {}, 0
for event in events:
exact(event, ['id', 'key', 'version', 'op', 'value'])
identifier(event['id'])
if event['op'] not in ('upsert', 'delete'):
raise ValueError('unknown operation')
r = row({'key': event['key'], 'version': event['version'],
'deleted': event['op'] == 'delete', 'value': event['value']})
if event['id'] in unique:
if unique[event['id']]!= r:
raise ValueError('event identity reused with different content')
duplicate_deliveries += 1
unique[event['id']] = r
versions = {}
for r in list(baseline.values) + list(unique.values):
versions.setdefault(r['key'], {}).setdefault(r['version'], set).add((r['deleted'], r['value']))
states, conflicts = {}, []
for key, history in sorted(versions.items):
bad = sorted(v for v, values in history.items if len(values) > 1)
if bad:
conflicts.append({'key': key, 'versions': bad})
continue
version = max(history)
deleted, value = next(iter(history[version]))
states[key] = {'key': key, 'version': version, 'deleted': deleted, 'value': value}
blocked = {r['key'] for r in conflicts}
differences = []
for key in sorted(set(states) | set(expected) | blocked):
if key in blocked:
reason = 'conflicting-version'
elif key not in states:
reason = 'missing-state'
elif key not in expected:
reason = 'unexpected-state'
elif states[key]!= expected[key]:
reason = 'state-mismatch'
else:
continue
differences.append({'key': key, 'reason': reason})
return {'states': [states[k] for k in sorted(states)], 'conflicts': conflicts,
'differences': differences, 'duplicateDeliveries': duplicate_deliveries,
'matchesDeclaredReference': not differences,
'referenceIndependentlyVerified': False, 'historyCompletenessProven': False,
'cloudSemanticsValidated': False, 'productionReplayApproved': False,
'stateLifetime': 'this invocation only'}
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')
print(json.dumps(evaluate(json.loads(raw, object_pairs_hook=unique_object)), sort_keys=True))
except (ValueError, TypeError, RecursionError) as error:
print(json.dumps({'error': str(error)}), file=sys.stderr)
sys.exit(2)
Fictional case: backfill restores balances, but the control team still requires intraday operations. The decision distinguishes current state, history and RUN handover.
Common pitfalls
Confusing Completed with applied data; using arrival as version; removing tombstones too early; accepting totals alone; assuming a Dataflow snapshot rewinds the sink.
Related topics: CDC continuity · Snapshots and retention · Idempotency and reconciliation
Replay must respect versions and earlier effects. Reconciled state proves neither complete history nor authorization for production handover.
Reference: Professional Data Engineer standard exam guide · Current linked standard guide (document title v4.2); edition date unconfirmed (2026-09-30 inspection)