← Professional Cloud Architect: architecture and operations
15 / 25 · 120 MIN

Migration, reconciliation and analytical results

Plan data cutoffs, preserve meaning and demonstrate acceptance criteria through identity-based reconciliation and interpretation of analytical results.

Define what must arrive and when

A migration needs a verifiable content definition, not just a calendar date. Specify tables or files, filters, history, later changes and transformation rules. In a fictional fund-reporting project, the data owner may require every operation accepted by 22:00, including a batch that started earlier. That statement requires identifying producers and demonstrating the cutoff. Size the transfer as well: 10 decimal TB at an effective 500 Mbit/s requires at least 160,000 seconds, about 44.4 hours, before validation and margin. If changes arrive at 25 MB/s but only 20 MB/s is applied, backlog grows by 5 MB/s. Rehearsal should measure extraction, transport and application; upgrading only the link may leave the actual bottleneck untouched. Document units and assumptions so the steering committee can distinguish a measured forecast from a hoped-for completion date.

Compare states at a demonstrable boundary

The initial copy provides a migration baseline; subsequent changes must be captured by the chosen mechanism. A source query at 10:00 and a target query at 10:05 may observe legitimately different states. That does not justify classifying every difference as lag or corruption. Retain results and establish a comparable cutoff through system capabilities and a rehearsed procedure. In this lesson’s active-passive case, the target remains closed to business writes until producers are controlled, accepted work finishes, changes are applied and reconciliation completes. A CFT file still processing belongs to that boundary even if the web interface has already stopped. Once the target accepts writes, the old source is no longer an automatically equivalent fallback. Record how those new operations will be preserved before promising a return. A name such as cut-9 in a report is merely a label; evidence must connect it to the observed state.

Preserve meaning during transformation

Write a comparison contract before normalizing values. In the exercise, IDs are text: 001 and 1 differ. A missing reference and a received empty reference also differ; do not convert both to empty text merely to make a check pass. Timestamps represent instants, so 11:00+01:00 and 10:00Z are equivalent after interpreting the offset. A time without a timezone is rejected because there is insufficient information to choose the instant. This simulation’s amounts have exactly two decimal places, a local assumption that does not apply to every currency. In a BigQuery architecture, choose a decimal type with suitable precision and scale and test its boundaries. SAFE_CAST can turn invalid input into NULL; record rejections separately instead of trusting a sum that observes only valid values. Every cleansing rule needs business justification and a version: trimming whitespace or ignoring case may change identifiers that appear to be ordinary text.

Combine checks that observe different failures

A transport checksum checks byte integrity; it does not confirm that a later transformation preserves meaning. Count, sum, uniqueness and identity-based comparison answer different questions. If A changes from 100 to 110 and B from 200 to 190, count and sum stay equal. If one ID disappears and another appears with the same value, even currency-level aggregates may match. Start by checking scope and keys, compare relevant values and use aggregates as supporting evidence. In Data Validation Tool, choose a key identifying one result row; trade_id may need leg_id. Inspect normalization in generated queries. A --dry-run displays SQL and does not demonstrate data equality. For a large dataset, plan coverage by partitions or ranges and retain which were actually checked. A sample with no discrepancies does not establish correctness of the whole dataset.

Exercise: find differences hidden by totals

Before running the code, predict three outcomes: reorder the rows, replace 001 with 1, and redistribute one euro between two IDs while preserving the sum. Run python3 content/labs/pca-data-reconciliation/run.py. Observed checks should show equivalence in the first case and differences in the other two. The canonical function validates the small schema, converts decimal amounts to integer cents and normalizes instants to UTC with precision up to microseconds. The index function rejects repeated IDs rather than silently removing rows. compare identifies missing and unexpected IDs and changed fields while keeping currency totals separate. Change one reference from None to empty text and inspect the reported field. The program compares only synthetic in-memory lists; it does not run DVT, capture snapshots or contact Google Cloud. Two empty lists with the same label pass, but that does not establish complete extraction. At work, also retain scope and cutoff evidence with appropriate access to results.

Design the meaning of a streaming result

A newer source update can reach a custom Datastream file consumer before an older update. Arrival order should not replace source-appropriate sequence metadata. Do not assign managed-destination merge behavior to your own consumer. If you then aggregate events, decide which timestamp represents the business: a payment at 20:59 arriving at 21:02 belongs to the occurrence window when that is the contract. The watermark helps determine when to emit, but late data can still arrive. Define what can be corrected and for how long. With accumulating panes, emitting 100 and then 130 for the same window does not mean producing 230. The destination must recognize revisions and prevent an older revision from replacing a newer one. These details affect the dashboard operations teams use for decisions; a fast result without a defined completeness meaning can lead to a wrong decision.

Align queries and estimates with use

Choose data organization from real queries. If reporting filters business date and desk, trial date partitioning and desk clustering within partitions. Measure the effect with representative data; a feature name does not guarantee savings. On a native nonclustered table, SELECT * LIMIT 10 does not replace a partition filter or column selection. Confirm that the predicate supports pruning and observe estimated and actual bytes. Likewise, do not interpret zero bytes in an external-table dry run as a zero-cost promise: it may merely be the lower bound the estimator can provide. If you declare BigQuery primary keys, keep data consistent with them because they are not automatically enforced and can influence optimizations. Architecture should identify who validates these properties when new data arrives, not only who created the table during the project.

Close the decision with criteria and owners

Prepare an acceptance table with requirement, observation, result, owner and outstanding action. In the three-trade case, state explicitly that totals match but T2 differs, T3 is missing and T9 is unexpected. The decision follows the approved per-ID equivalence rule, so acceptance remains blocked. An exception requires assessing impact and approval authority; it should not emerge from changing the verifier to obtain green results. For RUN handover, include how to repeat reconciliation, where to observe rejected inputs, when a report is provisional and how to escalate a discrepancy. If daily batch already meets the measured need, use it as the baseline and require a concrete benefit before adding streaming. Summarize the lesson’s principle: demonstrating transport, state, meaning and use requires different evidence. Architect and project manager should connect that evidence to the decision while preserving unresolved limitations.

"""Original local reconciliation exercise. Synthetic data; no cloud or DVT calls.
Contract: unique string IDs; exact two-decimal amounts; aware instants;
case-sensitive references including NULL versus empty; same declared checkpoint.
The checkpoint is an assertion by the caller, not proof of a consistent snapshot.
"""
import json
import re
from datetime import datetime, timezone
from decimal import Decimal

FIELDS = {'id', 'currency', 'amount', 'occurred_at', 'reference'}

def canonical(row):
 if not isinstance(row, dict) or set(row)!= FIELDS:
 raise ValueError('exact exercise schema required')
 if not isinstance(row['id'], str) or not row['id']:
 raise ValueError('nonempty string ID required')
 if not isinstance(row['currency'], str) or not re.fullmatch(r'[A-Z]{3}', row['currency']):
 raise ValueError('three uppercase letters required; not an ISO currency registry check')
 amount = row['amount']
 if not isinstance(amount, str) or not re.fullmatch(r'[+-]?\d{1,12}(?:\.\d{1,2})?', amount, flags=re.ASCII):
 raise ValueError('bounded exact decimal string required')
 cents = int(Decimal(amount) * 100)
 instant = row['occurred_at']
 if not isinstance(instant, str) or not re.fullmatch(r'\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}(?:\.\d{1,6})?(?:Z|[+-]\d{2}:\d{2})', instant, flags=re.ASCII):
 raise ValueError('timestamp string required')
 try:
 parsed = datetime.fromisoformat(instant)
 except ValueError as exc:
 raise ValueError('invalid ISO timestamp') from exc
 if parsed.tzinfo is None or parsed.utcoffset is None:
 raise ValueError('explicit UTC offset required')
 ref = row['reference']
 if ref is not None and not isinstance(ref, str):
 raise ValueError('reference must be a string or None')
 return {'id': row['id'], 'currency': row['currency'], 'cents': cents,
 'instant': parsed.astimezone(timezone.utc).isoformat(timespec='microseconds'),
 'reference': ref}

def index(rows):
 result = {}
 for row in rows:
 item = canonical(row)
 if item['id'] in result:
 raise ValueError('duplicate ID: ' + item['id'])
 result[item['id']] = item
 return result

def totals(items):
 result = {}
 for row in items.values:
 bucket = result.setdefault(row['currency'], {'count': 0, 'cents': 0})
 bucket['count'] += 1
 bucket['cents'] += row['cents']
 return result

def compare(source, target, source_checkpoint, target_checkpoint):
 if not isinstance(source_checkpoint, str) or not source_checkpoint or source_checkpoint!= target_checkpoint:
 raise ValueError('same nonempty declared checkpoint required')
 left, right = index(source), index(target)
 missing = sorted(left.keys - right.keys)
 unexpected = sorted(right.keys - left.keys)
 changed = {key: sorted(k for k in left[key] if left[key][k]!= right[key][k])
 for key in sorted(left.keys & right.keys) if left[key]!= right[key]}
 return {'same': not (missing or unexpected or changed), 'missing': missing,
 'unexpected': unexpected, 'changed': changed,
 'sourceTotals': totals(left), 'targetTotals': totals(right)}

def sample(key='001', amount='12.30', **changes):
 return dict(id=key, currency='EUR', amount=amount,
 occurred_at='2026-10-05T10:00:00+00:00', reference=None) | changes

def run:
 checks = []
 def check(name, condition):
 if not condition:
 raise AssertionError(name)
 checks.append(name)
 def reject(name, fn):
 try:
 fn
 except ValueError:
 checks.append(name)
 else:
 raise AssertionError(name)
 def same(a, b): return compare(a, b, 'cut-9', 'cut-9')
 a = sample
 check('decimal amount preserved', canonical(a)['cents'] == 1230)
 check('negative exact amount', canonical(sample(amount='-0.01'))['cents'] == -1)
 check('equivalent decimal spelling', canonical(sample(amount='12.3')) == canonical(a))
 check('same instant different offset', canonical(sample(occurred_at='2026-10-05T11:00:00+01:00')) == canonical(a))
 check('UTC Z accepted', canonical(sample(occurred_at='2026-10-05T10:00:00Z')) == canonical(a))
 check('identity retained as text', canonical(a)['id'] == '001')
 check('leading zeros remain significant', same([a], [sample(key='1')])['missing'] == ['001'])
 check('NULL differs from empty', same([a], [sample(reference='')])['changed'] == {'001': ['reference']})
 check('reference case significant', not same([sample(reference='ab')], [sample(reference='AB')])['same'])
 check('reference whitespace significant', not same([sample(reference='ab')], [sample(reference='ab ')])['same'])
 check('same rows reordered', same([a, sample('002')], [sample('002'), a])['same'])
 check('missing row identified', same([a, sample('002')], [a])['missing'] == ['002'])
 check('unexpected row identified', same([a], [a, sample('003')])['unexpected'] == ['003'])
 left = [sample('001', '10.00'), sample('002', '20.00')]
 right = [sample('001', '11.00'), sample('002', '19.00')]
 diff = same(left, right)
 check('equal count and sum hide changes', diff['sourceTotals'] == diff['targetTotals'] and not diff['same'])
 check('both cancelling changes identified', diff['changed'] == {'001': ['cents'], '002': ['cents']})
 replacement = same([sample('001', '10.00')], [sample('009', '10.00')])
 check('replacement identity not equivalent', replacement['missing'] == ['001'] and replacement['unexpected'] == ['009'])
 mixed = index([sample('001', '10.00'), sample('002', '10.00', currency='USD')])
 check('currencies never added together', totals(mixed) == {'EUR': {'count': 1, 'cents': 1000}, 'USD': {'count': 1, 'cents': 1000}})
 check('empty pair equivalent for supplied scope', same([], [])['same'])
 check('microseconds retained', not same([a], [sample(occurred_at='2026-10-05T10:00:00.000001Z')])['same'])
 check('Unicode preserved', same([sample(reference='ação')], [sample(reference='ação')])['same'])
 reject('duplicate source rejected', lambda: same([a, a], [a]))
 reject('duplicate target rejected', lambda: same([a], [a, a]))
 reject('different checkpoints rejected', lambda: compare([a], [a], 'cut-9', 'cut-10'))
 reject('empty checkpoint rejected', lambda: compare([], [], '', ''))
 reject('float amount rejected', lambda: canonical(sample(amount=12.3)))
 reject('extra precision rejected', lambda: canonical(sample(amount='12.301')))
 reject('NaN rejected', lambda: canonical(sample(amount='NaN')))
 reject('missing offset rejected', lambda: canonical(sample(occurred_at='2026-10-05T10:00:00')))
 reject('invalid timestamp rejected', lambda: canonical(sample(occurred_at='2026-99-05T10:00:00Z')))
 reject('missing field rejected', lambda: canonical({k: v for k, v in a.items if k!= 'reference'}))
 reject('numeric ID rejected', lambda: canonical(sample(key=1)))
 reject('lowercase currency rejected', lambda: canonical(sample(currency='eur')))
 reject('oversized decimal rejected', lambda: canonical(sample(amount='1000000000000.00')))
 reject('nontext reference rejected', lambda: canonical(sample(reference=0)))
 reject('extra field rejected', lambda: canonical(a | {'other': 1}))
 reject('empty ID rejected', lambda: canonical(sample(key='')))
 check('JSON retains null versus string', json.loads(json.dumps([None, ''])) == [None, ''])
 reject('submicrosecond input rejected', lambda: canonical(sample(occurred_at='2026-10-05T10:00:00.0000001Z')))
 check('input records unchanged', a == sample)
 return {'passed': len(checks), 'checks': checks, 'example': diff,
 'network': False, 'vendorExecution': False, 'persistentWrites': False}

if __name__ == '__main__':
 print(json.dumps(run, ensure_ascii=False, indent=2))
IN PRACTICE

A three-trade batch retains EUR 600 but changes one value and replaces one ID. Count and sum pass while record-level acceptance fails.

Common pitfalls

Confusing checksums with semantic equivalence, matching totals with matching rows, initial copy with final synchronization and cutoff labels with demonstrated snapshots.

Related topics: Data, consistency and event publication · Migration, costs and acceptance · Observability and completion evidence

Take this idea with you

Define scope and boundary; preserve identity and meaning; combine distinct checks and connect acceptance decisions to the evidence produced.

Create account

Reference: Transfer your large datasets · Current linked standard guide; 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.