이 문서에서는 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})
자바
dbClient.write(
Collections.singletonList(
Mutation.newSendBuilder("UserTasks")
.setKey(Key.of(123L, "some-unique-id-1"))
.setPayload(Value.bytes(ByteArray.copyFrom("message3")))
.setDeliveryTime(futureTime)
.build()));
메시지 수신하기
ExecuteStreamingSQL와 함께 RECEIVE_QUEUE_NAME() 테이블 값 함수 (TVF)를 사용하여 메시지를 수신합니다. 이는 장기 실행 호출입니다. 작업자당, 대기열당 이러한 호출 중 하나를 루프 방식으로 실행해야 합니다. 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"
임대 토큰은 다음 논리에 따라 반환됩니다.
- 파싱할 수 없는 리스 토큰은 행을 반환하지 않습니다.
- 이미 만료된 임대 토큰은 행을 반환하지 않습니다.
- 갱신할 수 없는 리스 토큰은
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}),
})
자바
dbClient.write(
Collections.singletonList(
Mutation.newAckBuilder("UserTasks")
.setKey(Key.of(2L))
.build()));
권장사항
다음은 Spanner 대기열 사용에 관한 권장사항입니다.
- 작은 페이로드: 대기열 메시지 페이로드를 4KB 미만으로 작게 유지합니다. 더 큰 데이터에는 대역 외 저장소 패턴을 사용합니다.
- 임대 관리: 기본 메시지 임대 기간인 10초를 초과할 수 있는 작업의 임대 기간을 연장합니다. 메시지 리스를 연장하지 않으면 재전송이 발생하고 이중 처리가 발생할 수 있습니다.
- 오류 처리: Spanner는 첫 시간 내에 처리에 실패한 재시도 메시지를 지수 백오프를 사용하여 대기열에 추가하며, 오래된 메시지는 시간당 한 번 재시도됩니다. 영구적으로 실패한 메시지를 별도의 큐로 이동하는 것이 좋습니다.
- 모니터링: 수신자가 파이프라인 인입을 제한하지 않도록 대기열 깊이와 확인되지 않은 가장 오래된 메시지 기간을 모니터링합니다.
- 멱등성: 최소 1회 전송은 간혹 중복이 발생할 수 있음을 의미하므로 반복 작업을 허용하도록 메시지 프로세서를 설계하세요. 자세한 내용은 정확히 한 번 처리 및 최대 한 번 확인 페이지를 참고하세요.
- 테이블 값 함수 기간: 테이블 값 함수의 기간이 너무 짧거나 너무 길지 않도록 합니다. 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 수를 확인하여 리더 작업자가 활성 상태로 연결되고 메시지를 가져오는지 확인합니다. 작업자가 연결이 끊어졌거나 충분한 동시
RECEIVE_QUEUE_NAME()쿼리를 실행하지 않으면 메시지가 버퍼에서 폴링되지 않은 상태로 유지됩니다.
다음 단계
- Spanner 큐 시나리오 및 예시를 자세히 알아보세요.
- 정확히 한 번 처리 및 최대 한 번 확인에 대해 알아봅니다.
- 대기열에 대한 세분화된 액세스 제어로 액세스 제어를 구성합니다.