使用 Spanner 队列

本文档介绍了如何使用 Spanner 队列。本文介绍了如何创建队列、发送和接收消息、延长消息租约以及确认消息。此外,还包括最佳实践、有关监控队列的信息和问题排查指南。

创建队列

如需创建队列,请使用 CREATE QUEUE 语句。

GoogleSQL

-- Example table for interleaving
CREATE TABLE Users (
  UserId   INT64 NOT NULL,
  UserName STRING(MAX)
) PRIMARY KEY (UserId);

-- Queue for processing user-related tasks
-- NOTE: Queues do not have to be interleaved, but for locality and
-- pre-warming purposes, it's recommended if you are inserting into a table
-- and a queue simultaneously.
CREATE QUEUE UserTasks (
  UserId     INT64 NOT NULL,
  MessageId  STRING(36) NOT NULL, -- UUID recommended
  Payload    BYTES(MAX) NOT NULL  -- Also: Proto, JSON, String are possible.
) PRIMARY KEY (UserId, MessageId),
INTERLEAVE IN PARENT Users ON DELETE CASCADE;

PostgreSQL

-- Example table for interleaving
CREATE TABLE users (
  userid   bigint NOT NULL,
  username varchar,
  PRIMARY KEY (userid)
);

-- Queue for processing user-related tasks
-- NOTE: Queues do not have to be interleaved, but for locality and
-- pre-warming purposes, it's recommended if you are inserting into a table
-- and a queue simultaneously.
CREATE QUEUE usertasks (
  userid     bigint NOT NULL,
  messageid  varchar(36) NOT NULL, -- UUID recommended
  payload    bytea NOT NULL, -- Also: text, varchar, jsonb are possible.
  PRIMARY KEY (userid, messageid)
) INTERLEAVE IN PARENT users ON DELETE CASCADE;

唯一不需要在 CREATE QUEUE 语句中明确创建的列在 GoogleSQL 中称为 DeliverTime,在 PostgreSQL 中称为 deliver_time。它们由 Spanner 自动创建。

队列支持存留时间 (TTL) 政策,这有助于管理旧的未确认消息的消息积压。

发送消息

如需向队列发送消息,请使用 INSERT DML 语句:

GoogleSQL

-- Send message immediately
INSERT INTO Users (UserId) VALUES (123);
INSERT INTO UserTasks (UserId, MessageId, Payload, DeliverTime)
VALUES (123, 'some-unique-id-1', b'Your task payload here', CURRENT_TIMESTAMP());

-- Schedule a message delivery
INSERT INTO UserTasks (UserId, MessageId, Payload, DeliverTime)
VALUES (123, 'some-unique-id-2', b'Scheduled task', TIMESTAMP_ADD(CURRENT_TIMESTAMP(), INTERVAL 1 HOUR));

PostgreSQL

-- Send message immediately
INSERT INTO users (userid) VALUES (123);
INSERT INTO usertasks (userid, messageid, payload, deliver_time)
VALUES (123, 'some-unique-id-1', CAST('Your task payload here' AS bytea), CURRENT_TIMESTAMP);

-- Schedule a message delivery
INSERT INTO usertasks (userid, messageid, payload, deliver_time)
VALUES (123, 'some-unique-id-2', CAST('Scheduled task' AS bytea), CURRENT_TIMESTAMP + INTERVAL '1 HOUR');

或者,也可以使用客户端库 Send 突变来插入消息:

Go

m := spanner.Send("UserTasks", spanner.Key{int64(123), "some-unique-id-1"}, []byte("Your task payload here"), spanner.WithDeliveryTime(futureTime))
_, err := client.Apply(ctx, []*spanner.Mutation{m})

Java

dbClient.write(
        Collections.singletonList(
            Mutation.newSendBuilder("UserTasks")
                .setKey(Key.of(123L, "some-unique-id-1"))
                .setPayload(Value.bytes(ByteArray.copyFrom("message3")))
                .setDeliveryTime(futureTime)
                .build()));

接收消息

RECEIVE_QUEUE_NAME() 表值函数 (TVF) 与 ExecuteStreamingSQL 搭配使用,以接收消息。这是一个长时间运行的调用。您必须以循环方式为每个 worker 和每个队列执行一次这些调用。使用强读取执行查询,因为 Spanner 会拒绝过时数据读取。

GoogleSQL

-- SQL query to stream messages
SELECT
  UserId,
  MessageId,
  Payload,
  DeliverTime,
  SpannerLeaseExpirationTimestamp,
  SpannerLeaseToken
FROM RECEIVE_UserTasks(max_duration=>'20m');

PostgreSQL

-- SQL query to stream messages
SELECT
  userid,
  messageid,
  payload,
  deliver_time,
  spanner_lease_expiration_timestamp,
  spanner_lease_token
FROM spanner.receive_usertasks(NULL, NULL, '20m');

您的客户端代码应使用流式查询来迭代处理结果。 返回的每一行都是一条消息。

批量接收消息

如需通过同时处理多条消息来提高吞吐量,您可以通过指定 max_batch_size 实参来批量接收消息:

GoogleSQL

-- SQL query to stream messages
SELECT
  UserId,
  MessageId,
  Payload,
  DeliverTime,
  SpannerLeaseExpirationTimestamp,
  SpannerLeaseToken,
  SpannerLastBatchMessage -- Special boolean column returns TRUE if the
                          -- last message is in a batch.
FROM RECEIVE_UserTasks(max_duration=>'20m', max_batch_size=>20);

PostgreSQL

-- SQL query to stream messages
SELECT
  userid,
  messageid,
  payload,
  deliver_time,
  spanner_lease_expiration_timestamp,
  spanner_lease_token,
  spanner_last_batch_message -- Special boolean column returns TRUE if the
                             -- last message is in a batch.
FROM spanner.receive_usertasks(20, NULL, '20m');

或者,您也可以使用客户端库。此 Go 示例演示了如何从队列中流式传输消息、验证租期到期情况以及异步确认消息:

// import "cloud.google.com/go/spanner"
// import "google.golang.org/api/iterator"

stmt := spanner.Statement{SQL: "SELECT * FROM RECEIVE_UserTasks(max_duration=>'20m')"}
iter := client.Single().Query(ctx, stmt)
defer iter.Stop()

for {
    row, err := iter.Next()
    if err == iterator.Done {
        break // Or potentially restart the query
    }
    if err != nil {
        // Handle error
        return err
    }

    var userId int64
    var messageId string
    var payload []byte
    var deliverTime time.Time
    var leaseExpiration time.Time
    var leaseToken string
    var lastBatchMessage bool

    if err := row.Columns(&userId, &messageId, &payload, &deliverTime, &leaseExpiration, &leaseToken, &lastBatchMessage); err != nil {
        // Handle column parsing error
        return err
    }

    if time.Now().After(leaseExpiration) {
        log.Printf("Lease expired for message %s, skipping", messageId)
        continue
    }

    // Process and acknowledge the message asynchronously so that we don't
    // block receiving subsequent messages.
    go func(userId int64, messageId string, payload []byte) {
        // Process the message (payload)
        // ... potentially long-running work ...
        // Need to extend the lease if processing is long

        // Acknowledge the message upon success
        _, err := client.Apply(ctx, []*spanner.Mutation{
            spanner.Ack("UserTasks", spanner.Key{userId, messageId}),
        })
        if err != nil {
            // Handle ack error
        }
    }(userId, messageId, payload)
}

延长消息租约

如果处理消息的时间超过初始租期(10 秒),请使用以下语法:

GoogleSQL

-- Extends the lease and returns new expiration time lease_token.
-- The lease will be extended by 10s (not configurable).
SELECT * FROM RENEWLEASE_UserTasks(lease_tokens => ['token1', ..., 'tokenN'])

-- Returns rows of tokens and whether they were successfully extended
SpannerOldLeaseToken   SpannerNewLeaseToken  SpannerLeaseExpirationTimestamp
        <old_token1>           <new_token1>  "2025-09-27T12:10:00.0Z"
            'token2'                 <NULL>  "2025-09-27T12:09:51.0Z"
...
            'tokenN'           'tokenN_new'  "2025-09-27T12:10:01.0Z"

PostgreSQL

-- Extends the lease and returns new expiration time lease_token.
-- The lease will be extended by 10s (not configurable).
SELECT * FROM spanner.renewlease_usertasks(lease_tokens => ARRAY['token1', ..., 'tokenN'])

-- Returns rows of tokens and whether they were successfully extended
spanner_old_lease_token   spanner_new_lease_token  spanner_lease_expiration_timestamp
           <old_token1>              <new_token1>  "2025-09-27T12:10:00.0Z"
               'token2'                    <NULL>  "2025-09-27T12:09:51.0Z"
...
               'tokenN'              'tokenN_new'  "2025-09-27T12:10:01.0Z"

系统会根据以下逻辑返回租约令牌:

  1. 无法解析的租约令牌不会返回任何行。
  2. 已过期的租约令牌不会返回任何行。
  3. 不可续期的租约令牌确实会返回 SpannerNewLeaseToken 为 NULL 的行。如果消息已被确认,但租约令牌尚未过期,则可能会发生这种情况。

或者,您也可以使用客户端库。此 Go 示例演示了如何延长消息租期:

// Inside message processing loop...

// Before leaseExpiration, for example, in a separate goroutine or timed check
// leaseToken is from the SELECT query
extendStmt := spanner.Statement{
    SQL: "SELECT * FROM RENEWLEASE_UserTasks(lease_tokens => [@token1])",
    Params: map[string]interface{}{
        "token1": leaseToken,
    },
}
_, err := client.Single().Query(ctx, extendStmt).Next() // Simplified call
if err != nil {
    log.Printf("Failed to extend lease for %s: %v", messageId, err)
    // Processing should probably stop as redelivery is likely
} else {
    // New lease expiration is typically approximately 10s from now
    log.Printf("Lease extended for %s", messageId)
    // Update local leaseExpiration time if needed
}

确认消息

使用 DELETE DML 确认消息。这必须与任何其他与消息处理相关的写入操作保持事务性。

GoogleSQL

-- DML for acknowledging
DELETE FROM UserTasks WHERE UserId = @userId AND MessageId = @messageId ASSERT_ROWS_MODIFIED 1;

PostgreSQL

-- DML for acknowledging
DELETE FROM usertasks WHERE userid = $1 AND messageid = $2 ASSERT_ROWS_MODIFIED 1;

或者,使用客户端库 Ack 突变来确认消息:

Go

_, err := client.Apply(ctx, []*spanner.Mutation{
    spanner.Ack("UserTasks", spanner.Key{1}),
})

Java

dbClient.write(
    Collections.singletonList(
        Mutation.newAckBuilder("UserTasks")
            .setKey(Key.of(2L))
            .build()));

最佳做法

以下是使用 Spanner 队列的最佳实践:

  • 小载荷:确保队列消息载荷小于 4 KB。对于较大的数据,请使用带外存储模式
  • 租约管理:延长可能超出默认消息租约(10 秒)的任务的租约。未能延长消息租约可能会导致重新传送消息,并可能导致重复处理。
  • 错误处理:Spanner 队列会在第一个小时内重试处理失败的消息,并采用退避机制,而较旧的消息将每小时重试一次。考虑将永久失败的消息移至单独的队列。
  • 监控:监控队列深度和最早的未确认消息存在时长,确保接收器不会限制流水线输入。
  • 幂等性:设计消息处理器时要使其能够容忍重复操作,因为“至少一次”传送意味着可能会出现重复。如需了解详情,请参阅“正好一次”处理和“最多一次”确认页面。
  • 表值函数时长:避免表值函数的时长过短或过长。建议设置适中的时长,例如 20 分钟。
  • 批次大小:根据工作负载调整 max_batch_size。对于高扇出事件,请使用较小的批次,以避免共享行上的锁争用。对于独立的高每秒查询次数任务,请使用更大的批次。在单个事务中确认或延长批次的消息租约,以获得最佳性能。

监控

您可以使用 Spanner 自省表来监控队列操作。虽然这些表格不包含特定于队列的列,但您可以在以下表格中搜索用户定义的队列名称,以识别队列活动:

以下 Spanner 队列指标位于 Cloud Monitoring 中的 spanner.googleapis.com/queue/* 前缀下:

  • buffered_ready_messages:(GAUGE,INT64,1)保存在内存中且已准备好递送给接收者的消息数量。
  • message_send_count:(DELTA, INT64, 1) 在相应时间间隔内,队列在 Spanner 中发送的消息数。
  • message_ack_count:(DELTA, INT64, 1) 在相应时间间隔内,Spanner 中针对队列确认的邮件数量。
  • oldest_unacked_message_age:(指标,INT64,1)队列中最旧的未确认消息的年龄(以秒为单位)。
  • lease_expiration_count:(DELTA, INT64, 1) 在相应时间间隔内,Spanner 中队列的租约过期次数。

上述所有指标的采样频率约为每 60 秒一次。采样后,数据在最长 120 秒的时间内可能不会显示。 Spanner 审核日志记录涵盖对队列的写入、读取和架构操作。

问题排查

以下部分介绍了如何识别和解决使用 Spanner 队列时的常见问题。

消息积压量不断增加

诊断

oldest_unacked_message_agebuffered_ready_messages 均处于高位。这表示您的消息发送速率与应用的消息处理能力之间存在不平衡。

解决方法

如需解决此问题,请执行以下操作:

  • 检查确认率:如果 message_ack_count 指标有所下降,请检查客户端工作器,确保它们正常运行,没有停滞或崩溃。
  • 检查发送速率:如果 message_send_count 出现峰值,请扩容消息处理工作器以应对增加的负载。
  • 确认资源耗尽:检查 buffered_ready_messageslease_expiration_count 计数是否升高。此组合表示活跃的表值函数 (TVF) 接收器不足或客户端处理速度较慢。

个别消息卡住

诊断

oldest_unacked_message_age 指标较高,但 buffered_ready_messages 指标较低或稳定。这表示单个消息无法处理或确认,而不是整体容量瓶颈。

解决方法

如需解决此问题,请执行以下操作:

  • 识别卡住的消息:查询队列中的表并按传送时间对结果排序,以查找最早的未确认消息:

    GoogleSQL

    SELECT *
    FROM UserTasks
    ORDER BY DeliverTime ASC
    LIMIT 10;
    

    PostgreSQL

    SELECT *
    FROM usertasks
    ORDER BY deliver_time ASC
    LIMIT 10;
    
  • 调查处理失败问题:检查应用日志,确定工作器未确认已识别消息的原因。 未确认的消息会在租约到期后自动重新传送。

消息租约在处理完成之前过期

诊断

lease_expiration_count 指标偏高或在增加。这表示消息处理时间超过了租期(默认值为 10 秒),然后工作器才能确认消息。

解决方法

如需解决此问题,请执行以下操作:

  • 主动续订租约:如果消息处理时间超过 10 秒,请定期调用 RENEWLEASE_QUEUE_NAME() TVF。在处理开始后大约 7-8 秒续订租约,以安全地考虑网络延迟。
  • 调查处理速度缓慢的问题:如果您的应用已主动续订租约,但 lease_expiration_count 仍然很高,请检查后端代码是否存在处理瓶颈、缓慢的 RPC 调用或死锁。

无法调整消息发送速率

诊断

当您尝试提高发送速率时,message_send_count 指标达到吞吐量平台期,或者发布请求遇到写入延迟时间增加的情况。

解决方法

如需解决此问题,请考虑进行以下结构性更改:

  • 纵向扩容计算资源:向 Spanner 实例添加节点或处理单元,以提高数据库的总体容量。
  • 检查分块:评估是否可以添加更多分块,以将写入负载分布到多个服务器上。
  • 优化路由:确保您的发布应用直接写入 Spanner 实例的主要区域,以最大限度地缩短写入延迟时间。

准备就绪的消息未得到处理

诊断

buffered_ready_messages 指标较高且在不断增加。这表示消息已在内存中缓冲并准备好进行传送,但接收器工作器未提取这些消息。

解决方法

如需解决此问题,请执行以下操作:

  • 检查活跃的 TVF 连接:检查活跃的接收器 TVF 数量,确保读取器工作线程正在积极连接并拉取消息。如果工作器已断开连接或未运行足够的并发 RECEIVE_QUEUE_NAME() 查询,则消息会保留在缓冲区中,不会被轮询。

后续步骤