使用 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 接收訊息。這是長期通話。您必須以迴圈方式,針對每個佇列中的每個工作人員執行其中一個呼叫。使用強式讀取執行查詢,因為 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. 無法續約的租約權杖傳回含有 NULL SpannerNewLeaseToken 的資料列。如果訊息已確認,但租約權杖尚未過期,就可能發生這種情況。

您也可以使用用戶端程式庫。這個 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 Monitoringspanner.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:(GAUGE, 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 數量,確保讀取器工作人員正在積極連線及提取訊息。如果 worker 已中斷連線,或並未執行足夠的並行 RECEIVE_QUEUE_NAME() 查詢,訊息就會留在緩衝區中,不會輪詢。

後續步驟