1. Define what resuming means
During a fictional banking close, the incident is over from an infrastructure perspective: the environment responds again and the scheduler has started. The team wants to release every pending job immediately. Before doing that, identify which state each component recovered. Orchestrator history might represent 20:00 while a destination table already contains the 21:00 publication. Resuming without comparing those facts can repeat effects and consume capacity required for the next close. Prepare a list of intervals, publication identifiers, destinations and evidence. Classify each interval as confirmed, absent or unresolved. Do not rewrite the meaning of an unknown state to make the plan easier. Investigation must determine whether data is missing, only confirmation is missing from history, or publication is partial. Each case can require a different action and owner. Keep the original observation alongside the conclusion so the handover remains understandable. Also define exit criteria: reconciled data, available dependencies, a still-achievable deadline, an authorized consumer and a prepared RUN team. A finished job provides only part of that evidence. Examples in this lesson are original and fictional; they do not represent BNP Paribas procedures. The local planner explores capacity and admission order without authorizing real recovery or reproducing cloud-service schedulers. Its assumptions must remain visible when discussing its output.
2. Recover the orchestrator without repeating effects
Managed Airflow documentation, on the Cloud Composer documentation pages, warns about new runs after loading a snapshot. History returns to the captured moment. A later publication can remain in its destination while becoming unknown to scheduling. The catchup=False setting does not eliminate every such case, including the documented example of a daily DAG executed and restored on the same day. During rehearsal, create a timeline with four events: capture, execution, destination confirmation and snapshot loading. Mark what each system knows after loading. This makes visible why a retry must not rely only on restored history. The publication contract might require a stable interval identifier, verification of an existing effect and an explicit decision for divergent results. Do not invent confirmation to avoid work: retain sufficient evidence to support classification, including where confirmation was obtained and which interval it covers. Review loading itself too. The destination must exist, be in a compatible state and use the same or a later Airflow version than the snapshot. A partially completed operation can appear failed while having installed packages or changed some configuration. Read the detail and inventory actual state before retrying. Technical restoration should precede dependency validation and operational acceptance, with named owners at each transition. Record remaining uncertainty instead of silently treating partial restoration as a clean starting point.
3. Distinguish observed from available capacity
A fast rehearsal does not guarantee a fast close. In BigQuery, some slots used might come from idle capacity in other reservations. Their owner can reclaim them when needed. The plan must distinguish owned capacity from borrowing opportunities, particularly if rehearsal happened outside the busiest period. Record concurrent load alongside every runtime measurement, rather than treating elapsed time as an independent property of query text. A query’s requested parallelism is also not the number of purchased slots. Units of work can wait for capacity; waiting does not automatically become extra on-demand charging for a query assigned to a reservation. To analyze cost and performance, separate demand, allocated capacity, utilization and duration. One average can hide the precise interval during which critical work was blocked. Build comparisons around equivalent inputs and a documented concurrency profile. The exercise uses abstract units and fixed durations. If a job requests one unit and takes twenty seconds, increasing capacity from one to two does not make it twice as fast. It might permit another concurrent job. Before applying this intuition to a real service, measure concurrency, destination limits and query behavior. An optimization that accelerates reads but saturates the destination database might merely move the incident. Define the limiting resource and the evidence that would justify changing capacity.
4. Read the planner contract
The run.py file receives capacity and up to fifty independent jobs. Each job has id, priority, releaseAt, duration, units and deadline. A smaller priority number means greater urgency. releaseAt is the first permitted start time. duration includes all work you decide to represent, including validation if it is required to meet the commitment. The program does not automatically add acceptance time or infer missing operational steps. At each relevant instant it first releases resources from completed jobs. It orders available candidates by priority, then deadline and finally id. It scans that order and admits requests that fit. A large request that cannot fit can be bypassed by a smaller one; no capacity is reserved in advance. It then advances to the next release or completion. A heap keeps completion events ordered, making simulation deterministic for the same declared workload. The policy does not interrupt started work or deliberately wait for a future critical job. Nor does it split requests, estimate runtimes, create retries or calculate dependencies. A request larger than capacity appears in blockedTooLarge while other jobs are still planned. This exposes the obstacle without hiding results for schedulable jobs. Times and units are bounded integers; validation rejects extra fields and repeated identifiers. An empty list meets every declared deadline trivially, which does not demonstrate recovery of any service.
5. Run the fixture and test an alternative
From the project root, run python 3 content/labs/pde-capacity-plan/run.py < content/labs/pde-capacity-plan/case.json. Code displayed in this lesson is the same executable. It uses no network, credentials or cloud services. It reads JSON, computes the plan in memory and writes output. Before running, draw the timeline: routine starts at 0, occupies the only unit until 10 and critical becomes available at 1 with duration 1 and deadline 2. Predict the delay. Critical has higher priority, but the unit is occupied and there is no preemption. It starts only at 10 and finishes at 11. lateJobIds contains critical. This does not prove the deadlines impossible: waiting until 1, running critical from 1 to 2 and routine from 2 to 12 meets both. The program does not search for that schedule because it immediately admits available work. Reproducibility of its answer should not be confused with optimality. In a fixture copy, postpone routine’s releaseAt to 2 and compare. Explain that you changed an admission condition to represent an operational decision; you did not automatically discover an optimal schedule. Next increase capacity to 2 in the original fixture and observe concurrency. Add a job requesting three units and inspect blockedTooLarge. For every variant retain input and output, identify the changed assumption and avoid comparing numbers without explaining the decision producing them.
6. Investigate delay before choosing mitigation
The planner only calculates consequences of supplied assumptions. In a real service, runtime can change because input grew, a join multiplied rows, other jobs consumed resources or the destination became slower. Investigate equivalent executions and identify where time increased. Use the execution graph and metrics from the affected interval instead of immediately concluding more capacity must be purchased. Preserve the comparison population so a faster result is not accidentally achieved by omitting required work. BigQuery performance insights provide different clues. Input growth should lead to reviewing volumes and filters. Shuffle pressure requires examining intermediate data, operations such as JOIN and GROUP BY and overlap with other queries. Reducing data before those operations or separating workloads in time might help. Measure the change and confirm the transformation still produces the correct result. A resource improvement is useful only within the required delivery contract. In the orchestrator, check what success means too. Airflow documentation describes a successful all_done leaf task that can leave a DAG Run green despite an intermediate failure. Resource cleanup is useful but must not hide delivery-relevant failure. Define explicit checks for critical tasks and reconciled outputs. In the incident report separate observed signal, hypothesis, test and result; this makes review by another team easier and prevents a tentative explanation becoming an accepted fact without evidence.
7. Stop processing without losing effect history
A stop decision must distinguish pipeline mode, observed state and destination effects. In Dataflow, drain applies to streaming and allows buffered data to finish while new ingestion stops. It is not a batch drain option. Cancellation can leave already written data accessible in the sink and can lose in-flight work. A stop request alone confirms neither rollback nor completion. Record the terminal state when it is observed rather than inferring it from request acceptance. Window closure needs attention. Drain can close a window before its usual time boundary and produce a partial result. If the new job writes another segment under the same filename based only on the hour, a destination collision can occur. Prepare an example with an hourly window and a mid-window stop; identify the segments and the rule combining them without overwriting valid data. The solution depends on the sink contract and needs its own rehearsal. Do not promise that drain fixes a stuck pipeline either. Timers that reschedule indefinitely can prevent completion. Collect state, identify behavior keeping work pending and assess a correction or another stop method with explicit impact. Communicate the decision together with what was processed, what remains unconfirmed and the restart plan. Preserve relationships among job, interval and artifacts so the next team does not have to reconstruct everything during the incident.
8. Deliver a plan RUN can execute
Finish the lesson with a one-page restart plan. List jobs that actually need to run, publications already present, dependencies still requiring validation and deadline commitments. For each job identify the estimate’s source and the conditions under which it was measured. If the estimate depends on borrowed capacity or an uncontended destination, write that condition beside the value. Assign an owner to revisit assumptions when workload changes. Attach planner input and output as discussion material. Explain that units are not BigQuery slots and local policy does not reproduce Airflow. If the local plan fails, revisit admission and assumptions; do not present the result as mathematical proof of impossibility. If it passes, verify real resources, runtimes and effects of mechanisms omitted by the model. Any safety margin must be explicit and agreed, including how validation time is represented and which deadline it protects. Define a decision point before each publication: required evidence, owner, an alternative if delayed and a condition for stopping restart. It might be acceptable to defer noncritical work or retain a previously approved publication, provided business owners accept freshness and impact. Do not hide a known delay behind green technical status. Summarize RUN handover with identifiers, blocking signals and a way to check whether an action already produced an effect before repeating it.
Prepare an attempt and manage time
Choose one of the three mocks and reserve 120 uninterrupted minutes. Each form contains 50 original course questions. The forms share no questions with each other, but you might recognize questions from lessons or earlier attempts. Record that familiarity: remembering the correct option can increase accuracy without showing that you can solve a new situation. The average of 144 seconds per question helps manage the session; it is not a mandatory duration for each decision. First read the requirements, answer when you can justify the choice, and flag uncertainties for review. Avoid using all your remaining time on one question while leaving several unread. Before comparing alternatives, identify the source, destination, execution identity and time boundary. In a pipeline, event occurrence can differ from publication, arrival and processing. During restoration, recovered orchestrator state can differ from effects already persisted at the destination. State the required guarantee to yourself: preserve valid movements, restrict access, meet a deadline or recover without repeating an external effect. The appropriate alternative must address that guarantee under the stated conditions. A familiar technology is insufficient justification. When asked for two answers, select both; on this platform, an incomplete selection or one with extra options earns no partial credit. Each fully correct answer earns one local point, up to 50. This training rule does not describe official scoring.
Separate recovery, reconciliation and publication
Consider a fictional case outside the scored questions. A pipeline prepares positions for a closing report. The orchestrator was restored to 20:00, but the destination contains results published at 20:15. The notification application confirms that some notices were sent. At 20:30, the team receives authorization to recover the service. Its plan proposes repeating everything since 20:00 and treating a green DAG status as acceptance. Begin by separating three questions: which data was calculated, which results were published, and which external effects occurred? Authorization permits execution of the approved plan; it does not establish that these states have been reconciled. A supported response identifies the population and comparison boundary. Compare relevant destination keys and values with the approved reference at the same logical point in time. Equal global totals can hide movements assigned to the wrong portfolio. Retain evidence of notices already sent and define how to prevent repetition before enabling that effect again. Calculate the available margin for required work using concurrency and observed capacity, rather than assuming the rehearsal duration is guaranteed. Also validate access with the identity that will consume the result. If reconciliation cannot finish before the committee meets, communicate the gap, its impact and the owner of the decision to postpone or limit publication. Apply this separation in the mocks without adding requirements absent from the question. If it asks for the next check, a complete architectural solution may go beyond what the evidence supports. If a specific discrepancy already exists, changing the status report does not repair the data. Look for the relationship between the symptom, likely cause, action and expected evidence.
Turn mistakes into a practice plan
After the attempt, review the explanation for the correct answer and every rejected alternative. For each mistake, record the domain, task, failed assumption and next practice action. Avoid writing only the product name. For example, if you confused the user identity with the worker identity, draw who launches the job and who reads the object; identify the permission needed at each connection. If you confused a cumulative result with an increment, construct two successive results and calculate the effect of both interpretations at the consumer. If you accepted a benchmark using cached data, define a comparison that separates result reuse from SQL execution. Return to the corresponding lesson and reference. Where a local exercise exists, run both a valid case and a condition that violates its contract. Explain the outcome before consulting the solution. These exercises use bounded models and fictional data; passing one does not validate a Google Cloud configuration. In an appropriate practice environment, separately identify the observation needed to confirm actual behavior. Then compare error types across forms. A useful improvement is being able to justify why your previous alternative failed and when it could be appropriate. Each form represents the guide’s 19 tasks and approximates the five domain weights with 11, 12, 10, 8 and 9 questions. This does not establish exhaustive subtopic coverage or equivalent psychometric difficulty. The published official exam languages are English and Japanese. Local accuracy has no official passing threshold, predicts no exam result and awards no certification. Use it to guide study and practice while keeping independent specialist review recorded as pending.
"""Bounded fictional, non-preemptive scheduler. Not a cloud scheduler or optimizer."""
import heapq
import json
import re
import sys
def keys(value, expected):
if not isinstance(value, dict) or set(value)!= set(expected):
raise ValueError('Unexpected object fields')
def integer(value, lo, hi):
if type(value) is not int or not lo <= value <= hi:
raise ValueError('Integer outside contract')
def plan(data):
keys(data, ['capacity', 'jobs'])
capacity = data['capacity']; integer(capacity, 1, 100)
if not isinstance(data['jobs'], list) or len(data['jobs']) > 50:
raise ValueError('Expected at most50jobs')
seen = set
for job in data['jobs']:
keys(job, ['id', 'priority', 'releaseAt', 'duration', 'units', 'deadline'])
if not isinstance(job['id'], str) or re.fullmatch(r'[A-Za-z0-9_-]{1,48}', job['id']) is None:
raise ValueError('Invalid job identifier')
if job['id'] in seen:
raise ValueError('Duplicate job identifier')
seen.add(job['id'])
for field, lo, hi in [('priority',0,10),('releaseAt',0,86400),
('duration',1,86400),('units',1,100),('deadline',0,172800)]:
integer(job[field],lo,hi)
blocked = sorted(j['id'] for j in data['jobs'] if j['units'] > capacity)
pending = [dict(j) for j in data['jobs'] if j['units'] <= capacity]
running = []; scheduled = []; now = 0; free = capacity; peak = 0
while pending or running:
while running and running[0][0] <= now:
_, _, units = heapq.heappop(running); free += units
ready = sorted((j for j in pending if j['releaseAt'] <= now),
key=lambda j:(j['priority'],j['deadline'],j['id']))
for job in ready:
if job['units'] > free:
continue
finish = now + job['duration']; free -= job['units']
heapq.heappush(running,(finish,job['id'],job['units']))
scheduled.append(dict(id=job['id'],start=now,finish=finish,
units=job['units'],wait=now-job['releaseAt'],
deadline=job['deadline'],deadlineMet=finish<=job['deadline']))
pending.remove(job); peak=max(peak,capacity-free)
events = [j['releaseAt'] for j in pending if j['releaseAt'] > now]
if running:
events.append(running[0][0])
if events:
now = min(events)
elif pending:
raise AssertionError('A fitting pending job must eventually be admitted')
scheduled.sort(key=lambda j:j['id'])
return dict(schedule=scheduled,blockedTooLarge=blocked,
lateJobIds=sorted(j['id'] for j in scheduled if not j['deadlineMet']),
peakUnits=peak,makespan=max((j['finish'] for j in scheduled),default=0),
allDeclaredDeadlinesMet=not blocked and all(j['deadlineMet'] for j in scheduled),
optimalityProven=False,cloudScheduleValidated=False,
runtimeEstimatesValidated=False,productionRecoveryApproved=False,
policy='ready priority,deadline,id; admit fitting jobs; no preemption or future reservation')
def unique_object(pairs):
obj={}
for k,v in pairs:
if k in obj:raise ValueError('Duplicate JSON key')
obj[k]=v
return obj
if __name__=='__main__':
try:
raw=sys.stdin.buffer.read(2000001)
if len(raw)>2000000:raise ValueError('Input exceeds2MB')
print(json.dumps(plan(json.loads(raw,object_pairs_hook=unique_object)),sort_keys=True))
except (ValueError,TypeError,UnicodeError,RecursionError) as exc:
print(json.dumps({'error':str(exc)}),file=sys.stderr);sys.exit(2)
Fictional close: restored history misses an existing publication and a routine job occupies capacity before a critical job arrives.
Common pitfalls
Repeating effects after restore; promising borrowed capacity; confusing priority with preemption; assuming planner optimality; accepting green cleanup as delivery proof; ignoring partial windows.
Related topics: History and idempotency · Capacity and deadlines · Streaming shutdown
Reconcile before rerunning, state admission policy and observe effects and completion before confirming recovery.
Reference: Professional Data Engineer standard exam guide · Current linked standard guide (document title v4.2); edition date unconfirmed (2026-09-30 inspection)