AI 智能体一旦从“回答问题”走向退款、库存调整、人工审批和多智能体协作,就会遇到一个经典的分布式系统难题:业务状态保存在数据库里,执行指令却发往另一套消息系统。两个提交点无法天然保持一致,可能出现“状态已更新但任务没发出去”,也可能出现“任务已执行但状态最终回滚”。
现已正式可用的 Spanner queues 将事务消息直接嵌入 Spanner。创建消息因此不再是事务结束后的额外动作,而是一次普通的事务写入:智能体的状态变化与后续执行意图要么一起提交,要么一起失败。
为什么智能体工作流尤其需要事务消息
传统服务同样会遭遇数据库与消息队列的双写问题,但智能体工作流会进一步放大风险:
- LLM 推理和外部工具调用可能耗时较长,任务租约可能在执行期间过期;
- 多个智能体可能同时读取、推测或修改同一份状态;
- 自动重试可能重复触发退款、发货或库存扣减;
- 人工审批、延迟检查和 SLA 升级会产生持续数小时甚至数天的任务;
- 智能体还要保存阶段性记忆、交接上下文和执行轨迹。
过去,团队通常用 transactional outbox、幂等键、定时扫描器和对账任务拼出可靠性。Spanner queues 的关键变化,是将队列定义为 Spanner 中的一等关系结构,并让入队、业务更新和 ACK 都能参与读写事务。
Spanner 的严格可串行化和全局外部一致性在这里提供了明确边界:如果订单被标记为“退款已批准”,对应的退款任务也一定已经提交;如果任务没有提交,状态更新同样不会生效。
一次事务完成“决定并行动”
下面用退款流程说明这种模型。示例假设已经存在 Orders 父表;执行前需要按实际结构补齐 Status、UpdatedAt 等字段,并由应用通过 Spanner 读写事务传入参数。
-- 假设父表已经存在:
-- CREATE TABLE Orders (
-- OrderId STRING(64) NOT NULL,
-- Status STRING(32),
-- UpdatedAt TIMESTAMP,
-- RefundTransactionId STRING(128)
-- ) PRIMARY KEY (OrderId);
CREATE QUEUE OrderAgentTasks (
OrderId STRING(64) NOT NULL,
TaskId STRING(64) NOT NULL,
TaskType STRING(64) NOT NULL,
Payload JSON NOT NULL
) PRIMARY KEY (OrderId, TaskId, TaskType),
INTERLEAVE IN Orders;
智能体批准退款时,在同一个读写事务中更新订单并写入任务:
-- 在同一个 Spanner Read-Write Transaction 中执行
UPDATE Orders
SET Status = 'REFUND_APPROVED',
UpdatedAt = PENDING_COMMIT_TIMESTAMP()
WHERE OrderId = @orderId;
INSERT INTO OrderAgentTasks (
OrderId,
TaskId,
TaskType,
Payload
)
VALUES (
@orderId,
@taskId,
'EXECUTE_REFUND',
JSON '{"action":"execute_refund","amount":49.99}'
);
这个事务建立了一个很有价值的不变量:EXECUTE_REFUND 任务存在,当且仅当订单成功进入 REFUND_APPROVED 状态。应用不再需要在提交数据库后调用另一个消息服务,也不必处理这个调用失败后留下的半完成状态。
同一种方式也适合多智能体交接。主智能体可以同时写入交接状态、上下文摘要和面向专业智能体的任务,使“记住了什么”和“下一步由谁执行”保持同步。
延迟投递让超时与人工审批进入事务
队列消息可以立即交付,也可以通过系统 DeliverTime 列安排未来投递。例如,在提交待审批状态时,同时安排 72 小时后的升级任务:
INSERT INTO OrderAgentTasks (
OrderId,
TaskId,
TaskType,
Payload,
DeliverTime
)
VALUES (
@orderId,
@escalationTaskId,
'ESCALATE_UNAPPROVED_ORDER',
JSON '{"action":"escalate_to_supervisor"}',
TIMESTAMP_ADD(CURRENT_TIMESTAMP(), INTERVAL 72 HOUR)
);
如果经理四小时后批准,可以在更新订单时原子删除尚未触发的升级任务:
-- 在同一个读写事务中执行
UPDATE Orders
SET Status = 'MANAGER_APPROVED',
ApprovedBy = @managerId
WHERE OrderId = @orderId;
DELETE FROM OrderAgentTasks
WHERE OrderId = @orderId
AND TaskId = @escalationTaskId
AND TaskType = 'ESCALATE_UNAPPROVED_ORDER',
ASSERT_ROWS_MODIFIED 1;
ASSERT_ROWS_MODIFIED 1 是这里的重要防线。如果待取消的任务已经不存在,语句会报错;应用应捕获该错误并中止整个事务,而不是继续提交一个可能与真实任务状态冲突的审批结果。
这类延迟消息可以覆盖审批超时、失败后的退避重试、定期回访和 SLA 升级,不再必须额外维护 cron 服务或轮询表。
Worker:流式拉取、续租与原子 ACK
Worker 通过 RECEIVE_<QueueName> 表值函数和 ExecuteStreamingSql 长连接消费任务。以下查询会持续接收可投递消息,并返回租约令牌与过期时间:
SELECT
OrderId,
TaskId,
TaskType,
Payload,
SpannerLeaseToken,
SpannerLeaseExpirationTimestamp
FROM RECEIVE_OrderAgentTasks(max_duration => '20m');
LLM 多轮推理或外部 API 调用可能超过默认租约窗口。Worker 应在租约到期前主动续租:
SELECT *
FROM RENEWLEASE_OrderAgentTasks(
lease_tokens => [@leaseToken]
);
任务完成后,将最终业务状态与删除队列消息放进同一个事务:
-- 在同一个 Spanner Read-Write Transaction 中执行
UPDATE Orders
SET RefundTransactionId = @externalRefundId,
Status = 'REFUND_COMPLETED'
WHERE OrderId = @orderId;
DELETE FROM OrderAgentTasks
WHERE OrderId = @orderId
AND TaskId = @taskId
AND TaskType = @taskType
ASSERT_ROWS_MODIFIED 1;
如果某个 Worker 因网络暂停而失去租约,另一个 Worker 可能已经完成并删除任务。此时旧 Worker 的删除操作无法影响一行,ASSERT_ROWS_MODIFIED 1 会触发错误。应用中止事务后,旧 Worker 就不能用过期结果覆盖新状态。
不过,“事务性 ACK”不等于任意外部副作用天然只执行一次。Spanner queues 提供至少一次投递与至多一次 ACK,可据此实现 exactly-once processing;如果任务调用支付、邮件或第三方物流 API,仍应把 TaskId 作为外部接口的幂等键。否则 Worker 在外部调用成功、但 ACK 前崩溃时,重试仍可能重复产生副作用。
一个稳妥的处理顺序是:
- 接收任务并记录
TaskId; - 调用支持幂等键的外部 API;
- 必要时续租,避免长任务失去所有权;
- 在读写事务中保存外部结果并删除消息;
- 若
ASSERT_ROWS_MODIFIED失败,中止事务并丢弃当前 Worker 的陈旧结果。
Queue 与 Change Stream 不应混用
Spanner change streams 和 Spanner queues 都能让下游响应数据变化,但它们解决的问题不同。
| 能力 | Spanner queues | Spanner change streams |
|---|---|---|
| 主要用途 | 事务任务编排 | CDC、审计和下游数据同步 |
| 消息租约 | 原生支持 | 不是其核心模型 |
| 延迟投递 | 支持 | 不面向任务调度 |
| 事务内 ACK | 支持 | 不以任务确认作为目标 |
| 消费方式 | SQL 拉取与流式接收 | 持续读取数据变更 |
需要“谁来执行、何时执行、执行后如何确认”时,应考虑 queue;需要把插入、更新和删除持续同步到分析、存储或其他系统时,change stream 更合适。不要仅因为两者都能产生事件,就把 CDC 当成任务队列使用。
上线前的工程检查表
Spanner queues 降低了数据库与消息基础设施分离带来的可靠性成本,但仍需要清晰的应用协议:
- 为每个任务生成稳定且全局唯一的
TaskId; - 外部工具调用必须尽量支持幂等键;
- 根据 LLM 和第三方 API 的最长耗时设计续租策略;
- 对
ASSERT_ROWS_MODIFIED失败采用“中止事务”,不要盲目忽略; - 用 SQL 监控积压、延迟任务和长期未完成任务;
- 为不可恢复错误设计终止状态、告警或人工处理流程;
- 明确队列数据的保留、审计与敏感信息访问策略;
- 区分任务编排与 CDC,避免让一种机制承担不适合的职责。
这项能力不只适合 AI 智能体。订单处理、库存工作流、金融通知、实时动态和高吞吐异步任务,同样能从“业务状态与执行意图共同提交”中获益。它真正减少的不是一套队列 API,而是双写、补偿和对账所造成的长期复杂度。