AI agents often need to coordinate multiple steps before completing a task. An agent might create a support ticket, update a customer record, request approval, and then trigger a follow-up action. If one step fails after another has already succeeded, the workflow can leave behind inconsistent state or duplicate operations.
Traditional queue-based architectures help decouple producers and consumers, but they introduce an additional consistency problem when application state and queued messages are stored in separate systems.
Google Cloud Spanner offers a different approach: use Spanner transactions to update application data and record messages that need to be processed. A worker can then consume those messages and execute the next step of an agent workflow.
This pattern is especially useful when an application needs transactional consistency between business state and work scheduling. It does not make external operations automatically transactional, but it can prevent a common failure in which an application commits a database change and loses the corresponding message.
For agent platforms, the key is to treat workflow state, queued work, retries, and external side effects as explicit parts of the system design.
Why Agent Workflows Need Transactional Messaging
Consider an AI agent that investigates an operational issue and decides to create a Jira ticket. The application first saves the investigation result in Spanner and then sends a message to a worker that calls the Jira API.
A naive implementation might look like this:
Save investigation result
|
v
Commit database transaction
|
v
Publish queue message
|
v
Create Jira ticket
The problem occurs when the database transaction succeeds but publishing the message fails. The investigation is recorded, yet no worker receives the instruction to create the ticket.
Reversing the order does not eliminate the problem. If the message is published first and the database transaction subsequently fails, the worker may process a message for a state change that never committed.
This is the dual-write problem: the application must update two independent systems, and there is no atomic transaction covering both operations.
A transactional messaging pattern addresses the database-to-message boundary by recording the business change and the message in the same Spanner transaction.
The worker processes the recorded message after the transaction commits. This ensures that the message cannot become visible as committed work without the corresponding database changes also committing.
Understand the Spanner Transactional Queue Pattern
A simple implementation uses two logical tables:
The application writes to both tables inside a single Spanner read-write transaction.
<box border radius="lg" padding={3} gap={2} align="center">
<box background="surface-secondary" radius="md" padding={3} width="100%" align="center" gap={1}>
**Agent application**
<text color="secondary" size="sm">Produces a workflow action</text>
</box>
<icon name="arrow-down" size="lg" color="secondary" />
<box border radius="md" padding={3} width="100%" align="center" gap={1}>
**Spanner read-write transaction**
<text color="secondary" size="sm">Update workflow state and insert the queue message atomically</text>
</box>
<icon name="arrow-down" size="lg" color="secondary" />
<box background="surface-secondary" radius="md" padding={3} width="100%" align="center" gap={1}>
**Queue worker**
<text color="secondary" size="sm">Find committed messages and claim work</text>
</box>
<icon name="arrow-down" size="lg" color="secondary" />
<box border radius="md" padding={3} width="100%" align="center" gap={1}>
**External action**
<text color="secondary" size="sm">Call an API, dispatch a task, or advance the workflow</text>
</box>
<text color="secondary" size="sm">Record the processing outcome and retry safely when necessary.</text>
</box>
The first transaction is the critical consistency boundary. It makes the application state and the queued instruction durable together.
The external action is a separate operation. A call to Jira, an AI model, or another API cannot normally participate in the same Spanner transaction. Consequently, the worker still needs idempotency, retry handling, and recovery logic.
The architecture provides transactional message creation, not automatic exactly-once execution across external systems.
Design the Queue Schema
The queue schema should represent the lifecycle of a message explicitly. At a minimum, a message needs a unique identifier, its type, its payload, its creation time, and a processing state.
The following Spanner SQL illustrates a queue table:
CREATE TABLE AgentQueue (
QueueId STRING(36) NOT NULL,
TaskId STRING(128) NOT NULL,
MessageType STRING(128) NOT NULL,
Payload JSON NOT NULL,
Status STRING(16) NOT NULL,
CreatedAt TIMESTAMP NOT NULL
OPTIONS (allow_commit_timestamp = true),
AvailableAt TIMESTAMP NOT NULL,
AttemptCount INT64 NOT NULL,
LeaseOwner STRING(128),
LeaseExpiresAt TIMESTAMP,
LastError STRING(MAX)
) PRIMARY KEY (QueueId);
This is an illustrative schema, not a complete production queue implementation. The exact types and constraints should reflect the application and the Spanner configuration.
The fields serve different purposes:
QueueId uniquely identifies a message.
TaskId links the message to the agent workflow.
MessageType identifies the operation to execute.
Payload contains the required input.
Status tracks whether the message is pending, processing, completed, or failed.
AvailableAt determines when the message can be processed.
AttemptCount supports retry limits.
LeaseOwner and LeaseExpiresAt support temporary worker ownership.
LastError records a bounded description of the latest failure.
A production design should also define how completed messages are retained or deleted, how failed messages are inspected, and how message ordering is handled.
Avoid storing unnecessary sensitive information in the payload. Where possible, store identifiers that allow the worker to retrieve the required data using its own authorized access.
Create the Business Change and Queue Message Atomically
Suppose an agent has completed an investigation and the application needs to schedule a follow-up action.
The application should update the investigation record and insert the queue message in the same Spanner read-write transaction.
The following Python example uses the Google Cloud Spanner client library. It illustrates the transactional boundary rather than a complete queue service.
import json
import uuid
from google.cloud import spanner
def record_investigation_and_enqueue(
database,
task_id,
investigation_result,
):
queue_id = str(uuid.uuid4())
payload = {
"taskId": task_id,
"action": "create_follow_up",
"result": investigation_result,
}
def transaction_body(transaction):
transaction.update(
table="Investigations",
columns=["TaskId", "Status", "Result"],
values=[[
task_id,
"COMPLETED",
json.dumps(investigation_result),
]],
)
transaction.insert(
table="AgentQueue",
columns=[
"QueueId",
"TaskId",
"MessageType",
"Payload",
"Status",
"CreatedAt",
"AvailableAt",
"AttemptCount",
],
values=[[
queue_id,
task_id,
"FOLLOW_UP",
json.dumps(payload),
"PENDING",
spanner.COMMIT_TIMESTAMP,
spanner.COMMIT_TIMESTAMP,
0,
]],
)
database.run_in_transaction(transaction_body)
return queue_id
The example assumes that an Investigations table exists with compatible columns and that the Spanner client library supports the shown transaction and commit-timestamp usage.
The Payload column is declared as JSON in the sample schema. Adapt the values to the precise serialization and type requirements of the library version used by your application.
If the transaction aborts, the business update and queue insertion do not commit independently. If the transaction commits, both become durable.
Spanner can retry a transaction callback when it encounters retryable transaction conflicts. For that reason, transaction callbacks must not perform external side effects such as sending an HTTP request or invoking an AI model. Those operations belong in the worker after the queue message has committed.
Claim Messages Safely With Multiple Workers
A queue becomes more difficult when multiple workers run concurrently. Without a claiming mechanism, two workers might read the same pending message and both execute the external action.
A worker should claim a message using a transactional state change or another concurrency-control mechanism supported by the queue design.
One practical approach is a lease. A worker claims a pending message by recording its identity and an expiration time. Other workers should not process the message while the lease is valid.
The state transition might look like this:
PENDING
|
v
PROCESSING
|
+------> COMPLETED
|
+------> PENDING (retry scheduled)
|
+------> FAILED (retry limit reached)
The claim operation must be concurrency-safe. A plain read followed by an update in separate transactions is not sufficient because two workers could observe the same pending state.
Use a Spanner read-write transaction that verifies the current state and changes it only when the message is eligible for claiming. If another worker changes the row concurrently, transaction conflict handling or a conditional state check must prevent both workers from successfully claiming the same message.
Leases also need expiration. A worker might crash after claiming a message but before completing it. Without lease expiration or a recovery process, the message could remain stuck in the processing state indefinitely.
When a lease expires, a recovery worker can make the message eligible for another attempt, subject to the application's retry policy.
Build a Reliable Worker
The worker reads eligible messages, claims them, executes the corresponding action, and records the result.
A simplified processing loop is:
Find messages that are pending and available.
Claim one message using a concurrency-safe transaction.
Load any additional data required by the action.
Execute the action outside the Spanner transaction.
Record the result in a new transaction.
Retry transient failures or move permanent failures to a terminal state.
The external action should not run inside the transaction used to claim the message. Keeping network calls outside database transactions avoids holding transaction resources while waiting for an external service.
For an agent workflow, the action might create a Jira issue, send a notification, invoke a model, or start another task. Each action type should have an explicit handler and validation rules.
A worker should also check whether the workflow is still in a state that permits the action. For example, a queued notification may no longer be appropriate if the task was canceled after the message was created.
This check does not replace idempotency, but it helps prevent stale work from producing unwanted effects.
Handle Retries and Duplicate Processing
A worker can fail after an external operation succeeds but before the queue message is marked completed.
For example:
The worker sends a request to create a Jira issue.
Jira creates the issue successfully.
The worker crashes before recording the Jira issue key in Spanner.
The message becomes eligible for retry.
The next worker sends the creation request again.
The result could be two Jira issues for the same task.
This is why transactional message creation alone cannot guarantee exactly-once side effects.
Use an idempotency strategy that matches the external service. Where supported, pass a stable idempotency key derived from the queue message or business operation. If the destination does not support idempotency keys, use a durable operation identifier and a reconciliation process to detect an earlier successful action before retrying.
The worker should classify errors rather than retry every failure indefinitely.
Transient failures: Network interruptions, throttling, and selected server errors may be retried with exponential backoff and jitter.
Permanent failures: Invalid payloads, unsupported actions, or authorization errors usually require intervention or configuration changes.
Ambiguous outcomes: A timeout after sending a request may mean the external operation succeeded even though the worker did not receive the response.
For ambiguous outcomes, reconcile the destination state before repeating an operation whenever possible.
Set a maximum attempt count, define a dead-letter or terminal-failure state, and provide a way for operators to inspect and replay failed work safely.
Choose a Queue Consumption Strategy
Spanner provides the transactional storage primitives, but the application still needs a way to discover and schedule pending messages.
There are two broad approaches.
Poll the queue table
Workers periodically query for eligible messages and claim them.
Polling is simple to understand and can be sufficient for moderate workloads. However, the polling interval introduces a trade-off between processing latency and database load.
A very short interval can generate unnecessary reads when the queue is mostly empty. A long interval can delay processing even when messages are available.
The query should be designed around the eligibility conditions and the table's primary-key structure. At scale, avoid repeatedly scanning the entire queue table. Use a schema and indexing strategy that supports the worker's access pattern, and validate it with realistic data volumes.
Use change-driven scheduling
A change-driven design can reduce unnecessary polling by reacting to changes that indicate new work is available.
This requires an explicit event or notification mechanism and a recovery strategy for missed notifications. The durable queue remains the source of truth; a notification should tell workers to look for work rather than replace the message record itself.
The appropriate approach depends on throughput, latency requirements, operational complexity, and the mechanisms supported by the surrounding infrastructure.
For either strategy, workers must be able to recover after downtime without losing committed messages.
Control Ordering and Workflow Dependencies
Not every agent workflow can process messages in arbitrary order.
Suppose an agent emits two actions: update a task's state and then notify an external system. If the notification executes before the state update is visible to the external integration, the resulting workflow may be inconsistent.
The queue design should make ordering requirements explicit.
For independent tasks, workers can process messages concurrently. For actions that must be sequential, use a stable partition or workflow identifier and ensure that only the next eligible action is claimed for that workflow.
A timestamp alone is not always enough to guarantee strict ordering. Concurrent transactions can commit in an order different from the order in which application requests were initiated.
If the workflow needs strict sequence numbers, store them with the task and enforce the transition rules transactionally.
Also consider cancellation and supersession. An agent might produce a queued action and then determine that the action is no longer required. The worker should check the current workflow state before executing irreversible operations.
Monitor Queue Health and Recovery
A queue that stores every message reliably can still fail operationally if processing falls behind or messages become stuck.
Track both message state and processing performance.
Metric | Why it matters |
|---|
Pending message count | Indicates the amount of unprocessed work |
Oldest pending message age | Reveals delays in processing |
Processing duration | Helps identify slow handlers |
Retry count | Highlights unstable dependencies |
Expired lease count | Reveals worker crashes or stuck processing |
Terminal failure count | Shows work requiring operator intervention |
Duplicate-action rate | Helps detect idempotency problems |
Also record the queue ID, task ID, handler name, attempt number, and external operation identifier in structured logs.
Avoid including sensitive payloads or credentials in those logs.
Define operational thresholds based on the workflow's requirements. A notification queue may tolerate a short delay, while an agent responsible for time-sensitive remediation may require a much tighter processing objective.
The queue should also have an explicit recovery process. Operators need to inspect failed messages, understand the reason for failure, correct the problem, and retry work without creating duplicate side effects.
Common Implementation Mistakes
Publishing messages after committing the business transaction: This creates a failure window in which the business state is saved but the message is lost. Insert the queue message in the same transaction.
Calling external services inside the transaction callback: Transaction retries can repeat the callback. Keep external side effects outside the transaction.
Assuming a transaction provides exactly-once execution: It provides atomicity for Spanner writes, not for external APIs. Use idempotency and reconciliation.
Claiming work without concurrency control: Multiple workers can execute the same action. Use a transactionally safe claim mechanism.
Ignoring expired leases: A crashed worker can leave messages stuck in processing. Implement lease expiration and recovery.
Polling the entire queue table: Unbounded scans become increasingly expensive as the queue grows. Design queries around eligibility and access patterns.
Retrying permanent errors indefinitely: Invalid requests and authorization failures need correction, not endless retries. Define retry limits and terminal failure handling.
Summary
Cloud Spanner can support transactional messaging for AI agent workflows by storing business-state changes and queue messages in the same read-write transaction.
This pattern prevents the database-to-queue dual-write failure that occurs when an application commits a state change but cannot publish the corresponding message. Workers can then process committed messages using a concurrency-safe claim mechanism, leases, retry policies, and durable processing states.
The external action remains a separate consistency boundary. Integrations with Jira, model APIs, and other services still require idempotency and reconciliation because an external operation may succeed even when the worker fails before recording its result.
For production systems, focus on transaction boundaries, safe message claiming, ordering requirements, bounded retries, queue observability, and recovery procedures. These controls make the workflow reliable without pretending that a database transaction can atomically include every external side effect.
Join the conversation! Your thoughts help the community grow.