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

Control attempts, resources and delivery deadlines

Decide whether to observe, reconcile, repair or retry a job, considering cost, concurrency and validation.

Reconstruct what happened before retrying

A fictional APS team is preparing a daily funds-position file. Extraction finished, transformation appears complete and loading lost its acknowledgement. The delivery deadline is approaching. The operator is asked to rerun everything. Before acting, identify the unit of work: business date, input partition, transformation version, destination and each remote job identifier. Without these details, repeating the operation can mean processing different data or appending rows already loaded. Separate three questions in the incident record: did the request reach the service, did the job finish and is its output correct? In BigQuery, DONE means the job is no longer running; check errorResult to determine whether it failed. Absence of that error can still coexist with nonfatal errors, which must be compared with acceptance criteria. A load that tolerated invalid rows may be technically complete yet remain incomplete for the business. Support’s objective is to recover a reconciled delivery while preserving evidence that explains each decision. Record the source and observation time of every status so that an old dashboard value cannot silently become current evidence.

Control resources before execution

An analyst proposes reducing an expensive query’s cost by adding LIMIT 10. On a non-clustered table, limiting returned rows does not by itself reduce bytes read. Examine selected columns, filters and the partitions actually required. Under on-demand pricing, maximum bytes billed can stop a query whose estimate exceeds the limit. Do not confuse that control with a duration promise or transfer it without analysis to capacity-based billing. In the example, the approved limit is 200 illustrative byte units and the estimate is 260. The appropriate action is to revise the plan or obtain the applicable budget change, not automatically increase the limit on every retry. For clustered tables, the estimate can be an upper bound and reject a query whose actual cost would be lower. Record that uncertainty when justifying an adjustment. If several steps recompute the same transformation, consider materializing an intermediate result with a version, retention and integrity check. Compare avoided processing with storage and maintenance costs; an intermediate table without an owner can create another operational problem.

Fix the meaning of another attempt

An Airflow task that always reads the latest partition can produce different results when retried. Using now to choose the business date creates the same problem. Use the data interval or another stable run identifier to select inputs and outputs. Record legitimate source changes as a new processing version rather than hiding them inside a supposedly equivalent retry. This choice makes attempts comparable and lets the team explain differences to the business. For a BigQuery load using WRITE_APPEND, retain the client-selected job ID. If submission loses its response and a retry with that ID returns duplicate, inspect the existing job. Do not immediately create a random ID to bypass the error. A new predictable ID can be appropriate after confirming that earlier attempts failed. The remote identifier, batch identifier and record key serve different purposes. The first helps manage submissions; it does not alone demonstrate that two different files contain no duplicate movements. Business reconciliation remains necessary before making the result available. Preserve the input version alongside the identifier so that a retry cannot silently change its payload.

Limit concurrency around the real dependency

Adding orchestrator workers can increase pressure on an already saturated database or API. An Airflow pool limits parallelism for a set of tasks. Assign tasks sharing the same dependency to the same control and estimate their weight. If a heavy task occupies two slots in a two-slot pool, light tasks in that pool wait. Available workers do not create additional pool slots. In the fictional service, the positions database supports either one heavy extraction or two small ones concurrently during the agreed window. This is an exercise assumption that would need measurement in the real system. Document how that capacity maps to pool_slots; do not assume a slot automatically represents one database connection. Also check whether deferred tasks count as occupied slots under the pool configuration. Releasing local resources while an external job remains active can admit additional work against the same dependency. The capacity plan must consider where work actually continues, the business deadline and queue priorities. A retry that is eligible can still have to wait for capacity; the local exercise below does not model that queue.

Attempt and backoff budgets

A retry policy has at least three decisions: which errors qualify, how many retries to allow and how long to wait between them. In Workflows, max_retries excludes the initial execution. The lab uses maxAttempts, which includes it. Thus max_retries=3 allows up to four executions, while maxAttempts=3 allows three. When translating configurations between tools, write out the attempt sequence to avoid this counting error. Choose the HTTP policy according to the step’s semantics. The default policy for non-idempotent steps covers fewer situations than the policy for idempotent steps; a timeout can leave the remote outcome unknown. The lab is deliberately conservative: unknown requires reconciliation even when safeRepeat is true. For a retryable state, it applies capped exponential backoff and calculates the next attempt’s estimated finish, including validation. It models neither jitter, queues nor future failures. In a real service, spreading retries helps prevent many clients from returning at the same instant. Increasing the maximum without considering the deadline may simply extend an incident that already needs another mitigation.

Run the decision exercise

Run python3 run.py < case.json in the pde-retry-window lab. Times are seconds on a shared teaching clock, not real timestamps. The fixture uses asOf=100, deadline=140 and a maximum observation age of ten seconds. Policy permits four total attempts, starts with four seconds of backoff, doubles the delay and caps it at twelve. Each task declares state, number of attempts started, observation time, confirmed last-attempt completion and estimated execution and validation durations. report failed twice and its last attempt ended at 95. Backoff is eight, making the earliest next start 103. With twenty seconds of execution and five of validation, estimated finish is 128 and the decision is wait-backoff. load can start at 100 and needs 35+5 seconds: it reaches the deadline exactly and receives retry-eligible. append has an unknown outcome and receives reconcile-remote. archive exhausted its attempts. monitor’s observation is stale and needs refreshing. transform succeeded, but its output still needs validation. The program only calculates decisions; it does not wait, observe services or submit jobs. Keep its output as an exercise result, not an execution record.

Test boundaries and challenge the forecast

Reduce deadline from 140 to 139: load no longer fits. Keep 140 but advance asOf to 101: it also no longer fits because human delay consumed one second. Validation cannot be removed from the calculation merely to produce a favorable result. For report, advancing asOf to 103 completes backoff and makes the attempt eligible, provided the observation remains within its allowed age. Also test the observation boundary: a sample from 90 is accepted at 100 with a limit of ten; a sample from 89 requires refresh. A future observation time is invalid input, not fresher evidence. Change safeRepeat to false for a retryable error: the model requires effect reconciliation. That boolean is a caller assertion; the program does not prove idempotency. Tests permute all six tasks into 720 orders and require identical output because each decision is independent. This does not demonstrate that every task can run concurrently. Dependencies, capacity and queues would need to enter another model and the service’s own trials. Preserve a failing boundary example as well as the happy path.

Mitigate while preserving RUN autonomy

A BigQuery cancellation request does not guarantee that the job was cancelled; the job may already have finished. Costs can also remain. Before starting a replacement, inspect state and check effects. Distinguish cancelling running work from reversing results already published. If the original delivery no longer fits the deadline, communicate impact, affected population and an operational alternative rather than hiding the forecast behind more retries. At RUN handover, provide a procedure containing identifiers, job locations, state sources, error classification, attempt limits, validation criteria and escalation owners. Include examples of duplicate after a lost response, DONE with an error and deadline-infeasible in the exercise. Ask another operator to explain the decision without help from the author. For the international meeting, prepare a short English update: confirmed state, outcome still unknown, next action and time of the next evidence. Closure should rely on reconciled data and consumer acceptance. The lab teaches how to make a decision explicit; validation of real behavior and specialist review remain separate work. Record unresolved assumptions so that a handover does not turn them into promises.

"""Original offline retry-decision model. Does not submit, cancel or observe jobs."""
import json
import re
import sys


def require(condition, message):
 if not condition:
 raise ValueError(message)


def integer(value, low, high):
 return type(value) is int and low <= value <= high


def keys(value, fields):
 require(type(value) is dict and set(value) == set(fields.split), 'Unexpected fields')


def evaluate(payload):
 keys(payload, 'asOf deadline maxObservationAge policy tasks')
 for name in ['asOf', 'deadline', 'maxObservationAge']:
 require(integer(payload[name], 0, 1000000000), 'Invalid global time')
 policy = payload['policy']
 keys(policy, 'maxAttempts initialDelay maxDelay multiplier')
 require(integer(policy['maxAttempts'], 1, 10), 'Invalid maxAttempts')
 require(integer(policy['initialDelay'], 1, 1000000), 'Invalid initialDelay')
 require(integer(policy['maxDelay'], policy['initialDelay'], 1000000), 'Invalid maxDelay')
 require(integer(policy['multiplier'], 1, 5), 'Invalid multiplier')
 tasks = payload['tasks']
 require(type(tasks) is list and 1 <= len(tasks) <= 20, 'Need 1..20 tasks')
 ids, results = set, []
 for task in tasks:
 keys(task, 'id state safeRepeat attempts observedAt lastFinishedAt duration validation')
 identifier = task['id']
 require(type(identifier) is str and re.fullmatch(r'[A-Za-z0-9_-]{1,48}', identifier) is not None, 'Invalid task id')
 require(identifier not in ids, 'Duplicate task id')
 ids.add(identifier)
 require(task['state'] in ['running', 'unknown', 'succeeded', 'retryable', 'terminal'], 'Invalid state')
 require(type(task['safeRepeat']) is bool, 'safeRepeat must be boolean')
 require(integer(task['attempts'], 1, 10), 'Invalid attempts')
 require(integer(task['observedAt'], 0, payload['asOf']), 'Invalid observation time')
 require(integer(task['duration'], 1, 1000000), 'Invalid duration')
 require(integer(task['validation'], 0, 1000000), 'Invalid validation duration')
 ended = task['lastFinishedAt']
 if task['state'] in ['running', 'unknown']:
 require(ended is None, 'Unconfirmed completion must use null')
 else:
 require(integer(ended, 0, task['observedAt']), 'Invalid completion time')
 result = {'id': identifier, 'decision': None, 'backoff': None,
 'earliestStart': None, 'estimatedFinish': None}
 if payload['asOf'] - task['observedAt'] > payload['maxObservationAge']:
 result['decision'] = 'refresh-observation'
 elif task['state'] == 'running':
 result['decision'] = 'observe-running'
 elif task['state'] == 'unknown':
 result['decision'] = 'reconcile-remote'
 elif task['state'] == 'succeeded':
 result['decision'] = 'validate-output'
 elif task['state'] == 'terminal':
 result['decision'] = 'repair-cause'
 elif task['attempts'] >= policy['maxAttempts']:
 result['decision'] = 'attempts-exhausted'
 elif not task['safeRepeat']:
 result['decision'] = 'reconcile-effects'
 else:
 delay = min(policy['maxDelay'], policy['initialDelay'] * policy['multiplier'] ** (task['attempts'] - 1))
 start = max(payload['asOf'], ended + delay)
 finish = start + task['duration'] + task['validation']
 result.update(backoff=delay, earliestStart=start, estimatedFinish=finish)
 if finish > payload['deadline']:
 result['decision'] = 'deadline-infeasible'
 elif start > payload['asOf']:
 result['decision'] = 'wait-backoff'
 else:
 result['decision'] = 'retry-eligible'
 results.append(result)
 return {'tasks': sorted(results, key=lambda item: item['id']),
 'deadlinePassed': payload['asOf'] > payload['deadline'],
 'jobsExecuted': False, 'idempotencyProven': False,
 'completionGuaranteed': False, 'productionApproval': False}


def main:
 raw = sys.stdin.read(1000001)
 require(len(raw) <= 1000000, 'Input too large')
 print(json.dumps(evaluate(json.loads(raw)), sort_keys=True))


if __name__ == '__main__':
 try:
 main
 except (ValueError, TypeError, RecursionError):
 print('Invalid retry fixture', file=sys.stderr)
 sys.exit(2)
IN PRACTICE

On the teaching clock, report waits until 103 and is estimated to finish at 128; load reaches deadline 140 including validation. append requires remote reconciliation.

Common pitfalls

Treating timeout as remote failure, DONE as success, max_retries as total attempts, free workers as free pool slots, or a favorable estimate as production authorization.

Related topics: Operations and capacity · Ingestion and quarantine · Dataset recovery

Take this idea with you

Retrying requires sufficiently recent state, understood effects, capacity and time to validate delivery.

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.