Spanner queues 2026-10-02 GA:事务内发消息,RECEIVE_* TVF 收消息与租约 ack
Google Cloud 于 2026 年 10 月 2 日宣布 Spanner queues 正式可用(GA),由产品经理 Nitin Sagar 与工程经理 Matthew Mucklo 发布。核心机制是原生事务性消息(native transactional messaging)直接内嵌在 Spanner 内:创建一条消息就是事务里的又一次写入,业务状态变更和下游动作要么一起提交,要么一起回滚。官方主要面向 AI agent / 智能体负载:agent 需要在运营数据库里维护内部状态,同时通过独立的消息/事件队列派发异步动作(退款、库存、多步交接、编排子代理)。两个系统提交点不一致,就会破坏事务一致性。
它要解决的“双写问题”
传统做法:数据库状态更新成功,但动作派发失败——agent 决定要动手却没动;或者消息派发成功,但本地事务回滚——发了不该发的动作。开发团队被迫自建 outbox 模式(把要发的消息先写进同库的一张表,再由独立进程捞出来发送)、幂等层和补偿 worker。官方称这是 agentic 架构上一笔沉重的可靠性税。Spanner queues 让发消息变成同一事务中的一次写:要么都提交,要么都不存在。
队列怎么定义
队列是关系结构,用 GoogleSQL DDL 定义,形似建表。必须包含一个 Payload 列(PostgreSQL 方言里是 payload)和主键。除了显式建列,Spanner 会自动加一个系统列 DeliverTime(PG 是 deliver_time)。
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;
在同一个读写事务里更新业务表并发消息:
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}');
发消息也可以用 mutation API(send),send 可指定 deliverTime。
收消息、租约与确认
接收端用 ExecuteStreamingSQL 调用表值函数 RECEIVE_<队列名>(),这是一个长连接流式调用,每个 worker、每个队列都要循环执行一个这样的调用。必须用强读(strong read),Spanner 拒绝 stale read。
SELECT UserId, MessageId, Payload, DeliverTime,
SpannerLeaseExpirationTimestamp, SpannerLeaseToken
FROM RECEIVE_UserTasks(max_duration=>'20m');
可以用 max_batch_size 批量接收提高吞吐。收到消息即获得一个租约(lease),默认 10 秒,Spanner 返回唯一的 SpannerLeaseToken 与过期时间戳。处理时间可能超过租约时,用 RENEWLEASE_<队列名>() TVF 续租;官方建议在处理进行到约 7–8 秒时续租,给网络延迟留余量。
处理完用 DELETE DML 确认(ack),且必须与本次处理的其它写入在同一个事务里:
DELETE FROM UserTasks
WHERE UserId = @userId AND MessageId = @messageId
ASSERT_ROWS_MODIFIED 1;
如果租约已过期、别的进程已经处理过这条消息,ASSERT_ROWS_MODIFIED 1 会报错,从而阻止旧进程覆盖新状态。也可以走客户端库的 Ack mutation(ignoreNotFound 可把确认不存在的消息当作成功)。
投递语义
官方表述是:至少一次投递(at-least-once delivery)+ 至多一次确认(at-most-once ACK),二者合起来实现 exactly-once 处理——前提是你要承认同一条消息可能被收到两次。调用外部 API 时,官方示例把 TaskId 当接收端的幂等键(idempotency key)传下去。付款、发货这类重复执行有严重后果的流程,这个设计是必须的。
常用模式
- 事务提交后触发:在同一事务里发消息,提交后接收端 stream 到消息、执行、ack。
- 延迟投递:设置系统列 DeliverTime,消息在该时间前对接收端不可见。示例是把升级动作排到 72 小时后。
- 取消已排期消息:如果审批先到,在更新业务状态的同一事务里 DELETE 掉那条排期消息。
- 大 payload:用 out-of-band 存储,把大 body 放另一张表,队列消息里只放引用;Payload 可以很小。
- checkpoint:最佳做法是原子地 ack 当前消息 + 发一条排在未来投递的新消息(携带或指向更新后的状态),避免立刻重投到别的 worker。
- 时间对齐批处理:发送端把消息的 DeliverTime 对齐到未来某个离散时间窗(例如向上取整到最近的 10 秒边界),让 Spanner 把同 split 的消息聚成一批投递。
- dirty flag:事务修改某表时,同事务往队列写一条轻量“脏位”消息,后台 worker 异步做昂贵处理。
- 显式重试:Spanner 会自动对失败/未确认消息做指数退避重试。但如果失败原因已知且时长已知(如外部限流 HTTP 429 带 Retry-After),自动退避可能过早重试浪费 CPU——此时应捕获该瞬时失败、ack 当前消息,并在同一事务里发一条 DeliverTime 设为指定未来时间(示例是 5 分钟后)的替代消息。
监控与排障
相关指标:buffered_ready_messages(内存中就绪待投递的消息数)、message ack_count、oldest_unacked_message_age(最老未确认消息的年龄)、lease_expiration_count。当 lease_expiration_count 升高,说明处理时间超过了 10 秒租约时长、worker 还没来得及 ack——典型解法就是前面说的主动续租。
与 change streams 的区别:change streams 面向持续 CDC、把变更流给下游分析/存储;queues 面向事务性任务编排,有原生消息租约、延迟投递、SQL 拉取与读写事务内的原子确认。
版本与判断
按报道,该功能面向 Google Spanner Enterprise 与 Enterprise Plus 版本开放。需要注意:如果你的业务系统本来跑在 PostgreSQL 之类的单库上,同样可以用“同库放一张 outbox 表”来避免双写——队列本身不太可能成为迁到 Spanner 的理由。Spanner queues 真正有价值的场景,是你已经在用(或将要用)Spanner,想省掉外部队列和 outbox 中继进程。
参考: