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

Operate pipelines, capacity and recovery

Connect cost, automation and observability to delivery deadlines; calculate net capacity and distinguish technical recovery from accepted data.

1. Operate the outcome the business uses

A pipeline is not recovered simply because its process is active. Define the outcome the consumer needs, its deadline and acceptance criteria. In our fictional APS case, the service delivers reconciled batches for reporting. A late batch or one with discrepancies does not count as the same accepted outcome. Record the business unit, window, consumers, acceptance owner and handling of provisional delivery. These examples do not represent BNP Paribas procedures. Compare cost and outcome over the same population. Moving from 120 cost units to 90 looks like improvement, but delivering only six accepted batches instead of twelve changes cost per outcome from 10 to 15. Investigate lost deliveries before presenting the change as efficiency. For RUN handover, connect each signal to a decision: who receives the alert, the first diagnostic action and when to escalate. Acting does not require a perfect metric; it requires carrying evidence limitations into the decision. A capacity indicator helps forecast; reconciliation confirms delivered content. Keep those observations separately traceable so a service recovery report can explain both timing and correctness.

2. Reduce waste while retaining control

Size compute from workload and useful execution window. A persistent cluster may serve frequent, predictable work; an isolated run may justify compute created for that work. In a Spark workflow with a managed cluster, prepare output to survive cluster completion. If the only file is on a worker’s local disk, shortening cluster lifetime can destroy the result the consumer needs. Include durable publication and validation in the workflow and measure complete delivery. Current documentation presents Managed Service for Apache Spark and Managed Service for Apache Airflow; the exam guide still uses Dataproc and Cloud Composer. Preserve that mapping without inferring a new exam edition. For financial control, distinguish Alerts-only budgets from Spend cap. Without additional automation, the former communicates deviation; it does not itself suspend usage. Always confirm the configured type and scope. Also validate the population of a usage report: summing a SCRIPT job and its children can duplicate aggregate values. A cost difference should lead to investigation of workload, configuration and measurement unit. Do not attribute it to one infrastructure change until those comparison conditions are understood.

3. Automate repeats with stable meaning

Design a task to repeat the same logical unit. Publication may have completed at the destination even if the orchestrator lost acknowledgment. Repeating an append under a new identity then turns recovery into duplication. Retain stable batch identity, persisted outcome and transformation version. The destination contract must allow recognition of the previous outcome or an idempotent operation. A timeout alone establishes neither success nor absence of commit. Control starts before execution. Network calls made while loading a DAG file can repeatedly delay its processing. Separate definition and execution, using prepared configuration or task-time queries where appropriate. Test dependencies: transform must not begin before extract has the required output. Textual order in a template is not that guarantee. Finally, fixing a template does not establish that an already active run received the fix. Record the affected instance’s version and the action applied to it. Recovery evidence must identify what ran, with which inputs and what outcome was published. Keep uncertain commits visible until destination evidence resolves them; a scheduler status alone may not answer the business question.

4. Connect capacity to the business deadline

Organize workloads by criticality, concurrency and acceptable cost. BigQuery reservations and assignments allow capacity management for different workload populations. Dataset names, labels and service accounts serve other purposes; alone they do not establish dedicated capacity. Validate the sizing and sharing behavior used in the environment. An undersized reservation can still miss the consumer’s objective. The deadline includes waiting, execution, publication and validation. Do not use only the duration of a previous run that started immediately. When a batch job is queued, estimate remaining margin and identify which change can actually affect it. Reassigning a project from reservation A to B does not retroactively move requests already queued or running in A; new requests go to B. This distinction prevents announcing mitigation that has not reached the affected work. In the incident report, separate old jobs from new ones and observe both. A priority or capacity change remains an improvement hypothesis until there is evidence of timely delivery with accepted content. Record the time of configuration changes so the job population affected by each decision remains clear.

5. Investigate signals before adding resources

Begin with symptoms and work distribution. Growing backlog with every worker busy calls for a different investigation from growing backlog with only one saturated worker. A frequent key can limit parallelism. In that case, evaluate redistribution compatible with required aggregation and ordering; adding workers does not automatically create useful capacity for that key. Confirm the hypothesis using stage and worker data, preventing an average from hiding the problem. Read error reasons. quotaExceeded does not universally mean insufficient slots: identify the quota, scope and affected operation before choosing a response. Also distinguish data freshness, processing latency and pending volume. None of these measures alone establishes that the destination received everything. If a series stops arriving, record an observation gap. Filling the graph with zero can hide that gap. Check missing-data handling in the alert policy and seek independent destination evidence. Automatic alert closure must not be used as certification of functional recovery. Keep a timeline connecting source arrivals, processing changes, monitoring gaps and destination observations; it helps separate a service change from a change in what the team can observe.

6. Lab: backlog, net capacity and sample age

Run python3 content/labs/pde-backlog-budget/run.py < content/labs/pde-backlog-budget/case.json from the project root. The original program uses local Python only. now and deadline are integer seconds on a fictional scale. maxSampleAge is the inclusive maximum accepted age. stages contains one to one hundred unique stages; samples accepts up to one thousand observations with id, stage, at, backlog, inRate and outRate. It uses no credentials and queries no cloud services. For each stage, select the latest instant not exceeding now. Without a sample, return unknown/missing; with excessive age, unknown/stale. Samples at that instant with conflicting values produce unknown/conflicting. Equal values retain all IDs as provenance. Define net=outRate−inRate and project backlogNow=max(0,backlog−net×age), assuming constant rates from sample time through deadline. within-model requires net>=0 and backlogNow<=net×(deadline−now). Integer comparison includes the exact boundary. An empty queue with arrivals exceeding departures remains at risk because it will grow again. The rounded seconds estimate only presents the same assumption. These rules are an explicit local model, not the implementation of a managed service’s backlog metric.

7. Interpret the exercise and challenge the forecast

In the fixture, now=1000 and deadline=1060. ingest has backlog 1200 and rates 80/100: net capacity 20 drains it in 60 seconds, exactly at the boundary. fold has the same backlog and rates 100/110: it needs 120 seconds and becomes risk. export has a sample aged 20 when the maximum is 10, so it remains unknown. A future fold sample is ignored. The overall result retains both the risk list and the unknown-stage list. Reduce the deadline by one second and observe ingest fail. Move the ingest sample to 990: projection subtracts ten seconds of draining without claiming a new measurement was taken. Create conflicting samples at the same instant and confirm that changing order does not select either one. Use backlog 3, arrivals 0, departures 2 and one available second to check rounding. Even if all stages become within-model, the end-to-end deadline is unproven: dependencies, variable rates, publication and reconciliation lie outside the model. The program keeps those limitations explicit. Explain which observation would be needed to replace each assumption before using a similar estimate in a real incident.

8. Recover data and restore RUN autonomy

Choose recovery according to failure type. Cloud SQL HA across zones in one region covers a different scope from regional unavailability. An asynchronous replica in another region can support recovery, but promotion requires considering transactions not yet replicated. Do not turn replica existence into a zero-RPO promise. Identify available evidence about the recovered point and communicate remaining uncertainty. An incorrect logical change may already have replicated. For Cloud SQL PostgreSQL PITR, plan a new instance, validation of the chosen instant and client routing. Technical recovery precedes functional acceptance: validate connectivity, permissions, expected data and capacity for RUN workload. In handover, provide the repeatable procedure, owners, authority boundaries, evidence from the latest exercise and escalation criteria. Update business communication with delivered outcome, delay and known discrepancies. The operational summary should let a colleague continue without guessing what was observed, what was inferred and what remains unconfirmed. This closes the loop between resource management and the consumer’s original acceptance criteria; a reachable replacement database is one milestone in that loop, not its complete evidence.

"""Original offline fluid-queue exercise; no cloud calls or recovery approval."""
import json
import re
import sys


def integer(value, maximum):
 if type(value) is not int or not 0 <= value <= maximum:
 raise ValueError('invalid non-negative integer')
 return value


def identifier(value):
 if not isinstance(value, str) or not re.fullmatch(r'[A-Za-z0-9_-]{1,64}', value):
 raise ValueError('invalid identifier')
 return value


def evaluate(data):
 if not isinstance(data, dict) or set(data)!= {'now', 'deadline', 'maxSampleAge', 'stages', 'samples'}:
 raise ValueError('invalid input shape')
 now = integer(data['now'], 10**9)
 deadline = integer(data['deadline'], 10**9)
 max_age = integer(data['maxSampleAge'], 10**6)
 if deadline < now:
 raise ValueError('deadline precedes evaluation')
 stages, samples = data['stages'], data['samples']
 if not isinstance(stages, list) or not 1 <= len(stages) <= 100:
 raise ValueError('stages must contain 1..100 identifiers')
 if not isinstance(samples, list) or len(samples) > 1000:
 raise ValueError('too many samples')
 for stage in stages:
 identifier(stage)
 if len(set(stages))!= len(stages):
 raise ValueError('duplicate stage')
 ids = set
 for row in samples:
 if not isinstance(row, dict) or set(row)!= {'id', 'stage', 'at', 'backlog', 'inRate', 'outRate'}:
 raise ValueError('invalid sample shape')
 identifier(row['id'])
 if row['id'] in ids:
 raise ValueError('duplicate sample identifier')
 ids.add(row['id'])
 identifier(row['stage'])
 if row['stage'] not in stages:
 raise ValueError('unknown stage')
 integer(row['at'], 10**9)
 for field in ['backlog', 'inRate', 'outRate']:
 integer(row[field], 10**6)
 results = []
 for stage in sorted(stages):
 rows = [r for r in samples if r['stage'] == stage and r['at'] <= now]
 result = {'stage': stage, 'status': 'unknown', 'reason': 'missing',
 'sampleIds': [], 'sampleAt': None, 'age': None,
 'projectedBacklogNow': None, 'netDrainRate': None, 'drainSecondsCeil': None}
 if rows:
 at = max(r['at'] for r in rows)
 latest = [r for r in rows if r['at'] == at]
 result.update(sampleIds=sorted(r['id'] for r in latest), sampleAt=at, age=now-at)
 if now - at > max_age:
 result['reason'] = 'stale'
 elif len({(r['backlog'], r['inRate'], r['outRate']) for r in latest}) > 1:
 result['reason'] = 'conflicting'
 else:
 row = latest[0]
 net = row['outRate'] - row['inRate']
 backlog_now = max(0, row['backlog'] - net * (now-at))
 seconds = ((backlog_now + net - 1) // net if net > 0
 else 0 if net == 0 and backlog_now == 0 else None)
 fits = net >= 0 and backlog_now <= net * (deadline-now)
 reason = ('within-window' if fits else 'growing' if net < 0
 else 'no-drain' if net == 0 else 'insufficient-window')
 result.update(status='within-model' if fits else 'risk', reason=reason,
 projectedBacklogNow=backlog_now, netDrainRate=net,
 drainSecondsCeil=seconds)
 results.append(result)
 risk = [r['stage'] for r in results if r['status'] == 'risk']
 unknown = [r['stage'] for r in results if r['status'] == 'unknown']
 return {'results': results, 'riskStages': risk, 'unknownStages': unknown,
 'status': 'risk' if risk else 'unknown' if unknown else 'within-model',
 'endToEndDeadlineProven': False, 'dataCompletenessProven': False,
 'cloudCapacityMeasured': False, 'productionApproval': 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(1000001)
 if len(raw) > 1000000:
 raise ValueError('input too large')
 payload = json.loads(raw, object_pairs_hook=unique_object)
 print(json.dumps(evaluate(payload), sort_keys=True))
 except (ValueError, TypeError) as error:
 print('invalid input: ' + str(error), file=sys.stderr)
 sys.exit(2)
IN PRACTICE

Backlog 1200, arrivals 100/s and departures 110/s require 120 seconds; a 60-second window is insufficient.

Common pitfalls

Gross output as drain; absence as zero; new template as fixed execution; zonal HA as regional DR; alert as spending cap.

Related topics: Ingestion and replay · Fidelity and reconciliation · Operational continuity

Take this idea with you

A forecast depends on rates and evidence; confirm accepted delivery before declaring recovery.

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.