1. Know what arrived and what remains unresolved
In a fictional APS example, the close pipeline succeeds but some movements were discarded during parsing. A green dashboard describes the process, not delivery completeness. Define three populations: received inputs, accepted logical events and rejected inputs. An event can arrive several times, so comparing only final counts with transport counts leads to incorrect conclusions. These cases do not represent BNP Paribas procedures. Before selecting services, fix the source and destination contract: input identity, event identity, field representation, permitted transformations and acceptance criteria. Silently changing a decimal comma to a point may seem convenient, but it requires an agreed rule. For every rejection, retain a reference that locates evidence, an understandable reason and a treatment owner. The production team should answer concrete questions: how many inputs arrived, how many contributed to accepted results and how many remain unresolved? Within what period can they be recovered? Who accepts partial delivery? Technical success is only part of this answer. Record observation limits and avoid presenting an unknown quantity as zero.
2. Separate filtering, delivery and validation
A filtered Pub/Sub subscription delivers only messages matching the selected attributes. The service automatically acknowledges the remainder; the filter does not inspect JSON inside data. A consumer that validated 700 of 1000 published messages therefore cannot present every message as validated. Define expected scope before comparing numbers. A syntactically correct filter can exclude a producer that stopped supplying an attribute. The dead-letter policy belongs to the subscription. The service can forward messages the subscriber cannot acknowledge, but the maximum attempt count is approximate. Counting depends on correct configuration and permissions and can reset under some conditions. Do not use it as a hard budget for business operations. A forwarded message has a new envelope and information identifying the source subscription; the quarantine reader should retain that correlation. Confirm the Pub/Sub managed identity's permission to publish to the destination and acknowledge on the subscription. An operator's successful manual publication does not prove that access. The destination also needs a consumption, observation and treatment path. Creating a topic does not automatically create a complete support operation.
3. Recover errors without inventing data
Classify the failure before choosing retries. A temporary service outage may clear; an impossible date remains impossible when the same parsing is repeated. Do not replace a rejected value with zero or the current date merely to obtain a green execution. A record error can use a separate output; a systemic authentication or storage failure needs operational treatment instead of being hidden as thousands of invalid records. An application quarantine pattern retains unprocessed elements for investigation and correction. A transport dead-letter policy does not automatically replace that output, especially after the application has acknowledged delivery. Confirm that the rejection destination actually receives records and that its own failure remains visible. Account separately for classification and persistence failures. Backoff reduces pressure during certain transient failures but is not an exact scheduler. It also does not make an operation idempotent. A request whose response was lost may have produced its effect. Look up state by reference where possible and define a budget, concurrency limit and escalation condition. Repeating retries across several layers can multiply calls and cost.
4. Run the local quarantine gate
The Python exercise below classifies a fictional batch in memory. It calls no APIs, stores no files and acknowledges no messages. The envelope contains records, with one to 100 inputs. Each input has a unique recordId and a payload. An ambiguous envelope stops the program; data problems go to quarantine. This distinction avoids emitting an apparently valid plan when inputs cannot even be referenced uniquely. The payload requires exactly eventId, amount and currency. Identifiers use ASCII alphanumeric characters, hyphens or underscores, up to 64 characters. amount is text with two decimal places, without spaces, a plus sign or redundant leading zeroes; a minus sign and up to nine digits before the point are allowed. currency accepts only EUR or USD in this teaching contract. The program converts to integer cents without rounding. These are local bounds, not Google Cloud limits. Run python3 content/labs/pde-quarantine-gate/run.py < content/labs/pde-quarantine-gate/case.json from the project root. Predict the result first. The fixture has two equivalent A copies, two conflicting B inputs, one valid C with an invalid peer, one invalid identifier and one valid D. Explain each classification before consulting the output.
5. Interpret identities and counts
recordId identifies a transport input; eventId identifies an immutable logical event in this contract. Two inputs with the same event and equal normalized values produce one accepted event while retaining both references. Normalization treats 0.00 and -0.00 as equivalent. It does not normalize commas, scientific notation or lowercase currencies. If valid values differ, the entire group conflicts without selecting the first or largest. There is another deliberately conservative rule: when an input has a valid eventId but another invalid field, its valid peers receive peer-invalid. The model does not choose an authoritative version. In the fixture, all of C remains unresolved even though one input has a valid amount. This policy should be discussed with the contract owner; it is not a universal rule for every pipeline. The result is two accepted events representing three inputs, with five inputs quarantined. Three plus five accounts for all eight inputs, while two counts logical outcomes. The duplicates list explains the difference without classifying an equivalent copy as invalid. Change batch order and confirm the result stays the same. Then change an amount and identify which counts and reasons change.
6. Understand the exercise limits
Output contains references and reasons without copying rejected payloads. This reduces information exposed by the report but does not prove anonymization: identifiers themselves may require protection in a real system. The model implements no retention, access control or recovery. rawPayloadIncluded, durablyStored, messagesAcknowledged and productionApproval remain false. Do not confuse the word quarantined with confirmation of a durable destination write. Comparison only knows the supplied batch. If A:1.00 appears in one process and A:2.00 in another, each execution may accept what it sees. Detecting the global conflict requires shared state, a time boundary, concurrency handling, consistent updates and late-data treatment. Raising the input limit or sorting each batch does not provide that missing context. Build three counterexamples: an unauthorized extra field, a repeated recordId and an event with two currencies. Explain why the first can be a data rejection, the second stops envelope processing and the third is a conflict. Then describe a real architecture that persists outcomes without promising atomicity across services that do not provide it. Identify where a failure would require reconciliation.
7. Diagnose startup and external pressure
Locate the phase in which the pipeline fails. In a Flex Template, the pipeline construction program must exit to allow launch. A wait_until_finish in that program can block startup. Also inspect the image entrypoint and launcher network dependencies. Current documentation describes a logging image pulled from gcr.io even when the custom image resides in Artifact Registry; access to one image does not prove access to all of them. Compare a reference template using the same identity and network, but treat the result as diagnostic evidence, not universal proof. Preinstalling tested dependencies in an image can reduce startup installation and avoid unavailable public repositories. Image access, library compatibility and the configuration actually promoted still need validation. After startup, measure useful throughput and rejections. More workers can increase 429 responses from a rate-limited API. Control concurrency, batch size and retries according to destination capacity. For join enrichment, an oversized reference can make side inputs inefficient; evaluate keyed distribution and shuffle cost. Every adjustment needs comparable outcomes, not just CPU utilization figures.
8. Hand over a usable procedure to RUN
Handover should connect signals to decisions. Define who observes rejection growth, how they distinguish producer errors from destination outages and where rule versions are recorded. Prepare known fictional samples: valid, duplicate, conflicting and invalid. The support team should predict outcomes and repeat the exercise without relying on the pipeline author. For real recovery, fix the input set, parser version, destination and permitted external effects. An approved correction should not erase previous evidence. Confirm recoverable storage before relying on quarantine and reconcile outcomes after replay. If a creation request loses its response, look up existing state before creating another batch under a different reference. As this lesson's final exercise, prepare a decision note for the fixture: which events are accepted, which five inputs remain unresolved, what evidence is missing and who decides publication? Add a rejection-destination failure and explain how the decision changes. The summary is straightforward: account, classify, preserve evidence, correct with authorization and verify delivery. Connect these steps to contracts, replay and capacity from the preceding lessons.
"""Bounded fictional batch gate; no storage, cloud calls or acknowledgement."""
import json
import re
import sys
IDENTIFIER = re.compile(r'[A-Za-z0-9_-]{1,64}', re.ASCII)
AMOUNT = re.compile(r'-?(?:0|[1-9][0-9]{0,8})\.[0-9]{2}', re.ASCII)
def identifier(value):
return isinstance(value, str) and IDENTIFIER.fullmatch(value) is not None
def inspect(payload):
issues = []
if not isinstance(payload, dict):
return None, None, ['payload-shape']
event = payload.get('eventId')
if not identifier(event):
event = None
issues.append('event-id')
if set(payload)!= {'eventId', 'amount', 'currency'}:
issues.append('payload-fields')
amount = payload.get('amount')
cents = None
if not isinstance(amount, str) or AMOUNT.fullmatch(amount) is None:
issues.append('amount-format')
else:
negative = amount.startswith('-')
major, minor = amount.lstrip('-').split('.')
cents = (int(major) * 100 + int(minor)) * (-1 if negative else 1)
currency = payload.get('currency')
if not isinstance(currency, str) or currency not in ('EUR', 'USD'):
issues.append('currency')
return event, (cents, currency) if not issues else None, sorted(issues)
def evaluate(data):
if not isinstance(data, dict) or set(data)!= {'records'}:
raise ValueError('expected records envelope')
records = data['records']
if not isinstance(records, list) or not 1 <= len(records) <= 100:
raise ValueError('expected 1..100 records')
groups, rejected, ids = {}, [], set
for record in records:
if not isinstance(record, dict) or set(record)!= {'recordId', 'payload'}:
raise ValueError('invalid record envelope')
rid = record['recordId']
if not identifier(rid) or rid in ids:
raise ValueError('invalid or duplicate recordId')
ids.add(rid)
event, value, reasons = inspect(record['payload'])
row = {'recordId': rid, 'eventId': event, 'reasons': reasons}
if event is None:
rejected.append(row)
else:
groups.setdefault(event, []).append((row, value))
accepted, duplicates = [], []
for event, group in sorted(groups.items):
invalid_peer = any(row['reasons'] for row, _ in group)
values = {value for _, value in group if value is not None}
conflict = len(values) > 1
if invalid_peer or conflict:
for row, _ in group:
reasons = set(row['reasons'])
if invalid_peer and not reasons:
reasons.add('peer-invalid')
if conflict:
reasons.add('event-conflict')
rejected.append({**row, 'reasons': sorted(reasons)})
else:
cents, currency = next(iter(values))
source_ids = sorted(row['recordId'] for row, _ in group)
accepted.append({'eventId': event, 'cents': cents,
'currency': currency, 'recordIds': source_ids})
if len(source_ids) > 1:
duplicates.append({'eventId': event, 'recordIds': source_ids,
'collapsed': len(source_ids) - 1})
rejected.sort(key=lambda row: row['recordId'])
return {'accepted': accepted, 'quarantined': rejected,
'duplicates': duplicates, 'inputRecords': len(records),
'acceptedSourceRecords': sum(len(x['recordIds']) for x in accepted),
'quarantinedRecords': len(rejected),
'rawPayloadIncluded': False, 'durablyStored': False,
'messagesAcknowledged': False, 'productionApproval': False}
def unique_object(pairs):
result = {}
for key, value in pairs:
if key in result:
raise ValueError('duplicate JSON key')
result[key] = value
return result
if __name__ == '__main__':
try:
raw = sys.stdin.read(1000001)
if len(raw) > 1000000:
raise ValueError('input too large')
print(json.dumps(evaluate(json.loads(raw, object_pairs_hook=unique_object)), sort_keys=True))
except (ValueError, TypeError, RecursionError) as error:
print('invalid input: ' + str(error), file=sys.stderr)
sys.exit(2)
Eight inputs produce two accepted events from three inputs and five rejections; no durable write is performed.
Common pitfalls
Empty backlog as complete validation; retry as data correction; zero as substitution; classification as persistence; more workers as a response to 429.
Related topics: Contracts and consumers · Replay and external effects · Capacity and recovery
A green process does not prove complete delivery: retain counts, reasons and recoverable evidence.
Reference: Professional Data Engineer standard exam guide · Current linked standard guide (document title v4.2); edition date unconfirmed (2026-09-30 inspection)