本文說明如何使用 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"
系統會根據下列邏輯傳回租約權杖:
- 無法剖析的租約權杖不會傳回資料列。
- 已過期的租約權杖不會傳回資料列。
- 無法續約的租約權杖會傳回含有 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 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:(GAUGE, INT64, 1) 佇列中未確認訊息的年齡 (以秒為單位)。lease_expiration_count:(DELTA、INT64、1) 佇列間隔期間,Spanner 中租約到期的次數。
系統大約每 60 秒會對上述所有指標取樣一次。取樣完畢後,最多會有 120 秒無法查看資料。 Spanner 稽核記錄涵蓋佇列的寫入、讀取和結構定義作業。
疑難排解
下列各節說明如何找出及解決使用 Spanner 佇列時的常見問題。
待處理訊息數不斷增加
診斷
oldest_unacked_message_age 和 buffered_ready_messages 都已提升。這表示訊息傳送速率與應用程式的訊息處理容量之間存在不平衡。
解析度
如要解決這個問題,請按照下列步驟操作:
- 檢查確認率:如果
message_ack_count指標下降,請檢查用戶端工作站,確保工作站正常運作,沒有停滯或當機。 - 檢查傳送速率:如果
message_send_count突然飆升,請調度更多訊息處理工作站,以因應增加的負載。 - 確認資源耗盡:檢查
buffered_ready_messages和lease_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()查詢,訊息就會留在緩衝區中,不會輪詢。
後續步驟
- 探索更多 Spanner 佇列情境和範例。
- 瞭解僅須處理一次的作業和最多確認一次。
- 使用佇列的精細存取控管設定存取控管。