← Professional Cloud Architect: architecture and operations
09 / 14 · 120 MIN

Data, consistency and event publication

Design transaction boundaries, retries and data publication without confusing local success, message delivery and business effect.

Define the rule and transaction boundary

Start with the business rule that must not be violated. A reservation needs available units and must update both its record and availability together. Within one Spanner database, a serializable read-write transaction can include the critical read and both writes. One separate strong read does not turn later requests into a single atomic decision. For read-only reconciliation spanning several queries, use a shared snapshot through a read-only transaction. Independent reads may observe different instants even when each is current. Choosing stale reads accepts a past instant: that may suit a delay-tolerant report but not a screen required to show a just-confirmed commit. Record these choices in the design, linking each requirement to a concrete operation and isolation mode rather than assigning every guarantee to every product configuration.

Distinguish order, retry and uncertain outcomes

For serializable transactions, external consistency preserves the precedence of nonoverlapping commits: if a policy revision committed before its dependent approval began, a snapshot containing the approval cannot omit that revision. This does not create a transaction with an external processor. A read-write callback can execute several times, and an external call inside it can produce repeated effects. Likewise, DEADLINE_EXCEEDED on a write proves neither rollback nor success. Before resending, locate state associated with the original identity and apply the idempotent recovery procedure. In the settlement exercise, two external receipts need reconciliation even if the database contains one row. Do not delete evidence or invent financial compensation; identify owners, confirm each outcome and retain the relationship between intent, attempt and observed effect.

Distribute load without losing identity

A stable business key recognizes the same intent; a physical key also influences where writes concentrate. Do not confuse these roles. In Spanner, a monotonic timestamp as the first component can direct new inserts into the same range. Reversing sort order changes the end, not the pattern. A distributed prefix or sharding design should be evaluated alongside time queries, indexes and operations. A global index starting with the same timestamp can reproduce concentration. Before choosing, collect expected rate, actual key distribution, critical queries and the cost of searching several ranges. In a migration project, this analysis belongs in load rehearsal and acceptance criteria. Extra capacity does not establish a well-distributed access model; retain measurements showing how the proposed key serves both writes and required reads.

Choose the smallest sequence the business requires

Pub/Sub ordering is per key, with subscription ordering enabled and the same key published within one region. It is not global ordering across different accounts. If the business only needs per-account sequencing, one key for the whole bank creates unnecessary serialization. Choose granularity matching the rule and measure backlog. Ordering also does not eliminate redelivery: in an at-least-once sequence without a dead-letter topic, a redelivered message can bring back later messages for the key, even acknowledged ones. Consumers need to recognize effects already applied. If a callback schedules asynchronous work and returns, receipt order does not guarantee that work’s completion order. Explicitly design effect sequencing and acknowledgment timing so a concurrency optimization does not change the business result or hide unfinished processing.

Prepare failure handling and replay as operational capabilities

A dead-letter topic is more than a configured name. The Pub/Sub service agent needs to publish at the destination and acknowledge on the source subscription. A human operator’s permissions do not replace that identity’s permissions. Maximum delivery attempts is approximate; a requirement for exactly five business effects needs a different control. To replay acknowledged messages, plan topic retention or retention of those messages on the subscription before they are needed. Seek changes acknowledgment state but does not undo debits, emails or other completed actions. Delivery also does not converge instantly after seek. Prepare a procedure identifying the interval, consumer version, completion criteria and protection against repeated effects. Rehearsal should include returning messages whose work partly exists, so the team can distinguish recovery from duplication.

Publish file sets by version identity

Cloud Storage provides strong read-after-write and listing consistency, but public caches may still serve an old copy. Distinguish the origin from the layer answering the request. For concurrency, use preconditions: ifGenerationMatch=0 prevents replacing a live version during creation; deletion should retain the approved generation to avoid targeting a later replacement. For ranged reads, pin the generation to avoid assembling parts from different versions. None of this makes a batch of operations across several objects atomic. An accounts-and-positions dataset needs run identity and verified generations. An application can publish an acceptance manifest after validating the set; that is a design decision, not an automatic bucket guarantee. Also confirm who can change the manifest and how consumers reject incomplete sets before they produce reports.

Separate append, finalization and analytical visibility

In BigQuery’s Storage Write API gRPC, the default stream has at-least-once semantics. An application-created stream can use offsets to recognize appends within that stream. Two different streams at offset zero do not represent one position, even if they contain the same orderId. Business identity still needs handling. For a batch, pending streams defer visibility: workers write and finalize; the coordinator commits the group when everything is ready. Finalize stops further appends but does not publish data. BatchCommitWriteStreams covers finalized streams of one parent table. If stream_errors is nonempty, no stream in that group committed through the call. Retain the response and commit_time when present rather than treating receipt of any response as success of the whole batch or assuming partial success without evidence.

Lab: compare local state with the external effect

Run the code below with Python 3.12 or later and the sqlite3 module available. Before running it, predict three outcomes: what remains after unsafe_attempt aborts, what happens when relay loses the response and what changes when S2 is created instead of retrying S1. The lab uses an in-memory SQLite database and a dictionary as a fictional receiver. It demonstrates local atomicity and sequential identity recognition without running Spanner, Pub/Sub or payments. The dictionary is neither durable storage nor tested protection against concurrent consumers. In a real system, identity checking and effect application need suitable receiver guarantees. Connect the results to both final cases and write acceptance criteria: recoverable intent, stable references, inputs from the same run and explicit confirmation of complete publication.

"""Original in-memory SQLite teaching lab, not a Spanner or Pub/Sub emulator.

All external calls below are Python dictionary operations. No network, credentials,
files, concurrency, vendor services or actual payments are used.
"""
import json
import platform
import sqlite3

checks = []


def check(name, condition):
 if not condition:
 raise AssertionError(name)
 checks.append(name)


def fails(call, exception):
 try:
 call
 except exception:
 return True
 return False


db = sqlite3.connect(":memory:", isolation_level=None,
 autocommit=sqlite3.LEGACY_TRANSACTION_CONTROL)
db.executescript("""
CREATE TABLE orders (id TEXT PRIMARY KEY, amount INTEGER NOT NULL);
CREATE TABLE outbox (id TEXT PRIMARY KEY, amount INTEGER NOT NULL, sent INTEGER NOT NULL);
""")
unsafe_receipts = []


def unsafe_attempt(order_id, amount, abort=False):
 db.execute("BEGIN")
 try:
 # This fictional external list is outside SQLite's transaction boundary.
 unsafe_receipts.append((order_id, amount))
 db.execute("INSERT INTO orders VALUES (?,?)", (order_id, amount))
 if abort:
 raise RuntimeError("injected failure before commit")
 db.execute("COMMIT")
 except Exception:
 db.execute("ROLLBACK")
 raise


def enqueue(order_id, amount, abort=False):
 db.execute("BEGIN")
 try:
 existing = db.execute("SELECT amount FROM orders WHERE id=?", (order_id,)).fetchone
 if existing:
 if existing[0]!= amount:
 raise ValueError("same identity with different intent")
 db.execute("COMMIT")
 return "existing"
 db.execute("INSERT INTO orders VALUES (?,?)", (order_id, amount))
 if abort:
 raise RuntimeError("injected failure before outbox write")
 db.execute("INSERT INTO outbox VALUES (?,?, 0)", (order_id, amount))
 db.execute("COMMIT")
 return "created"
 except Exception:
 db.execute("ROLLBACK")
 raise


external_effects = {}
delivery_calls = []


def fictional_receiver(order_id, amount):
 delivery_calls.append((order_id, amount))
 if order_id in external_effects:
 if external_effects[order_id]!= amount:
 raise ValueError("receiver identity collision")
 return "existing"
 external_effects[order_id] = amount
 return "applied"


def relay(order_id, lose_response=False):
 row = db.execute("SELECT amount, sent FROM outbox WHERE id=?", (order_id,)).fetchone
 if row is None:
 raise LookupError("no publication intent")
 if row[1]:
 return "already-sent"
 receipt = fictional_receiver(order_id, row[0])
 if lose_response:
 raise TimeoutError("injected response loss after receiver applied effect")
 db.execute("UPDATE outbox SET sent=1 WHERE id=?", (order_id,))
 return receipt


try:
 check("unsafe attempt aborts", fails(lambda: unsafe_attempt("U1", 80, True), RuntimeError))
 check("aborted database row absent", db.execute("SELECT * FROM orders WHERE id='U1'").fetchone is None)
 check("external effect survives rollback", unsafe_receipts == [("U1", 80)])
 unsafe_attempt("U1", 80)
 check("retry creates second external effect", len(unsafe_receipts) == 2)
 check("one local row does not prove one external effect", db.execute("SELECT count(*) FROM orders WHERE id='U1'").fetchone[0] == 1)
 check("safe enqueue abort injected", fails(lambda: enqueue("S1", 120, True), RuntimeError))
 check("order rolls back with failed intent", db.execute("SELECT * FROM orders WHERE id='S1'").fetchone is None)
 check("no partial outbox row", db.execute("SELECT * FROM outbox WHERE id='S1'").fetchone is None)
 check("successful atomic enqueue", enqueue("S1", 120) == "created")
 check("order and intent both present", db.execute("SELECT o.amount, x.amount FROM orders o JOIN outbox x USING(id) WHERE id='S1'").fetchone == (120, 120))
 check("repeat original identity returns existing", enqueue("S1", 120) == "existing")
 check("different amount under same identity rejected", fails(lambda: enqueue("S1", 121), ValueError))
 check("collision does not change intent", db.execute("SELECT amount FROM outbox WHERE id='S1'").fetchone[0] == 120)
 check("response loss follows external effect", fails(lambda: relay("S1", True), TimeoutError))
 check("effect exists despite timeout", external_effects == {"S1": 120})
 check("intent remains retryable", db.execute("SELECT sent FROM outbox WHERE id='S1'").fetchone[0] == 0)
 check("receiver recognizes retry", relay("S1") == "existing")
 check("two calls one effect", len(delivery_calls) == 2 and len(external_effects) == 1)
 check("successful response marks sent", db.execute("SELECT sent FROM outbox WHERE id='S1'").fetchone[0] == 1)
 check("sent intent is not sent again", relay("S1") == "already-sent" and len(delivery_calls) == 2)
 check("new ID creates distinct intent", enqueue("S2", 120) == "created")
 check("new ID is a new external effect", relay("S2") == "applied" and len(external_effects) == 2)
 check("receiver rejects mismatched retry payload", fails(lambda: fictional_receiver("S1", 999), ValueError))
 check("receiver preserved original effect", external_effects["S1"] == 120)
 check("relay rejects missing intent", fails(lambda: relay("missing"), LookupError))
 print(json.dumps({"passed": len(checks), "checks": checks,
 "python": platform.python_version, "sqlite": sqlite3.sqlite_version,
 "database": ":memory:", "network": False,
 "vendorExecution": False,
 "scope": "Sequential SQLite atomic writes and synthetic receiver identity checks; no concurrency, actual transport, Spanner isolation or durable external idempotency tested."}))
finally:
 db.close
IN PRACTICE

An order committed once in the database, but a retried callback produced two external processor receipts. Reconciliation needs both systems.

Common pitfalls

Equating timeout with rollback, acknowledgment with one effect, object name with version and FinalizeWriteStream with batch publication.

Related topics: Data migration and reconciliation · Idempotency and incident recovery · Acceptance criteria and RUN handover

Take this idea with you

Identify each guarantee’s boundary, retain stable identity and require completion evidence across every relevant system.

Create account

Reference: Spanner transactions overview · 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.